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 { }
}
}