using WxAgent.Service; using Xunit; namespace WxAgent.Service.Tests; public sealed class EventHubTests { [Fact] public async Task ReplayAndGapAreExplicitAndSubscriptionsAreOwnerScoped() { var hub = new EventHub(); var first = new AgentEvent("one", "account", "chat", "message", "summary", DateTimeOffset.UtcNow); hub.Publish("alice", first); var replay = hub.Subscribe("alice", "one"); Assert.Empty(replay.Replay); Assert.False(replay.Gap); var gap = hub.Subscribe("alice", "gone"); Assert.True(gap.Gap); Assert.False(hub.TryGet(replay.SubscriptionId, "bob", out _)); hub.Publish("alice", new AgentEvent("two", "account", "chat", "message", "summary", DateTimeOffset.UtcNow)); Assert.True(hub.TryGet(replay.SubscriptionId, "alice", out var subscription)); Assert.True(await subscription.Channel.Reader.WaitToReadAsync()); } [Fact] public void SlowConsumerGetsAnExplicitGapInsteadOfSilentDrop() { var hub = new EventHub(); var subscription = hub.Subscribe("alice", null); for (var i = 0; i < 102; i++) hub.Publish("alice", new AgentEvent(i.ToString(), "account", "chat", "message", "summary", DateTimeOffset.UtcNow)); Assert.True(hub.TryGet(subscription.SubscriptionId, "alice", out var current)); var foundGap = false; while (current.Channel.Reader.TryRead(out var item)) foundGap |= item.Kind == "gap"; Assert.True(foundGap); } }