using System.Net.Http;
using System.Net.Http.Json;
using System.Text.Json;
using ytLive.Models;
using AppLog = ytLive.Helpers.AppLog;
namespace ytLive.Services;
///
/// Manages YouTube live stream lifecycle — create broadcasts,
/// bind stream keys, monitor health.
///
public class YouTubeStreamService
{
private readonly YouTubeAuthService _auth;
private readonly HttpClient _http;
private const string ApiBase = "https://www.googleapis.com/youtube/v3";
public YouTubeStreamService(YouTubeAuthService auth, HttpClient? http = null)
{
_auth = auth;
_http = http ?? new HttpClient();
}
private async Task EnsureToken()
{
if (_auth.CurrentChannel == null) return false;
if (_auth.CurrentChannel.TokenExpiry <= DateTime.UtcNow.AddMinutes(5))
return await _auth.RefreshToken();
return true;
}
public async Task CreateBroadcast(string title, string description, DateTime scheduledStartTime, string? streamId = null)
{
if (!await EnsureToken()) return null;
var broadcast = new
{
snippet = new
{
title,
description,
scheduledStartTime = scheduledStartTime.ToString("o"),
categoryId = "22" // People & Blogs
},
status = new
{
// Private-only by enforcement (ship step 7) — the Go Live dialog
// is locked to Private and the service refuses anything else.
privacyStatus = "private",
selfDeclaredMadeForKids = false
},
// One-click go-live (TASK 5 design decision 1): auto start/stop with
// no monitor stream and low latency. A reusable stream, when given,
// binds here (boundStreamId) so no second bind round-trip is needed.
contentDetails = BuildContentDetails(streamId)
};
_http.DefaultRequestHeaders.Authorization = new("Bearer", _auth.CurrentChannel!.AccessToken);
var response = await _http.PostAsJsonAsync(
$"{ApiBase}/liveBroadcasts?part=snippet,status,contentDetails", broadcast);
if (!response.IsSuccessStatusCode) return null;
var json = await response.Content.ReadAsStringAsync();
var data = JsonSerializer.Deserialize(json);
return data.GetProperty("id").GetString();
}
private static Dictionary BuildContentDetails(string? streamId)
{
var details = new Dictionary
{
["enableAutoStart"] = true,
["enableAutoStop"] = true,
["enableMonitorStream"] = false,
["latencyPreference"] = "low",
};
if (streamId != null) details["boundStreamId"] = streamId;
return details;
}
///
/// Pushes the creator-editable broadcast fields to YouTube via
/// liveBroadcasts.update (part=snippet,status). Valid any time, including
/// while live. Returns null on success, otherwise a human-readable error.
/// Note: liveBroadcasts.update REPLACES the snippet part, so scheduledStartTime
/// is re-sent unchanged from the stored value — omitting it would clear the
/// schedule server-side.
///
public async Task UpdateBroadcast(string broadcastId, BroadcastMetadata meta)
{
if (!await EnsureToken()) return "not signed in";
var tags = meta.TagsCsv
.Split(',', StringSplitOptions.RemoveEmptyEntries | StringSplitOptions.TrimEntries)
.ToList();
var broadcast = new
{
id = broadcastId,
snippet = new
{
title = meta.Title,
description = meta.Description,
tags = tags,
// update replaces the whole snippet part; echo the original schedule
scheduledStartTime = (meta.ScheduledStartTime ?? DateTime.UtcNow).ToString("o"),
categoryId = "22"
},
status = new
{
privacyStatus = string.IsNullOrWhiteSpace(meta.Visibility) ? "private" : meta.Visibility.ToLowerInvariant(),
selfDeclaredMadeForKids = meta.MadeForKids
}
};
_http.DefaultRequestHeaders.Authorization = new("Bearer", _auth.CurrentChannel!.AccessToken);
var response = await _http.PutAsJsonAsync(
$"{ApiBase}/liveBroadcasts?part=snippet,status", broadcast);
if (response.IsSuccessStatusCode) return null;
var body = await response.Content.ReadAsStringAsync();
AppLog.Write($"Broadcast update failed ({(int)response.StatusCode}): {body}");
return $"YouTube rejected the update ({(int)response.StatusCode})";
}
/// Proper close-out (TASK 9 design decision 6 — implemented 2026-09-01;
/// until then we relied entirely on enableAutoStop, leaving viewers on a frozen
/// "stream offline" screen for ~a minute): POST liveBroadcasts.transition
/// broadcastStatus=complete. MUST be called AFTER the encoder closed the RTMP
/// push so no frames post-date the end. **Pre-flight (2026-09-22):** the POST is
/// now gated on the broadcast's OWN lifeCycleStatus — every 2026-09-22 end
/// logged 403 invalidTransition because enableAutoStop/YouTube had already marked
/// the broadcast complete; a blind complete only earned the 403 + log noise. We
/// skip ONLY on a confirmed already-ended status (complete/revoked) — an
/// inconclusive check still posts (old behavior) rather than silently stranding a
/// live broadcast. The POST is never-throwing: a rare still-failing transition
/// surfaces as an error string, never a throw — the stop path must not fail over a
/// cosmetic close-out.
public async Task EndBroadcastAsync(string broadcastId)
{
if (!await EnsureToken()) return "not signed in";
_http.DefaultRequestHeaders.Authorization = new("Bearer", _auth.CurrentChannel!.AccessToken);
var lifeCycle = await GetLifeCycleStatusAsync(broadcastId);
if (lifeCycle is "complete" or "revoked")
{
// Already ended (autoStop/YouTube raced us) — a blind complete only
// earns a 403 + noise. enableAutoStop finished it.
AppLog.Write($"End close-out: {broadcastId} already at lifeCycleStatus '{lifeCycle}' — skipping complete");
return null;
}
var response = await _http.PostAsync(
$"{ApiBase}/liveBroadcasts/transition?broadcastStatus=complete&id={Uri.EscapeDataString(broadcastId)}&part=status",
content: null);
if (response.IsSuccessStatusCode) return null;
var body = await response.Content.ReadAsStringAsync();
AppLog.Write($"Broadcast transition(complete) failed ({(int)response.StatusCode}): {body}");
return $"YouTube rejected the end transition ({(int)response.StatusCode})";
}
/// One liveBroadcasts.list (part=status) read of the broadcast's own
/// lifeCycleStatus — the only reliable "can we complete?" signal. Null when the
/// list fails or the status is missing (caller degrades gracefully).
private async Task GetLifeCycleStatusAsync(string broadcastId)
{
var response = await _http.GetAsync(
$"{ApiBase}/liveBroadcasts?part=status&id={Uri.EscapeDataString(broadcastId)}");
if (!response.IsSuccessStatusCode) return null;
var json = JsonSerializer.Deserialize(await response.Content.ReadAsStringAsync());
if (!json.TryGetProperty("items", out var items) || items.GetArrayLength() == 0) return null;
if (!items[0].TryGetProperty("status", out var status)) return null;
return status.TryGetProperty("lifeCycleStatus", out var lifeCycle)
? lifeCycle.GetString()
: null;
}
/// Returns the channel's reusable stream (TASK 5 design decision 2):
/// lists existing streams first and reuses the one with cdn.isReusable=true,
/// creating it with variable resolution/frame rate on first use. Binding to a
/// broadcast happens at broadcast insert (boundStreamId), so one reusable
/// stream serves every broadcast without recreation.
public async Task GetOrCreateReusableStreamAsync()
{
if (!await EnsureToken()) return null;
_http.DefaultRequestHeaders.Authorization = new("Bearer", _auth.CurrentChannel!.AccessToken);
var listResponse = await _http.GetAsync(
$"{ApiBase}/liveStreams?mine=true&part=snippet,cdn,status");
if (!listResponse.IsSuccessStatusCode) return null;
var listJson = JsonSerializer.Deserialize(
await listResponse.Content.ReadAsStringAsync());
if (listJson.TryGetProperty("items", out var items))
{
foreach (var item in items.EnumerateArray())
{
if (item.TryGetProperty("cdn", out var cdn) &&
cdn.TryGetProperty("isReusable", out var reusable) &&
reusable.GetBoolean())
{
var parsed = ParseStream(item);
if (parsed != null) return parsed;
}
}
}
var stream = new
{
snippet = new { title = "LlamaCasty Reusable Stream" },
cdn = new
{
ingestionType = "rtmp",
resolution = "variable",
frameRate = "variable",
isReusable = true
}
};
var response = await _http.PostAsJsonAsync(
$"{ApiBase}/liveStreams?part=snippet,cdn", stream);
if (!response.IsSuccessStatusCode) return null;
var json = JsonSerializer.Deserialize(await response.Content.ReadAsStringAsync());
return ParseStream(json);
}
private static ReusableStream? ParseStream(JsonElement item)
{
if (!item.TryGetProperty("id", out var id) ||
!item.TryGetProperty("cdn", out var cdn) ||
!cdn.TryGetProperty("ingestionInfo", out var info))
{
return null;
}
var streamId = id.GetString();
var address = info.TryGetProperty("ingestionAddress", out var addr) ? addr.GetString() : null;
var name = info.TryGetProperty("streamName", out var nameEl) ? nameEl.GetString() : null;
if (string.IsNullOrWhiteSpace(streamId) || string.IsNullOrWhiteSpace(address) || string.IsNullOrWhiteSpace(name))
return null;
return new ReusableStream(streamId, address, name);
}
/// Polls the reusable stream's health (TASK 5 item 3) via
/// liveStreams.status — report-by-exception: good/ok/noData yield an empty
/// issue list, warning/error entries in configurationIssues[] drive the
/// banner. Null on failure or an empty response, never a throw.
public async Task GetStreamHealthAsync(string streamId)
{
if (!await EnsureToken()) return null;
_http.DefaultRequestHeaders.Authorization = new("Bearer", _auth.CurrentChannel!.AccessToken);
var response = await _http.GetAsync(
$"{ApiBase}/liveStreams?part=status&id={streamId}");
if (!response.IsSuccessStatusCode) return null;
var json = JsonSerializer.Deserialize(await response.Content.ReadAsStringAsync());
var items = json.GetProperty("items");
if (items.GetArrayLength() == 0) return null;
var status = items[0].GetProperty("status");
// The real API nests health as status.healthStatus = { status, lastUpdateTimeSeconds,
// configurationIssues[] }. The flat "healthStatus":"bad" shape is our old wrong
// assumption — accept both so neither crashes the report-by-exception poll.
var health = new StreamHealth();
var issueElements = new List();
if (status.TryGetProperty("healthStatus", out var healthStatus))
{
if (healthStatus.ValueKind == JsonValueKind.String)
{
health.HealthStatus = healthStatus.GetString();
}
else if (healthStatus.TryGetProperty("status", out var inner) &&
inner.ValueKind == JsonValueKind.String)
{
health.HealthStatus = inner.GetString();
}
if (healthStatus.TryGetProperty("configurationIssues", out var nested))
issueElements.AddRange(nested.EnumerateArray());
}
if (status.TryGetProperty("configurationIssues", out var flat))
issueElements.AddRange(flat.EnumerateArray());
foreach (var issue in issueElements)
{
var severity = issue.TryGetProperty("severity", out var sev) && sev.ValueKind == JsonValueKind.String
? sev.GetString()
: null;
var type = issue.TryGetProperty("type", out var t) && t.ValueKind == JsonValueKind.String
? t.GetString()
: null;
health.ConfigurationIssues.Add(new StreamConfigurationIssue
{
Severity = severity switch
{
"error" => StreamIssueSeverity.Error,
"warning" => StreamIssueSeverity.Warning,
_ => StreamIssueSeverity.Info,
},
Type = type,
});
}
return health;
}
/// TASK 41 — mock chat input: posts a real text message into the live
/// chat via liveChat/messages.insert. Requires the same OAuth scopes the
/// app already holds (youtube.force-ssl). Returns itself or a human-readable
/// error; never throws. The inserted message round-trips back through the normal
/// chat poll (~2s) and renders through the live overlay path.
public async Task InsertChatMessageAsync(string liveChatId, string messageText)
{
if (!await EnsureToken()) return "not signed in";
var body = new
{
snippet = new
{
liveChatId,
type = "textMessageEvent",
textMessageDetails = new { messageText }
}
};
_http.DefaultRequestHeaders.Authorization = new("Bearer", _auth.CurrentChannel!.AccessToken);
var response = await _http.PostAsJsonAsync(
$"{ApiBase}/liveChat/messages?part=snippet", body);
if (response.IsSuccessStatusCode)
{
AppLog.Write($"Test chat message inserted into {liveChatId}");
return null;
}
var errorBody = await response.Content.ReadAsStringAsync();
AppLog.Write($"Chat message insert failed ({(int)response.StatusCode}): {errorBody}");
return $"YouTube rejected the message ({(int)response.StatusCode})";
}
/// Resolves the broadcast's liveChatId (needed to poll chat). Two facts
/// pin this down (2026-09-25, from the official liveBroadcasts reference + the
/// GetLiveChatId.java sample):
/// 1. It lives under snippet.liveChatId — contentDetails has no such
/// property, so part=contentDetails could NEVER resolve it (the "Chat polling
/// couldn't start" report).
/// 2. YouTube only populates it once the broadcast is live (the sample lists
/// broadcastStatus=active) — fetching right after insert (lifecycleStatus
/// ready) returns no id. So this polls with a bounded retry, which the
/// caller must run AFTER the encoder starts pushing RTMP (enableAutoStart flips
/// the broadcast to live). Returns the id or null after the retries exhaust;
/// never throws.
public async Task GetBroadcastLiveChatIdAsync(
string broadcastId, int maxAttempts = 10, int delayMs = 2000)
{
for (var attempt = 0; attempt < maxAttempts; attempt++)
{
if (attempt > 0) await Task.Delay(delayMs);
var liveChatId = await FetchLiveChatIdAsync(broadcastId);
if (liveChatId != null) return liveChatId;
}
return null;
}
private async Task FetchLiveChatIdAsync(string broadcastId)
{
if (!await EnsureToken()) return null;
_http.DefaultRequestHeaders.Authorization = new("Bearer", _auth.CurrentChannel!.AccessToken);
var response = await _http.GetAsync(
$"{ApiBase}/liveBroadcasts?part=snippet&id={broadcastId}");
if (!response.IsSuccessStatusCode) return null;
var json = await response.Content.ReadAsStringAsync();
var data = JsonSerializer.Deserialize(json);
var items = data.GetProperty("items");
if (items.GetArrayLength() == 0) return null;
var snippet = items[0].GetProperty("snippet");
if (snippet.TryGetProperty("liveChatId", out var liveChatId))
return liveChatId.GetString();
return null;
}
}