diff --git a/.changeset/durable-retry-queue.md b/.changeset/durable-retry-queue.md new file mode 100644 index 0000000..e6e39a4 --- /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, 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..b074db3 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,80 @@ 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 success = false; + var statusCode = 0; - yield return _networkClient.SendBatch( + yield return _sendBatch( payload, - (result, statusCode) => + (result, code) => { success = result; - shouldDeleteEvents = ShouldDeleteEventsOnError(statusCode); + statusCode = code; } ); if (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 (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 + ); + _adjustedFlushAt = Math.Min(_adjustedFlushAt, _adjustedMaxBatchSize); + PostHogLogger.Warning( + $"Payload too large, reducing batch size to {_adjustedMaxBatchSize}" ); - _pausedUntil = DateTime.UtcNow.AddSeconds(delay); - PostHogLogger.Warning($"Flush failed, retrying in {delay}s"); + PauseForRetry("Flush failed"); } - break; } + else if (RetryQueuePolicy.ShouldDelete(statusCode)) + { + DeleteEntries(entries); + ResetRetryState(); + } + else + { + PauseForRetry("Flush failed"); + } + break; } } finally @@ -294,27 +372,55 @@ public IEnumerator FlushCoroutine() } } - List LoadEvents(List eventIds) + void PauseForRetry(string message) { - var events = new List(); + if (_retryCount < int.MaxValue) + { + _retryCount++; + } + var delay = RetryQueuePolicy.GetRetryDelay(_retryCount); + _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 +469,56 @@ 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) + { + var seconds = Math.Min((long)retryCount * RetryDelaySeconds, MaxRetryDelaySeconds); + return TimeSpan.FromSeconds(seconds); + } + + 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/SessionReplay/ReplayQueue.cs b/com.posthog.unity/Runtime/SessionReplay/ReplayQueue.cs index 4533ae7..cb9868f 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,10 +272,9 @@ IEnumerator FlushCoroutine() PostHogLogger.Debug($"Flushing batch of {batch.Count} replay events"); - bool success = false; - int statusCode = 0; - - yield return SendBatch( + var success = false; + var statusCode = 0; + yield return _sendBatch( batch, (result, code) => { @@ -258,44 +285,41 @@ IEnumerator FlushCoroutine() if (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 (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(); } - break; } + else if (RetryQueuePolicy.ShouldDelete(statusCode)) + { + RemoveBatch(batch); + ResetRetryState(); + } + else + { + PauseForRetry(); + } + break; } } finally @@ -304,6 +328,34 @@ IEnumerator FlushCoroutine() } } + void RemoveBatch(List batch) + { + lock (_lock) + { + foreach (var sentEvent in batch) + { + _queue.Remove(sentEvent); + } + } + } + + void PauseForRetry() + { + if (_retryCount < int.MaxValue) + { + _retryCount++; + } + var delay = RetryQueuePolicy.GetRetryDelay(_retryCount); + _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/"; @@ -380,27 +432,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..21b1ca3 --- /dev/null +++ b/tests/PostHog.Unity.Tests/EventQueueTests.cs @@ -0,0 +1,332 @@ +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((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((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((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((false, 400)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + 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((false, 413)); + harness.Results.Enqueue((false, 413)); + harness.Results.Enqueue((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((true, 200)); + harness.Results.Enqueue((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((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<(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) + { + Attempts.Add(new List(payload.Batch)); + OnSend?.Invoke(payload.Batch); + var result = Results.Dequeue(); + onComplete(result.Success, result.StatusCode); + 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/ReplayQueueTests.cs b/tests/PostHog.Unity.Tests/ReplayQueueTests.cs index e88b374..3b0aff2 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,149 @@ 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((false, statusCode)); + + RunCoroutine(harness.Queue.FlushCoroutine()); + + Assert.Equal(uuid, Assert.Single(GetQueuedEvents(harness.Queue)).Uuid); + Assert.Single(harness.Attempts); + } + + [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((true, 200)); + harness.Results.Enqueue((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((false, 413)); + harness.Results.Enqueue((false, 413)); + harness.Results.Enqueue((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<(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; } + + 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); + var result = Results.Dequeue(); + onComplete(result.Success, result.StatusCode); + yield break; + } + } } }