Files
LlamaCasty/Services/YouTubeStreamService.cs
gramps 0981a72eea fix(stream): chat insert 400 — insert body must declare snippet.type
The liveChatId fix (TASK 44) worked on the next Test Stream, but every Mock
Chat Input send returned 'YouTube rejected the message (error 400)'. The
runtime log's body: 400 MISSING_REQUIRED_FIELD, domain
youtube.api.v3.LiveChatMessageInsertResponse.Error.

The liveChat/messages.insert snippet requires type ('textMessageEvent' or
'pollEvent') alongside liveChatId and textMessageDetails.messageText; the
TASK 41 body omitted it, and the Good Dog test asserted only liveChatId +
messageText were present — false-green while real YouTube rejected every
send. Fix body + assert type in the same test so the field can never drop
silently again.

Reference: https://developers.google.com/youtube/v3/live/docs/liveChatMessages/insert

Good Dog: ONE integration test (strengthened TestStream_DockTooling...).
Gate: clean build 0 warnings; full vstest 317/318 (the 1 failure is the
known environmental RealMouseDrag flake — passes 3/3 in isolation).
2026-09-25 08:12:05 -07:00

388 lines
17 KiB
C#

using System.Net.Http;
using System.Net.Http.Json;
using System.Text.Json;
using ytLive.Models;
using AppLog = ytLive.Helpers.AppLog;
namespace ytLive.Services;
/// <summary>
/// Manages YouTube live stream lifecycle — create broadcasts,
/// bind stream keys, monitor health.
/// </summary>
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<bool> EnsureToken()
{
if (_auth.CurrentChannel == null) return false;
if (_auth.CurrentChannel.TokenExpiry <= DateTime.UtcNow.AddMinutes(5))
return await _auth.RefreshToken();
return true;
}
public async Task<string?> 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<JsonElement>(json);
return data.GetProperty("id").GetString();
}
private static Dictionary<string, object?> BuildContentDetails(string? streamId)
{
var details = new Dictionary<string, object?>
{
["enableAutoStart"] = true,
["enableAutoStop"] = true,
["enableMonitorStream"] = false,
["latencyPreference"] = "low",
};
if (streamId != null) details["boundStreamId"] = streamId;
return details;
}
/// <summary>
/// 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.
/// </summary>
public async Task<string?> 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})";
}
/// <summary>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.</summary>
public async Task<string?> 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})";
}
/// <summary>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).</summary>
private async Task<string?> 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<JsonElement>(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;
}
/// <summary>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.</summary>
public async Task<ReusableStream?> 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<JsonElement>(
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<JsonElement>(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);
}
/// <summary>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.</summary>
public async Task<StreamHealth?> 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<JsonElement>(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<JsonElement>();
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;
}
/// <summary>TASK 41 — mock chat input: posts a real text message into the live
/// chat via <c>liveChat/messages.insert</c>. 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.</summary>
public async Task<string?> 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})";
}
/// <summary>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 <c>snippet.liveChatId</c> — <c>contentDetails</c> 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 <c>live</c> (the sample lists
/// broadcastStatus=active) — fetching right after insert (lifecycleStatus
/// <c>ready</c>) 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.</summary>
public async Task<string?> 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<string?> 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<JsonElement>(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;
}
}