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, 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})"; } /// Returns the liveChatId for a broadcast by fetching its contentDetails. /// The liveChatId is needed to poll chat messages. Called after broadcast creation /// in PrepareAndStartLiveAsync. public async Task GetBroadcastLiveChatIdAsync(string broadcastId) { if (!await EnsureToken()) return null; _http.DefaultRequestHeaders.Authorization = new("Bearer", _auth.CurrentChannel!.AccessToken); var response = await _http.GetAsync( $"{ApiBase}/liveBroadcasts?part=contentDetails&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 contentDetails = items[0].GetProperty("contentDetails"); if (contentDetails.TryGetProperty("liveChatId", out var liveChatId)) return liveChatId.GetString(); return null; } }