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 ReplayNeverCrossesPrincipalBoundary() { var hub = new EventHub(); hub.Publish("alice", new AgentEvent("alice-1", "account-a", "chat-a", "message", "summary", DateTimeOffset.UtcNow)); hub.Publish("bob", new AgentEvent("bob-1", "account-b", "chat-b", "message", "summary", DateTimeOffset.UtcNow)); var alice = hub.Subscribe("alice", "alice-1"); var bob = hub.Subscribe("bob", "bob-1"); Assert.Empty(alice.Replay); Assert.Empty(bob.Replay); Assert.False(alice.Gap); Assert.False(bob.Gap); var aliceAfterBob = hub.Subscribe("alice", "bob-1"); Assert.Empty(aliceAfterBob.Replay); Assert.True(aliceAfterBob.Gap); } [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); } }