Game/ReplayReactionRequests.cs
using System;
using System.Collections.Generic;
using System.Linq;
using System.Net.Http;
using System.Threading;
using System.Threading.Tasks;
using Sandbox;

namespace BlockParty;

/// <summary>One JSON read per replay, with a bounded session cache and shared in-flight reads.</summary>
internal static class ReplayReactionRequests
{
	// Empty disables the feature. Deploy Services/reactions before publishing the game.
	internal static string ServiceUrl { get; set; } = "https://blockparty-reactions.ryly.workers.dev";
	internal static bool Enabled => !string.IsNullOrWhiteSpace( ServiceUrl );
	internal const int MaxReactions = 1000;
	private sealed record Cached( Snapshot Value, DateTime Expires );
	private static readonly Dictionary<string, Cached> Cache = new();
	private static readonly Dictionary<string, Task<Snapshot>> Reads = new();
	private static readonly Dictionary<string, Task> Writes = new();
	private static long _generation;
	internal sealed class Entry
	{
		public string SteamId { get; set; }
		public string Name { get; set; }
		public ReplayReaction Reaction { get; set; }
	}
	internal sealed class Snapshot
	{
		public int Version { get; set; }
		public string ReplayId { get; set; }
		public Entry[] Entries { get; set; }
	}
	private sealed class SaveResult
	{
		public bool Saved { get; set; }
		public bool? Changed { get; set; }
		public string ReplayId { get; set; }
	}

	internal static bool IsSaving( string replay ) => Writes.ContainsKey( replay );

	internal static async Task<Snapshot> ReadAsync( string replay, int length, CancellationToken token )
	{
		if ( !Enabled ) throw new InvalidOperationException( "Reaction service is not configured." );
		// A returning viewer waits for an already-submitted edit, including a failed one,
		// before reading. Leaving the previous visit does not cancel an accepted submission.
		if ( Writes.TryGetValue( replay, out var write ) )
		{
			try { await WaitAsync( write, token ); }
			catch when ( !token.IsCancellationRequested ) { }
		}
		token.ThrowIfCancellationRequested();
		if ( Cache.TryGetValue( replay, out var cache ) && cache.Expires > DateTime.UtcNow ) return cache.Value;
		if ( !Reads.TryGetValue( replay, out var read ) )
		{
			long generation = _generation;
			read = FetchAsync( replay, length, generation );
			Reads[replay] = read;
			_ = RetireReadAsync( replay, read );
		}
		await WaitAsync( read, token );
		return await read;
	}

	private static async Task<Snapshot> FetchAsync( string replay, int length, long generation )
	{
		using var timeout = new CancellationTokenSource( TimeSpan.FromSeconds( 15 ) );
		var snapshot = await Http.RequestJsonAsync<Snapshot>( $"{ServiceUrl}/v1/replays/{replay}/reactions", cancellationToken: timeout.Token );
		timeout.Token.ThrowIfCancellationRequested();
		if ( snapshot?.Version != 1 || snapshot.ReplayId != replay || snapshot.Entries is null
			|| snapshot.Entries.Length > MaxReactions ) throw new Exception( "Invalid reaction snapshot." );
		var seen = new HashSet<long>();
		var valid = new List<Entry>();
		foreach ( var entry in snapshot.Entries )
		{
			if ( entry is null || entry.SteamId?.Length != 17 || !long.TryParse( entry.SteamId, out var id )
				|| id <= 0 || entry.Name?.Length > 80
				|| entry.Reaction?.IsValid( length ) != true || entry.Reaction.Removed )
				continue;
			if ( seen.Add( id ) ) valid.Add( entry );
		}
		// Submitted replay lengths are untrusted. One invalid annotation must not hide
		// other players' reactions; only the local recording supplies the real bounds.
		if ( valid.Count != snapshot.Entries.Length )
			Log.Warning( $"BlockParty: ignored {snapshot.Entries.Length - valid.Count} invalid or duplicate reaction entries." );
		snapshot.Entries = valid.ToArray();
		// A GET that started before a write must not refill the cache with the old data.
		if ( _generation == generation && !IsSaving( replay ) )
		{
			if ( Cache.Count >= 32 && !Cache.ContainsKey( replay ) )
				Cache.Remove( Cache.OrderBy( item => item.Value.Expires ).First().Key );
			Cache[replay] = new Cached( snapshot, DateTime.UtcNow.AddSeconds( 60 ) );
		}
		return snapshot;
	}

	private static async Task RetireReadAsync( string replay, Task<Snapshot> task )
	{
		try { await task; } catch { }
		finally
		{
			if ( Reads.TryGetValue( replay, out var current ) && current == task ) Reads.Remove( replay );
		}
	}

	internal static Task<bool?> SaveAsync( string replay, int length, ReplayReaction reaction )
	{
		if ( !Enabled || IsSaving( replay ) ) throw new InvalidOperationException( "Reaction service is busy." );
		// Invalidate both before and after saving; existing readers can finish for their old
		// visit, but a new visit cannot attach to a GET that predates the write.
		Invalidate( replay );
		var task = SendAsync( replay, length, reaction );
		Writes[replay] = task;
		_ = RetireWriteAsync( replay, task );
		return task;
	}

	private static void Invalidate( string replay )
	{
		Cache.Remove( replay );
		_generation++;
		Reads.Remove( replay );
	}

	private static async Task<bool?> SendAsync( string replay, int length, ReplayReaction reaction )
	{
		string step = "obtaining s&box auth token";
		try
		{
			using var timeout = new CancellationTokenSource( TimeSpan.FromSeconds( 20 ) );
			var tokenRequest = Sandbox.Services.Auth.GetToken( "blockparty-reactions", timeout.Token );
			// Auth.GetToken may ignore cancellation. Bound our wait independently so
			// a stalled token request cannot keep this replay's write lock forever.
			await WaitAsync( tokenRequest, timeout.Token );
			var token = await tokenRequest;
			timeout.Token.ThrowIfCancellationRequested();
			if ( string.IsNullOrWhiteSpace( token ) ) throw new InvalidOperationException( "No s&box auth token was issued." );
			step = "creating request body";
			using var body = Http.CreateJsonContent( new
			{
				SteamId = Game.SteamId.ToString(), Token = token,
				Name = new Friend( (long)Game.SteamId ).Name, Length = length, Reaction = reaction
			} );
			step = "posting to reaction service";
			var result = await Http.RequestJsonAsync<SaveResult>( $"{ServiceUrl}/v1/replays/{replay}/reactions", "POST", body, cancellationToken: timeout.Token );
			timeout.Token.ThrowIfCancellationRequested();
			if ( result?.Saved != true || result.ReplayId != replay ) throw new Exception( "Reaction was not acknowledged." );
			return result.Changed;
		}
		catch ( Exception error )
		{
			// Only log the stage, exception type and HTTP status, never token/body contents.
			string status = error is HttpRequestException http && http.StatusCode.HasValue
				? $"; HTTP {(int)http.StatusCode.Value}" : "";
			Log.Warning( $"BlockParty: reaction save failed while {step} ({error.GetType().Name}{status})." );
			throw;
		}
	}

	private static async Task RetireWriteAsync( string replay, Task task )
	{
		try { await task; } catch { }
		finally
		{
			Invalidate( replay );
			Writes.Remove( replay );
			// The generation also invalidates any detached reads still in flight.
		}
	}

	private static async Task WaitAsync( Task task, CancellationToken token )
	{
		try
		{
			while ( !task.IsCompleted ) await Task.Delay( 50, token );
			token.ThrowIfCancellationRequested();
			await task;
		}
		catch ( OperationCanceledException ) when ( token.IsCancellationRequested )
		{
			// Observe late failures without resuming a timed-out submission.
			_ = ObserveAsync( task );
			throw;
		}
	}

	private static async Task ObserveAsync( Task task )
	{
		try { await task; } catch { }
	}
}