From d7367e6881291c8e650fb97d32c72b4b82b859cc Mon Sep 17 00:00:00 2001 From: Dustin Byrne Date: Tue, 1 Sep 2026 12:14:54 -0400 Subject: [PATCH 1/2] fix: preserve durable queue entries across retries --- .changeset/durable-retry-queue.md | 5 + com.posthog.unity/Runtime/Core/EventQueue.cs | 317 +++++++++++----- .../Runtime/Core/NetworkClient.cs | 151 +++++++- .../Runtime/SessionReplay/ReplayQueue.cs | 171 +++++---- tests/PostHog.Unity.Tests/EventQueueTests.cs | 353 ++++++++++++++++++ .../PostHog.Unity.Tests/NetworkClientTests.cs | 148 ++++++++ tests/PostHog.Unity.Tests/ReplayQueueTests.cs | 176 ++++++++- 7 files changed, 1125 insertions(+), 196 deletions(-) create mode 100644 .changeset/durable-retry-queue.md create mode 100644 tests/PostHog.Unity.Tests/EventQueueTests.cs diff --git a/.changeset/durable-retry-queue.md b/.changeset/durable-retry-queue.md new file mode 100644 index 0000000..3935299 --- /dev/null +++ b/.changeset/durable-retry-queue.md @@ -0,0 +1,5 @@ +--- +'com.posthog.unity': patch +--- + +Preserve analytics and session replay events across transient ingestion failures, honor server retry delays, safely handle oversized payloads, and prevent in-flight queue replacements from being acknowledged as sent. diff --git a/com.posthog.unity/Runtime/Core/EventQueue.cs b/com.posthog.unity/Runtime/Core/EventQueue.cs index 0309ac9..b24b51b 100644 --- a/com.posthog.unity/Runtime/Core/EventQueue.cs +++ b/com.posthog.unity/Runtime/Core/EventQueue.cs @@ -12,7 +12,9 @@ class EventQueue { readonly PostHogConfig _config; readonly IStorageProvider _storage; - readonly NetworkClient _networkClient; + readonly Func, IEnumerator> _sendBatch; + readonly Func _utcNow; + readonly Func _isConnected; readonly object _lock = new(); bool _isRunning; @@ -27,24 +29,46 @@ class EventQueue int _adjustedMaxBatchSize; int _adjustedFlushAt; - const int RetryDelaySeconds = 5; - const int MaxRetryDelaySeconds = 30; - public EventQueue( PostHogConfig config, IStorageProvider storage, NetworkClient networkClient + ) + : this( + config, + storage, + networkClient.SendBatch, + () => DateTime.UtcNow, + () => Application.internetReachability != NetworkReachability.NotReachable + ) { } + + internal EventQueue( + PostHogConfig config, + IStorageProvider storage, + Func, IEnumerator> sendBatch, + Func utcNow, + Func isConnected ) { _config = config; _storage = storage; - _networkClient = networkClient; + _sendBatch = sendBatch; + _utcNow = utcNow; + _isConnected = isConnected; // Initialize adjusted values from config _adjustedMaxBatchSize = config.MaxBatchSize; _adjustedFlushAt = config.FlushAt; + + lock (_lock) + { + var dropped = TrimQueueToSize(QueueCapacity); + LogCapacityDrops(dropped); + } } + int QueueCapacity => Math.Max(1, _config.MaxQueueSize); + /// /// Starts the automatic flush timer. /// @@ -73,21 +97,13 @@ public void Enqueue(PostHogEvent evt) { lock (_lock) { - // Check queue size limit - var eventIds = _storage.GetEventIds(); - if (eventIds.Count >= _config.MaxQueueSize) - { - // Drop oldest event - MaxQueueSize is validated to be >= 1, - // so eventIds.Count > 0 is guaranteed here - _storage.DeleteEvent(eventIds[0]); - PostHogLogger.Warning( - $"Queue full ({_config.MaxQueueSize}), dropped oldest event" - ); - } + var dropped = TrimQueueToSize(QueueCapacity - 1); + LogCapacityDrops(dropped); - // Save event + // Queue identity is deliberately independent from the mutable payload UUID. + var entryId = GenerateUniqueEntryId(); var json = JsonSerializer.SerializeEvent(evt); - _storage.SaveEvent(evt.Uuid, json); + _storage.SaveEvent(entryId, json); PostHogLogger.Debug($"Enqueued event: {evt.Event}"); } @@ -138,6 +154,56 @@ public int Count } } + int TrimQueueToSize(int targetSize) + { + var dropped = 0; + while (_storage.GetEventCount() > targetSize) + { + var eventIds = _storage.GetEventIds(); + if (eventIds.Count == 0) + { + break; + } + + _storage.DeleteEvent(eventIds[0]); + dropped++; + } + return dropped; + } + + void LogCapacityDrops(int dropped) + { + if (dropped > 0) + { + PostHogLogger.Warning( + $"Queue full ({QueueCapacity}), dropped {dropped} oldest event(s)" + ); + } + } + + string GenerateUniqueEntryId() + { + while (true) + { + var entryId = UuidV7.Generate(); + var eventIds = _storage.GetEventIds(); + var exists = false; + for (int i = 0; i < eventIds.Count; i++) + { + if (eventIds[i] == entryId) + { + exists = true; + break; + } + } + + if (!exists) + { + return entryId; + } + } + } + void FlushIfOverThreshold() { int count = Count; @@ -195,18 +261,18 @@ public IEnumerator FlushCoroutine() _isFlushing = true; } - // Check if paused due to errors - if (_pausedUntil.HasValue && DateTime.UtcNow < _pausedUntil.Value) + // Known-offline periods do not consume a request or retry state. + if (!_isConnected()) { - PostHogLogger.Debug($"Queue paused until {_pausedUntil.Value}"); + PostHogLogger.Debug("No network connectivity, skipping flush"); _isFlushing = false; yield break; } - // Check network connectivity - if (Application.internetReachability == NetworkReachability.NotReachable) + // Check if paused due to errors + if (_pausedUntil.HasValue && _utcNow() < _pausedUntil.Value) { - PostHogLogger.Debug("No network connectivity, skipping flush"); + PostHogLogger.Debug($"Queue paused until {_pausedUntil.Value}"); _isFlushing = false; yield break; } @@ -224,68 +290,72 @@ public IEnumerator FlushCoroutine() // Create batch list without LINQ allocation int batchSize = Math.Min(eventIds.Count, _adjustedMaxBatchSize); - var batchIds = new List(batchSize); + var candidateIds = new List(batchSize); for (int i = 0; i < batchSize; i++) { - batchIds.Add(eventIds[i]); + candidateIds.Add(eventIds[i]); } - var events = LoadEvents(batchIds); + var entries = LoadEvents(candidateIds); - if (events.Count == 0) + if (entries.Count == 0) { - break; + // Missing or corrupt entries were removed; continue with the next records. + continue; + } + + var events = new List(entries.Count); + foreach (var entry in entries) + { + events.Add(entry.Event); } PostHogLogger.Debug($"Flushing batch of {events.Count} events"); var payload = new BatchPayload(_config.ApiKey, events); - bool success = false; - bool shouldDeleteEvents = true; + var uploadResult = new BatchUploadResult(false, 0); - yield return _networkClient.SendBatch( - payload, - (result, statusCode) => - { - success = result; - shouldDeleteEvents = ShouldDeleteEventsOnError(statusCode); - } - ); + yield return _sendBatch(payload, result => uploadResult = result); - if (success) + if (uploadResult.Success) { - // Delete successfully sent events - foreach (var id in batchIds) - { - _storage.DeleteEvent(id); - } - - _retryCount = 0; - _pausedUntil = null; + DeleteEntries(entries); + ResetRetryState(); PostHogLogger.Debug($"Successfully sent {events.Count} events"); + continue; } - else + + if (uploadResult.StatusCode == 413) { - if (shouldDeleteEvents) + if (entries.Count == 1) { - // Delete events on non-retryable errors - foreach (var id in batchIds) - { - _storage.DeleteEvent(id); - } + DeleteEntries(entries); + ResetRetryState(); + PostHogLogger.Warning( + "Dropped oversized event after a singleton batch received HTTP 413" + ); } else { - // Pause for retry - _retryCount++; - int delay = Math.Min( - _retryCount * RetryDelaySeconds, - MaxRetryDelaySeconds + _adjustedMaxBatchSize = RetryQueuePolicy.ReducedBatchSize( + entries.Count ); - _pausedUntil = DateTime.UtcNow.AddSeconds(delay); - PostHogLogger.Warning($"Flush failed, retrying in {delay}s"); + _adjustedFlushAt = Math.Min(_adjustedFlushAt, _adjustedMaxBatchSize); + PostHogLogger.Warning( + $"Payload too large, reducing batch size to {_adjustedMaxBatchSize}" + ); + PauseForRetry(uploadResult.RetryAfter, "Flush failed"); } - break; } + else if (RetryQueuePolicy.ShouldDelete(uploadResult.StatusCode)) + { + DeleteEntries(entries); + ResetRetryState(); + } + else + { + PauseForRetry(uploadResult.RetryAfter, "Flush failed"); + } + break; } } finally @@ -294,27 +364,55 @@ public IEnumerator FlushCoroutine() } } - List LoadEvents(List eventIds) + void PauseForRetry(TimeSpan? retryAfter, string message) { - var events = new List(); + if (_retryCount < int.MaxValue) + { + _retryCount++; + } + var delay = RetryQueuePolicy.GetRetryDelay(_retryCount, retryAfter); + _pausedUntil = RetryQueuePolicy.AddDelay(_utcNow(), delay); + PostHogLogger.Warning($"{message}, retrying in {delay.TotalSeconds:0.###}s"); + } + + void ResetRetryState() + { + _retryCount = 0; + _pausedUntil = null; + } + + void DeleteEntries(List entries) + { + foreach (var entry in entries) + { + _storage.DeleteEvent(entry.StorageId); + } + } + + List LoadEvents(List eventIds) + { + var events = new List(); foreach (var id in eventIds) { try { var json = _storage.LoadEvent(id); - if (!string.IsNullOrEmpty(json)) + if (string.IsNullOrEmpty(json)) { - var evt = DeserializeEvent(json); - if (evt != null) - { - events.Add(evt); - } - else - { - // Corrupted event, delete it - _storage.DeleteEvent(id); - } + _storage.DeleteEvent(id); + continue; + } + + var evt = DeserializeEvent(json); + if (evt != null) + { + events.Add(new QueuedEvent(id, evt)); + } + else + { + // Corrupted event, delete it + _storage.DeleteEvent(id); } } catch (Exception ex) @@ -363,40 +461,59 @@ PostHogEvent DeserializeEvent(string json) } } - bool ShouldDeleteEventsOnError(int statusCode) + sealed class QueuedEvent { - // Retry on network errors (0) and redirects (3xx) - if (statusCode == 0 || (statusCode >= 300 && statusCode < 400)) + public QueuedEvent(string storageId, PostHogEvent evt) { - return false; + StorageId = storageId; + Event = evt; } - // Don't retry on client errors (4xx) except 413 (payload too large) - if (statusCode >= 400 && statusCode < 500 && statusCode != 413) - { - return true; - } + public string StorageId { get; } + public PostHogEvent Event { get; } + } + } - // Retry on 413 (handled separately by reducing batch size) - if (statusCode == 413) - { - // Reduce batch size for next attempt (use local adjusted values, not config) - _adjustedMaxBatchSize = Math.Max(1, _adjustedMaxBatchSize / 2); - _adjustedFlushAt = Math.Max(1, _adjustedFlushAt / 2); - PostHogLogger.Warning( - $"Payload too large, reducing batch size to {_adjustedMaxBatchSize}" - ); - return false; - } + static class RetryQueuePolicy + { + const int RetryDelaySeconds = 5; + const int MaxRetryDelaySeconds = 30; - // Retry on server errors (5xx) - if (statusCode >= 500) + public static bool ShouldDelete(int statusCode) + { + if ( + statusCode == 0 + || statusCode == 408 + || statusCode == 413 + || statusCode == 429 + || (statusCode >= 300 && statusCode < 400) + || statusCode >= 500 + ) { return false; } - // Default: don't delete (retry) - return false; + return statusCode >= 400 && statusCode < 500; + } + + public static int ReducedBatchSize(int failedBatchSize) + { + return Math.Max(1, failedBatchSize / 2); + } + + public static TimeSpan GetRetryDelay(int retryCount, TimeSpan? retryAfter) + { + var localSeconds = Math.Min((long)retryCount * RetryDelaySeconds, MaxRetryDelaySeconds); + var localDelay = TimeSpan.FromSeconds(localSeconds); + return retryAfter.HasValue && retryAfter.Value > localDelay + ? retryAfter.Value + : localDelay; + } + + public static DateTime AddDelay(DateTime now, TimeSpan delay) + { + var remaining = DateTime.MaxValue - now; + return delay >= remaining ? DateTime.MaxValue : now.Add(delay); } } } diff --git a/com.posthog.unity/Runtime/Core/NetworkClient.cs b/com.posthog.unity/Runtime/Core/NetworkClient.cs index 1292350..6d136dc 100644 --- a/com.posthog.unity/Runtime/Core/NetworkClient.cs +++ b/com.posthog.unity/Runtime/Core/NetworkClient.cs @@ -1,6 +1,7 @@ using System; using System.Collections; using System.Collections.Generic; +using System.Globalization; using System.Text; using UnityEngine; using UnityEngine.Networking; @@ -15,6 +16,7 @@ class NetworkClient readonly PostHogConfig _config; readonly FeatureFlagsRequestFactory _featureFlagsRequestFactory; readonly FeatureFlagsRetryDelayFactory _featureFlagsRetryDelayFactory; + readonly BatchRequestFactory _batchRequestFactory; const int TimeoutSeconds = 10; const float FeatureFlagsInitialRetryDelaySeconds = 0.3f; @@ -30,6 +32,17 @@ Dictionary> groupProperties internal delegate object FeatureFlagsRetryDelayFactory(int failedAttempt); + internal delegate IBatchRequest BatchRequestFactory(string url, byte[] body); + + internal interface IBatchRequest : IDisposable + { + UnityWebRequest.Result Result { get; } + long ResponseCode { get; } + string Error { get; } + object Send(); + string GetResponseHeader(string name); + } + internal interface IFeatureFlagsRequest : IDisposable { string Url { get; } @@ -44,52 +57,67 @@ public NetworkClient(PostHogConfig config) : this( config, CreateFeatureFlagsRequest, - failedAttempt => new WaitForSeconds(GetFeatureFlagsRetryDelaySeconds(failedAttempt)) + failedAttempt => new WaitForSeconds( + GetFeatureFlagsRetryDelaySeconds(failedAttempt) + ), + CreateBatchRequest ) { } internal NetworkClient( PostHogConfig config, FeatureFlagsRequestFactory featureFlagsRequestFactory, FeatureFlagsRetryDelayFactory featureFlagsRetryDelayFactory + ) + : this( + config, + featureFlagsRequestFactory, + featureFlagsRetryDelayFactory, + CreateBatchRequest + ) { } + + internal NetworkClient( + PostHogConfig config, + FeatureFlagsRequestFactory featureFlagsRequestFactory, + FeatureFlagsRetryDelayFactory featureFlagsRetryDelayFactory, + BatchRequestFactory batchRequestFactory ) { _config = config; _featureFlagsRequestFactory = featureFlagsRequestFactory; _featureFlagsRetryDelayFactory = featureFlagsRetryDelayFactory; + _batchRequestFactory = batchRequestFactory; } /// /// Sends a batch of events to the PostHog API. /// /// The batch payload to send - /// Callback with (success, statusCode) - public IEnumerator SendBatch(BatchPayload payload, Action onComplete) + /// Callback with the upload result and retry metadata. + public IEnumerator SendBatch(BatchPayload payload, Action onComplete) { var url = GetBatchUrl(); var json = JsonSerializer.SerializeBatch(payload); PostHogLogger.Debug($"Sending batch to {url}"); - using var request = new UnityWebRequest(url, "POST"); - var bodyRaw = Encoding.UTF8.GetBytes(json); - request.uploadHandler = new UploadHandlerRaw(bodyRaw); - request.downloadHandler = new DownloadHandlerBuffer(); - ApplyDefaultHeaders(request); - request.timeout = TimeoutSeconds; - - yield return request.SendWebRequest(); + using var request = _batchRequestFactory(url, Encoding.UTF8.GetBytes(json)); + yield return request.Send(); - int statusCode = (int)request.responseCode; + int statusCode = (int)request.ResponseCode; + var retryAfter = ParseRetryAfter( + request.GetResponseHeader("Retry-After"), + DateTimeOffset.UtcNow + ); - if (request.result == UnityWebRequest.Result.Success) + if (request.Result == UnityWebRequest.Result.Success) { PostHogLogger.Debug($"Batch sent successfully (status: {statusCode})"); - onComplete?.Invoke(true, statusCode); + onComplete?.Invoke(new BatchUploadResult(true, statusCode, retryAfter)); } else { - PostHogLogger.Warning($"Batch send failed: {request.error} (status: {statusCode})"); - onComplete?.Invoke(false, statusCode); + PostHogLogger.Warning($"Batch send failed: {request.Error} (status: {statusCode})"); + onComplete?.Invoke(new BatchUploadResult(false, statusCode, retryAfter)); } } @@ -168,6 +196,44 @@ string GetBatchUrl() return $"{host}/batch"; } + internal static TimeSpan? ParseRetryAfter(string value, DateTimeOffset now) + { + if (string.IsNullOrWhiteSpace(value)) + { + return null; + } + + if ( + long.TryParse( + value.Trim(), + NumberStyles.None, + CultureInfo.InvariantCulture, + out var seconds + ) + && seconds >= 0 + ) + { + var maxSeconds = (long)TimeSpan.MaxValue.TotalSeconds; + return TimeSpan.FromSeconds(Math.Min(seconds, maxSeconds)); + } + + if ( + DateTimeOffset.TryParse( + value, + CultureInfo.InvariantCulture, + DateTimeStyles.AllowWhiteSpaces + | DateTimeStyles.AssumeUniversal + | DateTimeStyles.AdjustToUniversal, + out var retryAt + ) + ) + { + return retryAt <= now ? TimeSpan.Zero : retryAt - now; + } + + return null; + } + internal static bool ShouldRetryFeatureFlagsRequest( UnityWebRequest.Result result, int statusCode, @@ -210,6 +276,45 @@ internal static float GetFeatureFlagsRetryDelaySeconds(int failedAttempt) return FeatureFlagsInitialRetryDelaySeconds * (1 << (failedAttempt - 1)); } + static IBatchRequest CreateBatchRequest(string url, byte[] body) + { + var request = new UnityWebRequest(url, "POST"); + request.uploadHandler = new UploadHandlerRaw(body); + request.downloadHandler = new DownloadHandlerBuffer(); + ApplyDefaultHeaders(request); + request.timeout = TimeoutSeconds; + return new UnityWebRequestBatchRequest(request); + } + + sealed class UnityWebRequestBatchRequest : IBatchRequest + { + readonly UnityWebRequest _request; + + public UnityWebRequestBatchRequest(UnityWebRequest request) + { + _request = request; + } + + public UnityWebRequest.Result Result => _request.result; + public long ResponseCode => _request.responseCode; + public string Error => _request.error; + + public object Send() + { + return _request.SendWebRequest(); + } + + public string GetResponseHeader(string name) + { + return _request.GetResponseHeader(name); + } + + public void Dispose() + { + _request.Dispose(); + } + } + static IFeatureFlagsRequest CreateFeatureFlagsRequest( string apiKey, string host, @@ -332,4 +437,18 @@ static void ApplyDefaultHeaders(UnityWebRequest request) request.SetRequestHeader("User-Agent", SdkInfo.UserAgent); } } + + class BatchUploadResult + { + public BatchUploadResult(bool success, int statusCode, TimeSpan? retryAfter = null) + { + Success = success; + StatusCode = statusCode; + RetryAfter = retryAfter; + } + + public bool Success { get; } + public int StatusCode { get; } + public TimeSpan? RetryAfter { get; } + } } diff --git a/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs b/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs index 4533ae7..a231847 100644 --- a/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs +++ b/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs @@ -20,6 +20,9 @@ class ReplayQueue readonly string _host; readonly Func _getDistinctId; readonly Func _getSessionId; + readonly Func, Action, IEnumerator> _sendBatch; + readonly Func _utcNow; + readonly Func _isConnected; readonly object _lock = new(); readonly List _queue = new(); @@ -29,10 +32,9 @@ class ReplayQueue MonoBehaviour _coroutineRunner; DateTime? _pausedUntil; int _retryCount; + int _adjustedMaxBatchSize = MaxBatchSize; const int TimeoutSeconds = 30; - const int RetryDelaySeconds = 5; - const int MaxRetryDelaySeconds = 30; const int MaxBatchSize = 10; // Snapshots are large, so smaller batches public ReplayQueue( @@ -48,6 +50,31 @@ Func getSessionId _host = host.TrimEnd('/'); _getDistinctId = getDistinctId; _getSessionId = getSessionId; + _sendBatch = SendBatch; + _utcNow = () => DateTime.UtcNow; + _isConnected = () => + Application.internetReachability != NetworkReachability.NotReachable; + } + + internal ReplayQueue( + PostHogSessionReplayConfig config, + string apiKey, + string host, + Func getDistinctId, + Func getSessionId, + Func, Action, IEnumerator> sendBatch, + Func utcNow, + Func isConnected + ) + { + _config = config; + _apiKey = apiKey; + _host = host.TrimEnd('/'); + _getDistinctId = getDistinctId; + _getSessionId = getSessionId; + _sendBatch = sendBatch; + _utcNow = utcNow; + _isConnected = isConnected; } /// @@ -199,7 +226,7 @@ IEnumerator FlushTimerCoroutine() } } - IEnumerator FlushCoroutine() + internal IEnumerator FlushCoroutine() { lock (_lock) { @@ -211,16 +238,17 @@ IEnumerator FlushCoroutine() _isFlushing = true; } - if (_pausedUntil.HasValue && DateTime.UtcNow < _pausedUntil.Value) + // Known-offline periods do not consume a request or retry state. + if (!_isConnected()) { - PostHogLogger.Debug($"Replay queue paused until {_pausedUntil.Value}"); + PostHogLogger.Debug("No network connectivity, skipping replay flush"); _isFlushing = false; yield break; } - if (Application.internetReachability == NetworkReachability.NotReachable) + if (_pausedUntil.HasValue && _utcNow() < _pausedUntil.Value) { - PostHogLogger.Debug("No network connectivity, skipping replay flush"); + PostHogLogger.Debug($"Replay queue paused until {_pausedUntil.Value}"); _isFlushing = false; yield break; } @@ -235,7 +263,7 @@ IEnumerator FlushCoroutine() if (_queue.Count == 0) break; - int batchSize = Math.Min(_queue.Count, MaxBatchSize); + int batchSize = Math.Min(_queue.Count, _adjustedMaxBatchSize); batch = new List(_queue.GetRange(0, batchSize)); } @@ -244,58 +272,46 @@ IEnumerator FlushCoroutine() PostHogLogger.Debug($"Flushing batch of {batch.Count} replay events"); - bool success = false; - int statusCode = 0; + var uploadResult = new BatchUploadResult(false, 0); + yield return _sendBatch(batch, result => uploadResult = result); - yield return SendBatch( - batch, - (result, code) => - { - success = result; - statusCode = code; - } - ); - - if (success) + if (uploadResult.Success) { - lock (_lock) - { - for (int i = 0; i < batch.Count && _queue.Count > 0; i++) - { - _queue.RemoveAt(0); - } - } - - _retryCount = 0; - _pausedUntil = null; + RemoveBatch(batch); + ResetRetryState(); PostHogLogger.Debug($"Successfully sent {batch.Count} replay events"); + continue; } - else + + if (uploadResult.StatusCode == 413) { - bool shouldDelete = ShouldDeleteEventsOnError(statusCode); - if (shouldDelete) + if (batch.Count == 1) { - lock (_lock) - { - for (int i = 0; i < batch.Count && _queue.Count > 0; i++) - { - _queue.RemoveAt(0); - } - } + RemoveBatch(batch); + ResetRetryState(); + PostHogLogger.Warning( + "Dropped oversized replay event after a singleton batch received HTTP 413" + ); } else { - // Pause for retry - _retryCount++; - int delay = Math.Min( - _retryCount * RetryDelaySeconds, - MaxRetryDelaySeconds + _adjustedMaxBatchSize = RetryQueuePolicy.ReducedBatchSize(batch.Count); + PostHogLogger.Warning( + $"Replay payload too large, reducing batch size to {_adjustedMaxBatchSize}" ); - _pausedUntil = DateTime.UtcNow.AddSeconds(delay); - PostHogLogger.Warning($"Replay flush failed, retrying in {delay}s"); + PauseForRetry(uploadResult.RetryAfter); } - break; } + else if (RetryQueuePolicy.ShouldDelete(uploadResult.StatusCode)) + { + RemoveBatch(batch); + ResetRetryState(); + } + else + { + PauseForRetry(uploadResult.RetryAfter); + } + break; } } finally @@ -304,7 +320,35 @@ IEnumerator FlushCoroutine() } } - IEnumerator SendBatch(List events, Action onComplete) + void RemoveBatch(List batch) + { + lock (_lock) + { + foreach (var sentEvent in batch) + { + _queue.Remove(sentEvent); + } + } + } + + void PauseForRetry(TimeSpan? retryAfter) + { + if (_retryCount < int.MaxValue) + { + _retryCount++; + } + var delay = RetryQueuePolicy.GetRetryDelay(_retryCount, retryAfter); + _pausedUntil = RetryQueuePolicy.AddDelay(_utcNow(), delay); + PostHogLogger.Warning($"Replay flush failed, retrying in {delay.TotalSeconds:0.###}s"); + } + + void ResetRetryState() + { + _retryCount = 0; + _pausedUntil = null; + } + + IEnumerator SendBatch(List events, Action onComplete) { var url = $"{_host}/s/"; @@ -338,18 +382,22 @@ IEnumerator SendBatch(List events, Action onComplete) yield return request.SendWebRequest(); int statusCode = (int)request.responseCode; + var retryAfter = NetworkClient.ParseRetryAfter( + request.GetResponseHeader("Retry-After"), + DateTimeOffset.UtcNow + ); if (request.result == UnityWebRequest.Result.Success) { PostHogLogger.Debug($"Replay batch sent successfully (status: {statusCode})"); - onComplete?.Invoke(true, statusCode); + onComplete?.Invoke(new BatchUploadResult(true, statusCode, retryAfter)); } else { PostHogLogger.Warning( $"Replay batch send failed: {request.error} (status: {statusCode})" ); - onComplete?.Invoke(false, statusCode); + onComplete?.Invoke(new BatchUploadResult(false, statusCode, retryAfter)); } } @@ -380,27 +428,6 @@ byte[] CompressGzip(byte[] data) } return output.ToArray(); } - - bool ShouldDeleteEventsOnError(int statusCode) - { - // Retry on network errors (0) and redirects (3xx) - if (statusCode == 0 || (statusCode >= 300 && statusCode < 400)) - return false; - - // Don't retry on client errors (4xx) except 413 - if (statusCode >= 400 && statusCode < 500 && statusCode != 413) - return true; - - // Retry on 413 (handled by reducing batch size) - if (statusCode == 413) - return false; - - // Retry on server errors (5xx) - if (statusCode >= 500) - return false; - - return false; - } } /// diff --git a/tests/PostHog.Unity.Tests/EventQueueTests.cs b/tests/PostHog.Unity.Tests/EventQueueTests.cs new file mode 100644 index 0000000..853dcba --- /dev/null +++ b/tests/PostHog.Unity.Tests/EventQueueTests.cs @@ -0,0 +1,353 @@ +using System.Collections; + +namespace PostHogUnity.Tests +{ + [Collection("UnityGlobals")] + public class EventQueueTests + { + [Fact] + public void ConstructorTrimsLoadedQueueToConfiguredCapacity() + { + var storage = new FakeStorage(); + storage.SaveEvent("oldest", EventJson("oldest-payload", "oldest")); + storage.SaveEvent("middle", EventJson("middle-payload", "middle")); + storage.SaveEvent("newest", EventJson("newest-payload", "newest")); + + _ = CreateHarness(storage, maxQueueSize: 2); + + Assert.Equal(new[] { "middle", "newest" }, storage.GetEventIds()); + } + + [Fact] + public void EnqueueTrimsAllExcessAfterCapacityIsLowered() + { + var storage = new FakeStorage(); + for (var i = 0; i < 5; i++) + { + storage.SaveEvent($"legacy-{i}", EventJson($"payload-{i}", $"event-{i}")); + } + var harness = CreateHarness(storage, maxQueueSize: 5); + harness.Config.MaxQueueSize = 2; + + harness.Queue.Enqueue(Event("replacement")); + + Assert.Equal(2, storage.GetEventCount()); + Assert.DoesNotContain("legacy-0", storage.GetEventIds()); + Assert.DoesNotContain("legacy-3", storage.GetEventIds()); + Assert.Contains("legacy-4", storage.GetEventIds()); + } + + [Fact] + public void DuplicatePayloadUuidsUseDistinctQueueOwnedIds() + { + var harness = CreateHarness(); + var first = Event("first"); + var second = Event("second"); + first.Uuid = "shared-payload-uuid"; + second.Uuid = "shared-payload-uuid"; + + harness.Queue.Enqueue(first); + harness.Queue.Enqueue(second); + + var ids = harness.Storage.GetEventIds(); + Assert.Equal(2, ids.Count); + Assert.NotEqual(ids[0], ids[1]); + Assert.DoesNotContain("shared-payload-uuid", ids); + Assert.All(ids, id => Assert.Equal('7', id.Split('-')[2][0])); + Assert.All( + ids, + id => Assert.Equal("shared-payload-uuid", ReadUuid(harness.Storage.LoadEvent(id))) + ); + } + + [Fact] + public void LegacyStorageIdIsLoadedAndAcknowledged() + { + var storage = new FakeStorage(); + storage.SaveEvent( + "legacy-payload-uuid", + EventJson("legacy-payload-uuid", "legacy-event") + ); + var harness = CreateHarness(storage); + harness.Results.Enqueue(new BatchUploadResult(true, 200)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Empty(storage.GetEventIds()); + Assert.Equal("legacy-payload-uuid", Assert.Single(harness.Attempts).Single().Uuid); + } + + [Fact] + public void CorruptLeadingEntryDoesNotShiftValidEntryAcknowledgement() + { + var storage = new FakeStorage(); + storage.SaveEvent("corrupt-storage-id", "not-json"); + storage.SaveEvent("valid-storage-id", EventJson("valid-payload-uuid", "valid-event")); + var harness = CreateHarness(storage); + harness.Results.Enqueue(new BatchUploadResult(true, 200)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + var sentEvent = Assert.Single(Assert.Single(harness.Attempts)); + Assert.Equal("valid-event", sentEvent.Event); + Assert.Equal("valid-payload-uuid", sentEvent.Uuid); + Assert.Equal( + new[] { "corrupt-storage-id", "valid-storage-id" }, + storage.DeletedEventIds + ); + Assert.Empty(storage.GetEventIds()); + } + + [Theory] + [InlineData(0)] + [InlineData(408)] + [InlineData(429)] + [InlineData(500)] + [InlineData(503)] + public void RetryableFailuresRetainEntries(int statusCode) + { + var harness = CreateHarness(); + harness.Queue.Enqueue(Event("retained")); + var id = Assert.Single(harness.Storage.GetEventIds()); + harness.Results.Enqueue(new BatchUploadResult(false, statusCode)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(id, Assert.Single(harness.Storage.GetEventIds())); + Assert.Single(harness.Attempts); + } + + [Fact] + public void TerminalClientFailureDeletesOnlySentEntries() + { + var harness = CreateHarness(); + harness.Queue.Enqueue(Event("terminal")); + harness.Results.Enqueue(new BatchUploadResult(false, 400)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Empty(harness.Storage.GetEventIds()); + } + + [Fact] + public void RetryAfterLongerThanLocalBackoffDelaysNextAttempt() + { + var harness = CreateHarness(); + harness.Queue.Enqueue(Event("retry-after")); + harness.Results.Enqueue(new BatchUploadResult(false, 503, TimeSpan.FromSeconds(60))); + harness.Results.Enqueue(new BatchUploadResult(true, 200)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + harness.Now = harness.Now.AddSeconds(30); + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Single(harness.Attempts); + Assert.Single(harness.Storage.GetEventIds()); + + harness.Now = harness.Now.AddSeconds(30); + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(2, harness.Attempts.Count); + Assert.Empty(harness.Storage.GetEventIds()); + } + + [Fact] + public void PayloadTooLargeShrinksToSingletonThenDropsOnlyPoisonEntry() + { + var harness = CreateHarness(maxQueueSize: 10); + for (var i = 0; i < 4; i++) + { + harness.Queue.Enqueue(Event($"event-{i}")); + } + var originalIds = new List(harness.Storage.GetEventIds()); + harness.Results.Enqueue(new BatchUploadResult(false, 413)); + harness.Results.Enqueue(new BatchUploadResult(false, 413)); + harness.Results.Enqueue(new BatchUploadResult(false, 413)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + harness.Now = harness.Now.AddSeconds(5); + RunCoroutine(harness.Queue.FlushCoroutine()); + harness.Now = harness.Now.AddSeconds(10); + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(new[] { 4, 2, 1 }, harness.Attempts.Select(batch => batch.Count)); + Assert.Equal(3, harness.Storage.GetEventCount()); + Assert.DoesNotContain(originalIds[0], harness.Storage.GetEventIds()); + Assert.Equal(originalIds.Skip(1), harness.Storage.GetEventIds()); + } + + [Fact] + public void SuccessfulInFlightBatchDoesNotAcknowledgeCapacityReplacements() + { + var harness = CreateHarness(maxQueueSize: 3); + harness.Queue.Enqueue(Event("a")); + harness.Queue.Enqueue(Event("b")); + harness.Queue.Enqueue(Event("c")); + var replaced = false; + harness.OnSend = _ => + { + if (replaced) + return; + replaced = true; + harness.Queue.Enqueue(Event("x")); + harness.Queue.Enqueue(Event("y")); + harness.Queue.Enqueue(Event("z")); + }; + harness.Results.Enqueue(new BatchUploadResult(true, 200)); + harness.Results.Enqueue(new BatchUploadResult(false, 503)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(2, harness.Attempts.Count); + Assert.Equal(new[] { "x", "y", "z" }, harness.Attempts[1].Select(evt => evt.Event)); + Assert.Equal(3, harness.Storage.GetEventCount()); + } + + [Fact] + public void KnownOfflineFlushDoesNotAttemptTransportOrConsumeRetryState() + { + var harness = CreateHarness(); + harness.Queue.Enqueue(Event("offline")); + harness.Connected = false; + + RunCoroutine(harness.Queue.FlushCoroutine()); + Assert.Empty(harness.Attempts); + + harness.Connected = true; + harness.Results.Enqueue(new BatchUploadResult(true, 200)); + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Single(harness.Attempts); + Assert.Empty(harness.Storage.GetEventIds()); + } + + static QueueHarness CreateHarness(FakeStorage storage = null, int maxQueueSize = 100) + { + return new QueueHarness(storage ?? new FakeStorage(), maxQueueSize); + } + + static PostHogEvent Event(string name) + { + return new PostHogEvent(name, "distinct-id"); + } + + static string EventJson(string uuid, string name) + { + var evt = Event(name); + evt.Uuid = uuid; + return JsonSerializer.SerializeEvent(evt); + } + + static string ReadUuid(string json) + { + return JsonSerializer.DeserializeDictionary(json)["uuid"].ToString(); + } + + static void RunCoroutine(IEnumerator coroutine) + { + while (coroutine.MoveNext()) + { + if (coroutine.Current is IEnumerator nestedCoroutine) + { + RunCoroutine(nestedCoroutine); + } + } + } + + sealed class QueueHarness + { + public QueueHarness(FakeStorage storage, int maxQueueSize) + { + Storage = storage; + Config = new PostHogConfig + { + ApiKey = "test-api-key", + MaxBatchSize = 10, + MaxQueueSize = maxQueueSize, + FlushAt = 10, + }; + Queue = new EventQueue(Config, Storage, SendBatch, () => Now, () => Connected); + } + + public PostHogConfig Config { get; } + public FakeStorage Storage { get; } + public EventQueue Queue { get; } + public Queue Results { get; } = new(); + public List> Attempts { get; } = new(); + public DateTime Now { get; set; } = new(2026, 1, 1, 0, 0, 0, DateTimeKind.Utc); + public bool Connected { get; set; } = true; + public Action> OnSend { get; set; } + + IEnumerator SendBatch(BatchPayload payload, Action onComplete) + { + Attempts.Add(new List(payload.Batch)); + OnSend?.Invoke(payload.Batch); + onComplete(Results.Dequeue()); + yield break; + } + } + + sealed class FakeStorage : IStorageProvider + { + readonly List _ids = new(); + readonly Dictionary _events = new(); + readonly Dictionary _state = new(); + + public List DeletedEventIds { get; } = new(); + + public void Initialize(string basePath) { } + + public void SaveEvent(string id, string jsonData) + { + if (!_events.ContainsKey(id)) + { + _ids.Add(id); + } + _events[id] = jsonData; + } + + public string LoadEvent(string id) + { + return _events.TryGetValue(id, out var value) ? value : null; + } + + public void DeleteEvent(string id) + { + DeletedEventIds.Add(id); + _events.Remove(id); + _ids.Remove(id); + } + + public IReadOnlyList GetEventIds() + { + return _ids.AsReadOnly(); + } + + public int GetEventCount() + { + return _ids.Count; + } + + public void Clear() + { + _events.Clear(); + _ids.Clear(); + } + + public void SaveState(string key, string jsonData) + { + _state[key] = jsonData; + } + + public string LoadState(string key) + { + return _state.TryGetValue(key, out var value) ? value : null; + } + + public void DeleteState(string key) + { + _state.Remove(key); + } + } + } +} diff --git a/tests/PostHog.Unity.Tests/NetworkClientTests.cs b/tests/PostHog.Unity.Tests/NetworkClientTests.cs index 1494cba..c283f57 100644 --- a/tests/PostHog.Unity.Tests/NetworkClientTests.cs +++ b/tests/PostHog.Unity.Tests/NetworkClientTests.cs @@ -382,5 +382,153 @@ public object Send() public void Dispose() { } } } + + [Collection("UnityGlobals")] + public class TheBatchRetryAfterHandling + { + static readonly DateTimeOffset Now = new(2026, 1, 1, 0, 0, 0, TimeSpan.Zero); + + [Fact] + public void ParsesDeltaSeconds() + { + Assert.Equal(TimeSpan.FromSeconds(60), NetworkClient.ParseRetryAfter("60", Now)); + } + + [Fact] + public void ParsesHttpDate() + { + Assert.Equal( + TimeSpan.FromSeconds(60), + NetworkClient.ParseRetryAfter("Thu, 01 Jan 2026 00:01:00 GMT", Now) + ); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData("not-a-delay")] + [InlineData("-1")] + public void IgnoresMissingOrInvalidValues(string value) + { + Assert.Null(NetworkClient.ParseRetryAfter(value, Now)); + } + + [Fact] + public void PropagatesDeltaSecondsFromProtocolErrorResponse() + { + var (result, request) = SendProtocolError("60"); + + Assert.False(result.Success); + Assert.Equal(503, result.StatusCode); + Assert.Equal(TimeSpan.FromSeconds(60), result.RetryAfter); + Assert.True(request.WasSent); + Assert.Equal("Retry-After", request.RequestedHeader); + } + + [Fact] + public void PropagatesHttpDateFromProtocolErrorResponse() + { + var retryAt = DateTimeOffset.UtcNow.AddMinutes(2); + + var (result, _) = SendProtocolError(retryAt.ToString("R")); + + Assert.NotNull(result.RetryAfter); + Assert.InRange(result.RetryAfter.Value.TotalSeconds, 118, 120); + } + + [Theory] + [InlineData(null)] + [InlineData("")] + [InlineData("not-a-delay")] + public void DoesNotPropagateMissingOrInvalidHeader(string header) + { + var (result, _) = SendProtocolError(header); + + Assert.Null(result.RetryAfter); + } + + static (BatchUploadResult Result, FakeBatchRequest Request) SendProtocolError( + string retryAfter + ) + { + var request = new FakeBatchRequest( + UnityWebRequest.Result.ProtocolError, + 503, + "Service Unavailable", + retryAfter + ); + var client = new NetworkClient( + new PostHogConfig { ApiKey = "test-api-key", Host = "https://example.com" }, + (_, _, _, _, _, _, _) => throw new InvalidOperationException(), + _ => EmptyCoroutine(), + (_, _) => request + ); + BatchUploadResult result = null; + + RunCoroutine( + client.SendBatch( + new BatchPayload("test-api-key", new List()), + response => result = response + ) + ); + + Assert.NotNull(result); + return (result, request); + } + + static void RunCoroutine(IEnumerator coroutine) + { + while (coroutine.MoveNext()) + { + if (coroutine.Current is IEnumerator nestedCoroutine) + { + RunCoroutine(nestedCoroutine); + } + } + } + + static IEnumerator EmptyCoroutine() + { + yield break; + } + + sealed class FakeBatchRequest : NetworkClient.IBatchRequest + { + readonly string _retryAfter; + + public FakeBatchRequest( + UnityWebRequest.Result result, + long responseCode, + string error, + string retryAfter + ) + { + Result = result; + ResponseCode = responseCode; + Error = error; + _retryAfter = retryAfter; + } + + public UnityWebRequest.Result Result { get; } + public long ResponseCode { get; } + public string Error { get; } + public bool WasSent { get; private set; } + public string RequestedHeader { get; private set; } + + public object Send() + { + WasSent = true; + return EmptyCoroutine(); + } + + public string GetResponseHeader(string name) + { + RequestedHeader = name; + return _retryAfter; + } + + public void Dispose() { } + } + } } } diff --git a/tests/PostHog.Unity.Tests/ReplayQueueTests.cs b/tests/PostHog.Unity.Tests/ReplayQueueTests.cs index e88b374..7c7567b 100644 --- a/tests/PostHog.Unity.Tests/ReplayQueueTests.cs +++ b/tests/PostHog.Unity.Tests/ReplayQueueTests.cs @@ -1,8 +1,10 @@ +using System.Collections; using System.Reflection; using PostHogUnity.SessionReplay; namespace PostHogUnity.Tests { + [Collection("UnityGlobals")] public class ReplayQueueTests { [Fact] @@ -18,14 +20,7 @@ public void EnqueueCreatesReplayEnvelopeWithUtcTimestamp() queue.Enqueue(new List { RREvent.CreateMeta(100, 200, "Home", 123L) }); - var queueField = typeof(ReplayQueue).GetField( - "_queue", - BindingFlags.Instance | BindingFlags.NonPublic - ); - Assert.NotNull(queueField); - var events = Assert.IsType>(queueField.GetValue(queue)); - var envelope = Assert.Single(events); - + var envelope = Assert.Single(GetQueuedEvents(queue)); Assert.Matches(@"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}\.\d{7}Z$", envelope.Timestamp); Assert.Equal(123L, Assert.Single(envelope.SnapshotData).Timestamp); } @@ -62,5 +57,170 @@ public void PreparePayloadFallsBackToUncompressedWhenCompressionFails() Assert.Same(bodyBytes, payloadBytes); Assert.False(useCompression); } + + [Theory] + [InlineData(408)] + [InlineData(429)] + [InlineData(500)] + [InlineData(503)] + public void RetryableFailuresRetainReplayEvents(int statusCode) + { + var harness = new ReplayHarness(); + harness.Enqueue("retained"); + var uuid = Assert.Single(GetQueuedEvents(harness.Queue)).Uuid; + harness.Results.Enqueue(new BatchUploadResult(false, statusCode)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(uuid, Assert.Single(GetQueuedEvents(harness.Queue)).Uuid); + Assert.Single(harness.Attempts); + } + + [Fact] + public void RetryAfterLongerThanLocalBackoffDelaysReplayRetry() + { + var harness = new ReplayHarness(); + harness.Enqueue("retry-after"); + harness.Results.Enqueue(new BatchUploadResult(false, 503, TimeSpan.FromSeconds(60))); + harness.Results.Enqueue(new BatchUploadResult(true, 200)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + harness.Now = harness.Now.AddSeconds(30); + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Single(harness.Attempts); + Assert.Equal(1, harness.Queue.Count); + + harness.Now = harness.Now.AddSeconds(30); + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(2, harness.Attempts.Count); + Assert.Equal(0, harness.Queue.Count); + } + + [Fact] + public void SuccessfulInFlightBatchDoesNotAcknowledgeReplayReplacements() + { + var harness = new ReplayHarness(maxQueueSize: 3); + harness.Enqueue("a"); + harness.Enqueue("b"); + harness.Enqueue("c"); + var replaced = false; + harness.OnSend = _ => + { + if (replaced) + return; + replaced = true; + harness.Enqueue("x"); + harness.Enqueue("y"); + harness.Enqueue("z"); + }; + harness.Results.Enqueue(new BatchUploadResult(true, 200)); + harness.Results.Enqueue(new BatchUploadResult(false, 503)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(2, harness.Attempts.Count); + Assert.Equal(3, harness.Queue.Count); + Assert.Equal( + harness.Attempts[1].Select(evt => evt.Uuid), + GetQueuedEvents(harness.Queue).Select(evt => evt.Uuid) + ); + } + + [Fact] + public void PayloadTooLargeShrinksReplayBatchThenDropsOnlySingleton() + { + var harness = new ReplayHarness(); + for (var i = 0; i < 4; i++) + { + harness.Enqueue($"event-{i}"); + } + harness.Results.Enqueue(new BatchUploadResult(false, 413)); + harness.Results.Enqueue(new BatchUploadResult(false, 413)); + harness.Results.Enqueue(new BatchUploadResult(false, 413)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + harness.Now = harness.Now.AddSeconds(5); + RunCoroutine(harness.Queue.FlushCoroutine()); + harness.Now = harness.Now.AddSeconds(10); + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(new[] { 4, 2, 1 }, harness.Attempts.Select(batch => batch.Count)); + var poisonUuid = harness.Attempts[2][0].Uuid; + Assert.Equal(3, harness.Queue.Count); + Assert.DoesNotContain(GetQueuedEvents(harness.Queue), evt => evt.Uuid == poisonUuid); + } + + [Fact] + public void KnownOfflineReplayFlushDoesNotAttemptTransport() + { + var harness = new ReplayHarness(); + harness.Enqueue("offline"); + harness.Connected = false; + + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Empty(harness.Attempts); + Assert.Equal(1, harness.Queue.Count); + } + + static List GetQueuedEvents(ReplayQueue queue) + { + var queueField = typeof(ReplayQueue).GetField( + "_queue", + BindingFlags.Instance | BindingFlags.NonPublic + ); + Assert.NotNull(queueField); + return Assert.IsType>(queueField.GetValue(queue)); + } + + static void RunCoroutine(IEnumerator coroutine) + { + while (coroutine.MoveNext()) + { + if (coroutine.Current is IEnumerator nestedCoroutine) + { + RunCoroutine(nestedCoroutine); + } + } + } + + sealed class ReplayHarness + { + public ReplayHarness(int maxQueueSize = 100) + { + Queue = new ReplayQueue( + new PostHogSessionReplayConfig { FlushAt = 20, MaxQueueSize = maxQueueSize }, + "test-api-key", + "https://example.com", + () => "distinct-id", + () => "session-id", + SendBatch, + () => Now, + () => Connected + ); + } + + public ReplayQueue Queue { get; } + public Queue Results { get; } = new(); + public List> Attempts { get; } = new(); + public DateTime Now { get; set; } = new(2026, 1, 1, 0, 0, 0, DateTimeKind.Utc); + public bool Connected { get; set; } = true; + public Action> OnSend { get; set; } + + public void Enqueue(string screenName) + { + Queue.Enqueue(new List { RREvent.CreateMeta(100, 200, screenName, 123L) }); + } + + IEnumerator SendBatch(List batch, Action onComplete) + { + Attempts.Add(new List(batch)); + OnSend?.Invoke(batch); + onComplete(Results.Dequeue()); + yield break; + } + } } } From d5498e04432055da929947efbab7db045c14be1b Mon Sep 17 00:00:00 2001 From: Dustin Byrne Date: Tue, 1 Sep 2026 13:08:32 -0400 Subject: [PATCH 2/2] refactor: defer retry-after handling to capture v1 --- .changeset/durable-retry-queue.md | 2 +- com.posthog.unity/Runtime/Core/EventQueue.cs | 39 +++-- .../Runtime/Core/NetworkClient.cs | 151 ++---------------- .../Runtime/SessionReplay/ReplayQueue.cs | 40 ++--- tests/PostHog.Unity.Tests/EventQueueTests.cs | 49 ++---- .../PostHog.Unity.Tests/NetworkClientTests.cs | 148 ----------------- tests/PostHog.Unity.Tests/ReplayQueueTests.cs | 41 ++--- 7 files changed, 85 insertions(+), 385 deletions(-) diff --git a/.changeset/durable-retry-queue.md b/.changeset/durable-retry-queue.md index 3935299..e6e39a4 100644 --- a/.changeset/durable-retry-queue.md +++ b/.changeset/durable-retry-queue.md @@ -2,4 +2,4 @@ 'com.posthog.unity': patch --- -Preserve analytics and session replay events across transient ingestion failures, honor server retry delays, safely handle oversized payloads, and prevent in-flight queue replacements from being acknowledged as sent. +Preserve analytics and session replay events across transient ingestion failures, safely handle oversized payloads, and prevent in-flight queue replacements from being acknowledged as sent. diff --git a/com.posthog.unity/Runtime/Core/EventQueue.cs b/com.posthog.unity/Runtime/Core/EventQueue.cs index b24b51b..b074db3 100644 --- a/com.posthog.unity/Runtime/Core/EventQueue.cs +++ b/com.posthog.unity/Runtime/Core/EventQueue.cs @@ -12,7 +12,7 @@ class EventQueue { readonly PostHogConfig _config; readonly IStorageProvider _storage; - readonly Func, IEnumerator> _sendBatch; + readonly Func, IEnumerator> _sendBatch; readonly Func _utcNow; readonly Func _isConnected; readonly object _lock = new(); @@ -45,7 +45,7 @@ NetworkClient networkClient internal EventQueue( PostHogConfig config, IStorageProvider storage, - Func, IEnumerator> sendBatch, + Func, IEnumerator> sendBatch, Func utcNow, Func isConnected ) @@ -312,11 +312,19 @@ public IEnumerator FlushCoroutine() PostHogLogger.Debug($"Flushing batch of {events.Count} events"); var payload = new BatchPayload(_config.ApiKey, events); - var uploadResult = new BatchUploadResult(false, 0); + var success = false; + var statusCode = 0; - yield return _sendBatch(payload, result => uploadResult = result); + yield return _sendBatch( + payload, + (result, code) => + { + success = result; + statusCode = code; + } + ); - if (uploadResult.Success) + if (success) { DeleteEntries(entries); ResetRetryState(); @@ -324,7 +332,7 @@ public IEnumerator FlushCoroutine() continue; } - if (uploadResult.StatusCode == 413) + if (statusCode == 413) { if (entries.Count == 1) { @@ -343,17 +351,17 @@ public IEnumerator FlushCoroutine() PostHogLogger.Warning( $"Payload too large, reducing batch size to {_adjustedMaxBatchSize}" ); - PauseForRetry(uploadResult.RetryAfter, "Flush failed"); + PauseForRetry("Flush failed"); } } - else if (RetryQueuePolicy.ShouldDelete(uploadResult.StatusCode)) + else if (RetryQueuePolicy.ShouldDelete(statusCode)) { DeleteEntries(entries); ResetRetryState(); } else { - PauseForRetry(uploadResult.RetryAfter, "Flush failed"); + PauseForRetry("Flush failed"); } break; } @@ -364,13 +372,13 @@ public IEnumerator FlushCoroutine() } } - void PauseForRetry(TimeSpan? retryAfter, string message) + void PauseForRetry(string message) { if (_retryCount < int.MaxValue) { _retryCount++; } - var delay = RetryQueuePolicy.GetRetryDelay(_retryCount, retryAfter); + var delay = RetryQueuePolicy.GetRetryDelay(_retryCount); _pausedUntil = RetryQueuePolicy.AddDelay(_utcNow(), delay); PostHogLogger.Warning($"{message}, retrying in {delay.TotalSeconds:0.###}s"); } @@ -501,13 +509,10 @@ public static int ReducedBatchSize(int failedBatchSize) return Math.Max(1, failedBatchSize / 2); } - public static TimeSpan GetRetryDelay(int retryCount, TimeSpan? retryAfter) + public static TimeSpan GetRetryDelay(int retryCount) { - var localSeconds = Math.Min((long)retryCount * RetryDelaySeconds, MaxRetryDelaySeconds); - var localDelay = TimeSpan.FromSeconds(localSeconds); - return retryAfter.HasValue && retryAfter.Value > localDelay - ? retryAfter.Value - : localDelay; + var seconds = Math.Min((long)retryCount * RetryDelaySeconds, MaxRetryDelaySeconds); + return TimeSpan.FromSeconds(seconds); } public static DateTime AddDelay(DateTime now, TimeSpan delay) diff --git a/com.posthog.unity/Runtime/Core/NetworkClient.cs b/com.posthog.unity/Runtime/Core/NetworkClient.cs index 6d136dc..1292350 100644 --- a/com.posthog.unity/Runtime/Core/NetworkClient.cs +++ b/com.posthog.unity/Runtime/Core/NetworkClient.cs @@ -1,7 +1,6 @@ using System; using System.Collections; using System.Collections.Generic; -using System.Globalization; using System.Text; using UnityEngine; using UnityEngine.Networking; @@ -16,7 +15,6 @@ class NetworkClient readonly PostHogConfig _config; readonly FeatureFlagsRequestFactory _featureFlagsRequestFactory; readonly FeatureFlagsRetryDelayFactory _featureFlagsRetryDelayFactory; - readonly BatchRequestFactory _batchRequestFactory; const int TimeoutSeconds = 10; const float FeatureFlagsInitialRetryDelaySeconds = 0.3f; @@ -32,17 +30,6 @@ Dictionary> groupProperties internal delegate object FeatureFlagsRetryDelayFactory(int failedAttempt); - internal delegate IBatchRequest BatchRequestFactory(string url, byte[] body); - - internal interface IBatchRequest : IDisposable - { - UnityWebRequest.Result Result { get; } - long ResponseCode { get; } - string Error { get; } - object Send(); - string GetResponseHeader(string name); - } - internal interface IFeatureFlagsRequest : IDisposable { string Url { get; } @@ -57,67 +44,52 @@ public NetworkClient(PostHogConfig config) : this( config, CreateFeatureFlagsRequest, - failedAttempt => new WaitForSeconds( - GetFeatureFlagsRetryDelaySeconds(failedAttempt) - ), - CreateBatchRequest + failedAttempt => new WaitForSeconds(GetFeatureFlagsRetryDelaySeconds(failedAttempt)) ) { } internal NetworkClient( PostHogConfig config, FeatureFlagsRequestFactory featureFlagsRequestFactory, FeatureFlagsRetryDelayFactory featureFlagsRetryDelayFactory - ) - : this( - config, - featureFlagsRequestFactory, - featureFlagsRetryDelayFactory, - CreateBatchRequest - ) { } - - internal NetworkClient( - PostHogConfig config, - FeatureFlagsRequestFactory featureFlagsRequestFactory, - FeatureFlagsRetryDelayFactory featureFlagsRetryDelayFactory, - BatchRequestFactory batchRequestFactory ) { _config = config; _featureFlagsRequestFactory = featureFlagsRequestFactory; _featureFlagsRetryDelayFactory = featureFlagsRetryDelayFactory; - _batchRequestFactory = batchRequestFactory; } /// /// Sends a batch of events to the PostHog API. /// /// The batch payload to send - /// Callback with the upload result and retry metadata. - public IEnumerator SendBatch(BatchPayload payload, Action onComplete) + /// Callback with (success, statusCode) + public IEnumerator SendBatch(BatchPayload payload, Action onComplete) { var url = GetBatchUrl(); var json = JsonSerializer.SerializeBatch(payload); PostHogLogger.Debug($"Sending batch to {url}"); - using var request = _batchRequestFactory(url, Encoding.UTF8.GetBytes(json)); - yield return request.Send(); + using var request = new UnityWebRequest(url, "POST"); + var bodyRaw = Encoding.UTF8.GetBytes(json); + request.uploadHandler = new UploadHandlerRaw(bodyRaw); + request.downloadHandler = new DownloadHandlerBuffer(); + ApplyDefaultHeaders(request); + request.timeout = TimeoutSeconds; + + yield return request.SendWebRequest(); - int statusCode = (int)request.ResponseCode; - var retryAfter = ParseRetryAfter( - request.GetResponseHeader("Retry-After"), - DateTimeOffset.UtcNow - ); + int statusCode = (int)request.responseCode; - if (request.Result == UnityWebRequest.Result.Success) + if (request.result == UnityWebRequest.Result.Success) { PostHogLogger.Debug($"Batch sent successfully (status: {statusCode})"); - onComplete?.Invoke(new BatchUploadResult(true, statusCode, retryAfter)); + onComplete?.Invoke(true, statusCode); } else { - PostHogLogger.Warning($"Batch send failed: {request.Error} (status: {statusCode})"); - onComplete?.Invoke(new BatchUploadResult(false, statusCode, retryAfter)); + PostHogLogger.Warning($"Batch send failed: {request.error} (status: {statusCode})"); + onComplete?.Invoke(false, statusCode); } } @@ -196,44 +168,6 @@ string GetBatchUrl() return $"{host}/batch"; } - internal static TimeSpan? ParseRetryAfter(string value, DateTimeOffset now) - { - if (string.IsNullOrWhiteSpace(value)) - { - return null; - } - - if ( - long.TryParse( - value.Trim(), - NumberStyles.None, - CultureInfo.InvariantCulture, - out var seconds - ) - && seconds >= 0 - ) - { - var maxSeconds = (long)TimeSpan.MaxValue.TotalSeconds; - return TimeSpan.FromSeconds(Math.Min(seconds, maxSeconds)); - } - - if ( - DateTimeOffset.TryParse( - value, - CultureInfo.InvariantCulture, - DateTimeStyles.AllowWhiteSpaces - | DateTimeStyles.AssumeUniversal - | DateTimeStyles.AdjustToUniversal, - out var retryAt - ) - ) - { - return retryAt <= now ? TimeSpan.Zero : retryAt - now; - } - - return null; - } - internal static bool ShouldRetryFeatureFlagsRequest( UnityWebRequest.Result result, int statusCode, @@ -276,45 +210,6 @@ internal static float GetFeatureFlagsRetryDelaySeconds(int failedAttempt) return FeatureFlagsInitialRetryDelaySeconds * (1 << (failedAttempt - 1)); } - static IBatchRequest CreateBatchRequest(string url, byte[] body) - { - var request = new UnityWebRequest(url, "POST"); - request.uploadHandler = new UploadHandlerRaw(body); - request.downloadHandler = new DownloadHandlerBuffer(); - ApplyDefaultHeaders(request); - request.timeout = TimeoutSeconds; - return new UnityWebRequestBatchRequest(request); - } - - sealed class UnityWebRequestBatchRequest : IBatchRequest - { - readonly UnityWebRequest _request; - - public UnityWebRequestBatchRequest(UnityWebRequest request) - { - _request = request; - } - - public UnityWebRequest.Result Result => _request.result; - public long ResponseCode => _request.responseCode; - public string Error => _request.error; - - public object Send() - { - return _request.SendWebRequest(); - } - - public string GetResponseHeader(string name) - { - return _request.GetResponseHeader(name); - } - - public void Dispose() - { - _request.Dispose(); - } - } - static IFeatureFlagsRequest CreateFeatureFlagsRequest( string apiKey, string host, @@ -437,18 +332,4 @@ static void ApplyDefaultHeaders(UnityWebRequest request) request.SetRequestHeader("User-Agent", SdkInfo.UserAgent); } } - - class BatchUploadResult - { - public BatchUploadResult(bool success, int statusCode, TimeSpan? retryAfter = null) - { - Success = success; - StatusCode = statusCode; - RetryAfter = retryAfter; - } - - public bool Success { get; } - public int StatusCode { get; } - public TimeSpan? RetryAfter { get; } - } } diff --git a/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs b/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs index a231847..cb9868f 100644 --- a/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs +++ b/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs @@ -20,7 +20,7 @@ class ReplayQueue readonly string _host; readonly Func _getDistinctId; readonly Func _getSessionId; - readonly Func, Action, IEnumerator> _sendBatch; + readonly Func, Action, IEnumerator> _sendBatch; readonly Func _utcNow; readonly Func _isConnected; readonly object _lock = new(); @@ -62,7 +62,7 @@ internal ReplayQueue( string host, Func getDistinctId, Func getSessionId, - Func, Action, IEnumerator> sendBatch, + Func, Action, IEnumerator> sendBatch, Func utcNow, Func isConnected ) @@ -272,10 +272,18 @@ internal IEnumerator FlushCoroutine() PostHogLogger.Debug($"Flushing batch of {batch.Count} replay events"); - var uploadResult = new BatchUploadResult(false, 0); - yield return _sendBatch(batch, result => uploadResult = result); + var success = false; + var statusCode = 0; + yield return _sendBatch( + batch, + (result, code) => + { + success = result; + statusCode = code; + } + ); - if (uploadResult.Success) + if (success) { RemoveBatch(batch); ResetRetryState(); @@ -283,7 +291,7 @@ internal IEnumerator FlushCoroutine() continue; } - if (uploadResult.StatusCode == 413) + if (statusCode == 413) { if (batch.Count == 1) { @@ -299,17 +307,17 @@ internal IEnumerator FlushCoroutine() PostHogLogger.Warning( $"Replay payload too large, reducing batch size to {_adjustedMaxBatchSize}" ); - PauseForRetry(uploadResult.RetryAfter); + PauseForRetry(); } } - else if (RetryQueuePolicy.ShouldDelete(uploadResult.StatusCode)) + else if (RetryQueuePolicy.ShouldDelete(statusCode)) { RemoveBatch(batch); ResetRetryState(); } else { - PauseForRetry(uploadResult.RetryAfter); + PauseForRetry(); } break; } @@ -331,13 +339,13 @@ void RemoveBatch(List batch) } } - void PauseForRetry(TimeSpan? retryAfter) + void PauseForRetry() { if (_retryCount < int.MaxValue) { _retryCount++; } - var delay = RetryQueuePolicy.GetRetryDelay(_retryCount, retryAfter); + var delay = RetryQueuePolicy.GetRetryDelay(_retryCount); _pausedUntil = RetryQueuePolicy.AddDelay(_utcNow(), delay); PostHogLogger.Warning($"Replay flush failed, retrying in {delay.TotalSeconds:0.###}s"); } @@ -348,7 +356,7 @@ void ResetRetryState() _pausedUntil = null; } - IEnumerator SendBatch(List events, Action onComplete) + IEnumerator SendBatch(List events, Action onComplete) { var url = $"{_host}/s/"; @@ -382,22 +390,18 @@ IEnumerator SendBatch(List events, Action onCo yield return request.SendWebRequest(); int statusCode = (int)request.responseCode; - var retryAfter = NetworkClient.ParseRetryAfter( - request.GetResponseHeader("Retry-After"), - DateTimeOffset.UtcNow - ); if (request.result == UnityWebRequest.Result.Success) { PostHogLogger.Debug($"Replay batch sent successfully (status: {statusCode})"); - onComplete?.Invoke(new BatchUploadResult(true, statusCode, retryAfter)); + onComplete?.Invoke(true, statusCode); } else { PostHogLogger.Warning( $"Replay batch send failed: {request.error} (status: {statusCode})" ); - onComplete?.Invoke(new BatchUploadResult(false, statusCode, retryAfter)); + onComplete?.Invoke(false, statusCode); } } diff --git a/tests/PostHog.Unity.Tests/EventQueueTests.cs b/tests/PostHog.Unity.Tests/EventQueueTests.cs index 853dcba..21b1ca3 100644 --- a/tests/PostHog.Unity.Tests/EventQueueTests.cs +++ b/tests/PostHog.Unity.Tests/EventQueueTests.cs @@ -69,7 +69,7 @@ public void LegacyStorageIdIsLoadedAndAcknowledged() EventJson("legacy-payload-uuid", "legacy-event") ); var harness = CreateHarness(storage); - harness.Results.Enqueue(new BatchUploadResult(true, 200)); + harness.Results.Enqueue((true, 200)); RunCoroutine(harness.Queue.FlushCoroutine()); @@ -84,7 +84,7 @@ public void CorruptLeadingEntryDoesNotShiftValidEntryAcknowledgement() storage.SaveEvent("corrupt-storage-id", "not-json"); storage.SaveEvent("valid-storage-id", EventJson("valid-payload-uuid", "valid-event")); var harness = CreateHarness(storage); - harness.Results.Enqueue(new BatchUploadResult(true, 200)); + harness.Results.Enqueue((true, 200)); RunCoroutine(harness.Queue.FlushCoroutine()); @@ -109,7 +109,7 @@ public void RetryableFailuresRetainEntries(int statusCode) var harness = CreateHarness(); harness.Queue.Enqueue(Event("retained")); var id = Assert.Single(harness.Storage.GetEventIds()); - harness.Results.Enqueue(new BatchUploadResult(false, statusCode)); + harness.Results.Enqueue((false, statusCode)); RunCoroutine(harness.Queue.FlushCoroutine()); @@ -122,35 +122,13 @@ public void TerminalClientFailureDeletesOnlySentEntries() { var harness = CreateHarness(); harness.Queue.Enqueue(Event("terminal")); - harness.Results.Enqueue(new BatchUploadResult(false, 400)); + harness.Results.Enqueue((false, 400)); RunCoroutine(harness.Queue.FlushCoroutine()); Assert.Empty(harness.Storage.GetEventIds()); } - [Fact] - public void RetryAfterLongerThanLocalBackoffDelaysNextAttempt() - { - var harness = CreateHarness(); - harness.Queue.Enqueue(Event("retry-after")); - harness.Results.Enqueue(new BatchUploadResult(false, 503, TimeSpan.FromSeconds(60))); - harness.Results.Enqueue(new BatchUploadResult(true, 200)); - - RunCoroutine(harness.Queue.FlushCoroutine()); - harness.Now = harness.Now.AddSeconds(30); - RunCoroutine(harness.Queue.FlushCoroutine()); - - Assert.Single(harness.Attempts); - Assert.Single(harness.Storage.GetEventIds()); - - harness.Now = harness.Now.AddSeconds(30); - RunCoroutine(harness.Queue.FlushCoroutine()); - - Assert.Equal(2, harness.Attempts.Count); - Assert.Empty(harness.Storage.GetEventIds()); - } - [Fact] public void PayloadTooLargeShrinksToSingletonThenDropsOnlyPoisonEntry() { @@ -160,9 +138,9 @@ public void PayloadTooLargeShrinksToSingletonThenDropsOnlyPoisonEntry() harness.Queue.Enqueue(Event($"event-{i}")); } var originalIds = new List(harness.Storage.GetEventIds()); - harness.Results.Enqueue(new BatchUploadResult(false, 413)); - harness.Results.Enqueue(new BatchUploadResult(false, 413)); - harness.Results.Enqueue(new BatchUploadResult(false, 413)); + harness.Results.Enqueue((false, 413)); + harness.Results.Enqueue((false, 413)); + harness.Results.Enqueue((false, 413)); RunCoroutine(harness.Queue.FlushCoroutine()); harness.Now = harness.Now.AddSeconds(5); @@ -193,8 +171,8 @@ public void SuccessfulInFlightBatchDoesNotAcknowledgeCapacityReplacements() harness.Queue.Enqueue(Event("y")); harness.Queue.Enqueue(Event("z")); }; - harness.Results.Enqueue(new BatchUploadResult(true, 200)); - harness.Results.Enqueue(new BatchUploadResult(false, 503)); + harness.Results.Enqueue((true, 200)); + harness.Results.Enqueue((false, 503)); RunCoroutine(harness.Queue.FlushCoroutine()); @@ -214,7 +192,7 @@ public void KnownOfflineFlushDoesNotAttemptTransportOrConsumeRetryState() Assert.Empty(harness.Attempts); harness.Connected = true; - harness.Results.Enqueue(new BatchUploadResult(true, 200)); + harness.Results.Enqueue((true, 200)); RunCoroutine(harness.Queue.FlushCoroutine()); Assert.Single(harness.Attempts); @@ -272,17 +250,18 @@ public QueueHarness(FakeStorage storage, int maxQueueSize) public PostHogConfig Config { get; } public FakeStorage Storage { get; } public EventQueue Queue { get; } - public Queue Results { get; } = new(); + public Queue<(bool Success, int StatusCode)> Results { get; } = new(); public List> Attempts { get; } = new(); public DateTime Now { get; set; } = new(2026, 1, 1, 0, 0, 0, DateTimeKind.Utc); public bool Connected { get; set; } = true; public Action> OnSend { get; set; } - IEnumerator SendBatch(BatchPayload payload, Action onComplete) + IEnumerator SendBatch(BatchPayload payload, Action onComplete) { Attempts.Add(new List(payload.Batch)); OnSend?.Invoke(payload.Batch); - onComplete(Results.Dequeue()); + var result = Results.Dequeue(); + onComplete(result.Success, result.StatusCode); yield break; } } diff --git a/tests/PostHog.Unity.Tests/NetworkClientTests.cs b/tests/PostHog.Unity.Tests/NetworkClientTests.cs index c283f57..1494cba 100644 --- a/tests/PostHog.Unity.Tests/NetworkClientTests.cs +++ b/tests/PostHog.Unity.Tests/NetworkClientTests.cs @@ -382,153 +382,5 @@ public object Send() public void Dispose() { } } } - - [Collection("UnityGlobals")] - public class TheBatchRetryAfterHandling - { - static readonly DateTimeOffset Now = new(2026, 1, 1, 0, 0, 0, TimeSpan.Zero); - - [Fact] - public void ParsesDeltaSeconds() - { - Assert.Equal(TimeSpan.FromSeconds(60), NetworkClient.ParseRetryAfter("60", Now)); - } - - [Fact] - public void ParsesHttpDate() - { - Assert.Equal( - TimeSpan.FromSeconds(60), - NetworkClient.ParseRetryAfter("Thu, 01 Jan 2026 00:01:00 GMT", Now) - ); - } - - [Theory] - [InlineData(null)] - [InlineData("")] - [InlineData("not-a-delay")] - [InlineData("-1")] - public void IgnoresMissingOrInvalidValues(string value) - { - Assert.Null(NetworkClient.ParseRetryAfter(value, Now)); - } - - [Fact] - public void PropagatesDeltaSecondsFromProtocolErrorResponse() - { - var (result, request) = SendProtocolError("60"); - - Assert.False(result.Success); - Assert.Equal(503, result.StatusCode); - Assert.Equal(TimeSpan.FromSeconds(60), result.RetryAfter); - Assert.True(request.WasSent); - Assert.Equal("Retry-After", request.RequestedHeader); - } - - [Fact] - public void PropagatesHttpDateFromProtocolErrorResponse() - { - var retryAt = DateTimeOffset.UtcNow.AddMinutes(2); - - var (result, _) = SendProtocolError(retryAt.ToString("R")); - - Assert.NotNull(result.RetryAfter); - Assert.InRange(result.RetryAfter.Value.TotalSeconds, 118, 120); - } - - [Theory] - [InlineData(null)] - [InlineData("")] - [InlineData("not-a-delay")] - public void DoesNotPropagateMissingOrInvalidHeader(string header) - { - var (result, _) = SendProtocolError(header); - - Assert.Null(result.RetryAfter); - } - - static (BatchUploadResult Result, FakeBatchRequest Request) SendProtocolError( - string retryAfter - ) - { - var request = new FakeBatchRequest( - UnityWebRequest.Result.ProtocolError, - 503, - "Service Unavailable", - retryAfter - ); - var client = new NetworkClient( - new PostHogConfig { ApiKey = "test-api-key", Host = "https://example.com" }, - (_, _, _, _, _, _, _) => throw new InvalidOperationException(), - _ => EmptyCoroutine(), - (_, _) => request - ); - BatchUploadResult result = null; - - RunCoroutine( - client.SendBatch( - new BatchPayload("test-api-key", new List()), - response => result = response - ) - ); - - Assert.NotNull(result); - return (result, request); - } - - static void RunCoroutine(IEnumerator coroutine) - { - while (coroutine.MoveNext()) - { - if (coroutine.Current is IEnumerator nestedCoroutine) - { - RunCoroutine(nestedCoroutine); - } - } - } - - static IEnumerator EmptyCoroutine() - { - yield break; - } - - sealed class FakeBatchRequest : NetworkClient.IBatchRequest - { - readonly string _retryAfter; - - public FakeBatchRequest( - UnityWebRequest.Result result, - long responseCode, - string error, - string retryAfter - ) - { - Result = result; - ResponseCode = responseCode; - Error = error; - _retryAfter = retryAfter; - } - - public UnityWebRequest.Result Result { get; } - public long ResponseCode { get; } - public string Error { get; } - public bool WasSent { get; private set; } - public string RequestedHeader { get; private set; } - - public object Send() - { - WasSent = true; - return EmptyCoroutine(); - } - - public string GetResponseHeader(string name) - { - RequestedHeader = name; - return _retryAfter; - } - - public void Dispose() { } - } - } } } diff --git a/tests/PostHog.Unity.Tests/ReplayQueueTests.cs b/tests/PostHog.Unity.Tests/ReplayQueueTests.cs index 7c7567b..3b0aff2 100644 --- a/tests/PostHog.Unity.Tests/ReplayQueueTests.cs +++ b/tests/PostHog.Unity.Tests/ReplayQueueTests.cs @@ -68,7 +68,7 @@ public void RetryableFailuresRetainReplayEvents(int statusCode) var harness = new ReplayHarness(); harness.Enqueue("retained"); var uuid = Assert.Single(GetQueuedEvents(harness.Queue)).Uuid; - harness.Results.Enqueue(new BatchUploadResult(false, statusCode)); + harness.Results.Enqueue((false, statusCode)); RunCoroutine(harness.Queue.FlushCoroutine()); @@ -76,28 +76,6 @@ public void RetryableFailuresRetainReplayEvents(int statusCode) Assert.Single(harness.Attempts); } - [Fact] - public void RetryAfterLongerThanLocalBackoffDelaysReplayRetry() - { - var harness = new ReplayHarness(); - harness.Enqueue("retry-after"); - harness.Results.Enqueue(new BatchUploadResult(false, 503, TimeSpan.FromSeconds(60))); - harness.Results.Enqueue(new BatchUploadResult(true, 200)); - - RunCoroutine(harness.Queue.FlushCoroutine()); - harness.Now = harness.Now.AddSeconds(30); - RunCoroutine(harness.Queue.FlushCoroutine()); - - Assert.Single(harness.Attempts); - Assert.Equal(1, harness.Queue.Count); - - harness.Now = harness.Now.AddSeconds(30); - RunCoroutine(harness.Queue.FlushCoroutine()); - - Assert.Equal(2, harness.Attempts.Count); - Assert.Equal(0, harness.Queue.Count); - } - [Fact] public void SuccessfulInFlightBatchDoesNotAcknowledgeReplayReplacements() { @@ -115,8 +93,8 @@ public void SuccessfulInFlightBatchDoesNotAcknowledgeReplayReplacements() harness.Enqueue("y"); harness.Enqueue("z"); }; - harness.Results.Enqueue(new BatchUploadResult(true, 200)); - harness.Results.Enqueue(new BatchUploadResult(false, 503)); + harness.Results.Enqueue((true, 200)); + harness.Results.Enqueue((false, 503)); RunCoroutine(harness.Queue.FlushCoroutine()); @@ -136,9 +114,9 @@ public void PayloadTooLargeShrinksReplayBatchThenDropsOnlySingleton() { harness.Enqueue($"event-{i}"); } - harness.Results.Enqueue(new BatchUploadResult(false, 413)); - harness.Results.Enqueue(new BatchUploadResult(false, 413)); - harness.Results.Enqueue(new BatchUploadResult(false, 413)); + harness.Results.Enqueue((false, 413)); + harness.Results.Enqueue((false, 413)); + harness.Results.Enqueue((false, 413)); RunCoroutine(harness.Queue.FlushCoroutine()); harness.Now = harness.Now.AddSeconds(5); @@ -203,7 +181,7 @@ public ReplayHarness(int maxQueueSize = 100) } public ReplayQueue Queue { get; } - public Queue Results { get; } = new(); + public Queue<(bool Success, int StatusCode)> Results { get; } = new(); public List> Attempts { get; } = new(); public DateTime Now { get; set; } = new(2026, 1, 1, 0, 0, 0, DateTimeKind.Utc); public bool Connected { get; set; } = true; @@ -214,11 +192,12 @@ public void Enqueue(string screenName) Queue.Enqueue(new List { RREvent.CreateMeta(100, 200, screenName, 123L) }); } - IEnumerator SendBatch(List batch, Action onComplete) + IEnumerator SendBatch(List batch, Action onComplete) { Attempts.Add(new List(batch)); OnSend?.Invoke(batch); - onComplete(Results.Dequeue()); + var result = Results.Dequeue(); + onComplete(result.Success, result.StatusCode); yield break; } }