package streamgate import ( "context" "fmt" "reflect" "strings" "testing" "time" ) // now is a deterministic timestamp shared across tests. var now = time.Date(2026, 7, 24, 12, 0, 0, 0, time.UTC) // --------------------------------------------------------------------------- // API-1: FilterHoldRequirement, EvidencePlan, EvidenceTail core state // --------------------------------------------------------------------------- // TestFilterHoldRequirementValidatesModeMatrix verifies the field matrix for // every FilterHoldMode. Each mode has required, forbidden, and optional // fields, plus hard rune bounds. The test covers the constructor surface and // the Validate path. func TestFilterHoldRequirementValidatesModeMatrix(t *testing.T) { textKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} toolKinds := []EventKind{EventKindToolCallFragment} allKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta, EventKindToolCallFragment} // --- none mode ----------------------------------------------------------- t.Run("none", func(t *testing.T) { // Valid: channel + subscribed kinds, zeroed rune/trigger/buffer. r, err := NewFilterHoldRequirementNone("ch-none", textKinds) if err != nil { t.Fatalf("unexpected error: %v", err) } if err := r.Validate(); err != nil { t.Fatalf("validate failed: %v", err) } if r.EvidenceRunes() != 0 { t.Errorf("evidenceRunes want 0, got %d", r.EvidenceRunes()) } if r.TriggerKind() != "" { t.Errorf("trigger want empty, got %q", r.TriggerKind()) } if r.MaxBufferRunes() != 0 { t.Errorf("maxBuffer want 0, got %d", r.MaxBufferRunes()) } if r.IsBlocking() { t.Errorf("none must not be blocking") } // Empty channel rejected. _, err = NewFilterHoldRequirementNone("", textKinds) if err == nil { t.Error("empty channel should be rejected") } // Empty subscribed kinds rejected. _, err = NewFilterHoldRequirementNone("ch", nil) if err == nil { t.Error("nil kinds should be rejected") } }) // --- rolling_window mode ------------------------------------------------- t.Run("rolling_window", func(t *testing.T) { // Valid defaults. r, err := NewFilterHoldRequirementRolling("ch-roll", textKinds, 500) if err != nil { t.Fatalf("unexpected error: %v", err) } if err := r.Validate(); err != nil { t.Fatalf("validate failed: %v", err) } if r.EvidenceRunes() != 500 { t.Errorf("evidenceRunes want 500, got %d", r.EvidenceRunes()) } if r.TriggerKind() != "" { t.Errorf("trigger must be empty for rolling, got %q", r.TriggerKind()) } if r.MaxBufferRunes() == 0 { t.Error("rolling must have a positive max_buffer_runes") } if !r.IsBlocking() { t.Error("rolling must be blocking") } // Below minimum evidence runes. _, err = NewFilterHoldRequirementRolling("ch", textKinds, 0) if err == nil { t.Error("evidenceRunes=0 should be rejected") } // Above maximum evidence runes. _, err = NewFilterHoldRequirementRolling("ch", textKinds, maxEvidenceRunes+1) if err == nil { t.Error("evidenceRunes>max should be rejected") } // Trigger kind is forbidden on rolling. bad := FilterHoldRequirement{ channel: "ch", mode: FilterHoldModeRolling, subscribedKinds: textKinds, evidenceRunes: 100, triggerKind: EventKindTerminal, maxBufferRunes: defaultMaxBufferRunes, } if err := bad.Validate(); err == nil { t.Error("rolling with trigger should be rejected") } // maxBufferRunes below minimum. bad2 := FilterHoldRequirement{ channel: "ch", mode: FilterHoldModeRolling, subscribedKinds: textKinds, evidenceRunes: 100, maxBufferRunes: minMaxBufferRunes - 1, } if err := bad2.Validate(); err == nil { t.Error("maxBuffer below minimum should be rejected") } }) // --- terminal_gate mode -------------------------------------------------- t.Run("terminal_gate", func(t *testing.T) { r, err := NewFilterHoldRequirementTerminalGate("ch-term", allKinds, EventKindTerminal) if err != nil { t.Fatalf("unexpected error: %v", err) } if err := r.Validate(); err != nil { t.Fatalf("validate failed: %v", err) } if r.TriggerKind() != EventKindTerminal { t.Errorf("trigger want terminal, got %q", r.TriggerKind()) } if r.EvidenceRunes() != 0 { t.Error("terminal_gate must not have evidenceRunes") } if !r.IsBlocking() { t.Error("terminal_gate must be blocking") } // Provider error also accepted as trigger. r2, err := NewFilterHoldRequirementTerminalGate("ch", allKinds, EventKindProviderError) if err != nil { t.Fatalf("provider_error trigger should be accepted: %v", err) } if err := r2.Validate(); err != nil { t.Fatalf("validate: %v", err) } // Non-terminal kind rejected. _, err = NewFilterHoldRequirementTerminalGate("ch", allKinds, EventKindTextDelta) if err == nil { t.Error("text_delta trigger should be rejected for terminal_gate") } // evidenceRunes is forbidden. bad := FilterHoldRequirement{ channel: "ch", mode: FilterHoldModeTerminalGate, subscribedKinds: allKinds, evidenceRunes: 100, triggerKind: EventKindTerminal, maxBufferRunes: defaultMaxBufferRunes, } if err := bad.Validate(); err == nil { t.Error("terminal_gate with evidenceRunes should be rejected") } }) // --- fragment_gate mode -------------------------------------------------- t.Run("fragment_gate", func(t *testing.T) { r, err := NewFilterHoldRequirementFragmentGate("ch-frag", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("unexpected error: %v", err) } if err := r.Validate(); err != nil { t.Fatalf("validate failed: %v", err) } if r.TriggerKind() != EventKindToolCallFragment { t.Errorf("trigger want tool_call_fragment, got %q", r.TriggerKind()) } if r.EvidenceRunes() != 0 { t.Error("fragment_gate must not have evidenceRunes") } if !r.IsBlocking() { t.Error("fragment_gate must be blocking") } // Wrong trigger kind. _, err = NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindTextDelta) if err == nil { t.Error("text_delta trigger should be rejected for fragment_gate") } // evidenceRunes is forbidden. bad := FilterHoldRequirement{ channel: "ch", mode: FilterHoldModeFragmentGate, subscribedKinds: toolKinds, evidenceRunes: 10, triggerKind: EventKindToolCallFragment, maxBufferRunes: defaultMaxBufferRunes, } if err := bad.Validate(); err == nil { t.Error("fragment_gate with evidenceRunes should be rejected") } }) // --- duplicate subscribed kinds rejected --------------------------------- t.Run("duplicateKinds", func(t *testing.T) { dupKinds := []EventKind{EventKindTextDelta, EventKindTextDelta} _, err := NewFilterHoldRequirementRolling("ch", dupKinds, 100) if err == nil { t.Error("duplicate kinds should be rejected") } }) // --- unknown mode rejected ----------------------------------------------- t.Run("unknownMode", func(t *testing.T) { bad := FilterHoldRequirement{ channel: "ch", mode: "not_a_mode", subscribedKinds: textKinds, } if err := bad.Validate(); err == nil { t.Error("unknown mode should be rejected") } }) } // TestEvidencePlanComposesOnlyBlockingRequirements verifies that observe-only // and none-mode bindings never create a channel hold, while blocking bindings // are kept and the strongest mode per channel wins. func TestEvidencePlanComposesOnlyBlockingRequirements(t *testing.T) { textKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} toolKinds := []EventKind{EventKindToolCallFragment} t.Run("observeOnlyDoesNotBlock", func(t *testing.T) { // Build a blocking rolling requirement. roll, err := NewFilterHoldRequirementRolling("ch", textKinds, 500) if err != nil { t.Fatalf("build rolling: %v", err) } // Observe-only binding for the same channel. bindObs, err := NewFilterHoldBinding("f-obs", roll, FilterEnforcementObserveOnly) if err != nil { t.Fatalf("build observe binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bindObs}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } if plan.HasBlockingRequirement("ch") { t.Error("observe-only must not create blocking requirement") } if plan.SubscribeKinds("ch") == nil { t.Error("subscribe kinds should still include observe kinds") } }) t.Run("noneModeDoesNotBlock", func(t *testing.T) { none, err := NewFilterHoldRequirementNone("ch", textKinds) if err != nil { t.Fatalf("build none: %v", err) } bind, err := NewFilterHoldBinding("f-none", none, FilterEnforcementObserveOnly) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } if plan.HasBlockingRequirement("ch") { t.Error("none mode must not create blocking requirement") } }) t.Run("multiRequirementPreserved", func(t *testing.T) { // Same channel: rolling + fragment_gate. Both bindings are preserved. roll, err := NewFilterHoldRequirementRolling("ch", textKinds, 500) if err != nil { t.Fatalf("build rolling: %v", err) } frag, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build fragment: %v", err) } b1, err := NewFilterHoldBinding("f1", roll, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", frag, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } if !plan.HasBlockingRequirement("ch") { t.Fatal("expected blocking requirement") } bindings := plan.BindingsForChannel("ch") if len(bindings) != 2 { t.Errorf("bindings count want 2, got %d", len(bindings)) } // Blocking kinds should be the union of both. blockingKinds := plan.BlockingKinds("ch") if len(blockingKinds) != 3 { t.Errorf("blocking kinds want 3 (textDelta, reasoningDelta, toolCallFragment), got %d", len(blockingKinds)) } }) t.Run("multiScheduleSupported", func(t *testing.T) { // fragment_gate (trigger=tool_call_fragment) + terminal_gate // (trigger=terminal) on same channel: multi-schedule is supported. frag, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build fragment: %v", err) } term, err := NewFilterHoldRequirementTerminalGate("ch", textKinds, EventKindTerminal) if err != nil { t.Fatalf("build terminal: %v", err) } b1, err := NewFilterHoldBinding("f1", frag, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", term, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("multi-schedule evidence plan failed: %v", err) } if len(plan.BindingsForChannel("ch")) != 2 { t.Errorf("bindings count want 2, got %d", len(plan.BindingsForChannel("ch"))) } }) t.Run("emptyBindings", func(t *testing.T) { plan, err := NewEvidencePlan(nil) if err != nil { t.Fatalf("empty bindings: %v", err) } if plan.HasBlockingRequirement("anything") { t.Error("empty plan must not have blocking requirements") } }) } // TestEvidenceTailRollingRuneThresholdAndLookBehind verifies rolling threshold // behaviour at the default 500-rune and a policy-override 200-rune window, // plus the bounded cross-boundary look-behind invariant. func TestEvidenceTailRollingRuneThresholdAndLookBehind(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} t.Run("default500RuneThreshold", func(t *testing.T) { req, err := NewFilterHoldRequirementRolling("ch", textKinds, 500) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Feed 499 'a' chars — should not trigger. ev, err := NewTextDeltaEvent("ch", repeatRune('a', 499), now) if err != nil { t.Fatalf("build event: %v", err) } epoch, signal, err := tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } if signal != EvidenceTailSignalNone { t.Errorf("signal want none before threshold, got %s", signal) } if epoch.Triggered() { t.Error("epoch must not be triggered below threshold") } if cs := tail.channelState["ch"]; cs == nil { t.Fatal("channel state should exist") } else if cs.pendingRunes != 499 { t.Errorf("pendingRunes want 499, got %d", cs.pendingRunes) } // 500th rune triggers the threshold. ev2, err := NewTextDeltaEvent("ch", "a", now) if err != nil { t.Fatalf("build event: %v", err) } _, signal, err = tail.Append(ev2) if err != nil { t.Fatalf("append 500th: %v", err) } if signal != EvidenceTailSignalThreshold { t.Errorf("signal want threshold at 500, got %s", signal) } }) t.Run("policy200RuneThreshold", func(t *testing.T) { req, err := NewFilterHoldRequirementRolling("ch", textKinds, 200) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Feed 200 runes in chunks — threshold fires on the 200th. var total int signal := EvidenceTailSignalNone for i := 0; i < 25; i++ { chunk := repeatRune('한', 8) // 8 korean chars = 24 runes ev, err := NewTextDeltaEvent("ch", chunk, now) if err != nil { t.Fatalf("build event: %v", err) } _, signal, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } total += 24 if signal == EvidenceTailSignalThreshold { break } } if signal != EvidenceTailSignalThreshold { t.Errorf("signal want threshold, got %s", signal) } }) t.Run("boundedLookBehind", func(t *testing.T) { req, err := NewFilterHoldRequirementRolling("ch", textKinds, 10) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Fill to threshold. for i := 0; i < 5; i++ { ev, err := NewTextDeltaEvent("ch", repeatRune('b', 2), now) if err != nil { t.Fatalf("build event: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } cs := tail.channelState["ch"] if len(cs.pendingEntries) != 5 { t.Errorf("pending want 5, got %d", len(cs.pendingEntries)) } if cs.pendingRunes != 10 { t.Errorf("pendingRunes want 10, got %d", cs.pendingRunes) } }) } // TestEvidenceTailPreservesKoreanRuneChunks verifies that Korean multi-byte // UTF-8 strings are counted by Unicode runes (not bytes) and that invalid // UTF-8 is rejected without mutating state. func TestEvidenceTailPreservesKoreanRuneChunks(t *testing.T) { allKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} req, err := NewFilterHoldRequirementRolling("ch", allKinds, 3) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } t.Run("koreanChunksCountedAsRunes", func(t *testing.T) { // "한글" = 2 runes, each 3 bytes in UTF-8. ev, err := NewTextDeltaEvent("ch", "한글", now) if err != nil { t.Fatalf("build event: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } cs := tail.channelState["ch"] if cs.pendingRunes != 2 { t.Errorf("pendingRunes want 2 for '한글', got %d", cs.pendingRunes) } // One more Korean char = 3 runes total → threshold. ev2, err := NewTextDeltaEvent("ch", "한", now) if err != nil { t.Fatalf("build event: %v", err) } _, signal, err := tail.Append(ev2) if err != nil { t.Fatalf("append: %v", err) } if signal != EvidenceTailSignalThreshold { t.Errorf("signal want threshold at 3 runes, got %s", signal) } }) t.Run("koreanSplitAcrossChunks", func(t *testing.T) { tail2, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Feed partial Korean characters — each Append validates UTF-8. // Split "한" (UTF-8: E1 8A 90) into two pieces. han := []byte{0xE1, 0x8A, 0x90} partial := string(han[:2]) // Invalid UTF-8 prefix. // Valid chunk first. ev, err := NewTextDeltaEvent("ch", "가", now) // 1 rune, 3 bytes valid UTF-8. if err != nil { t.Fatalf("build event: %v", err) } _, _, err = tail2.Append(ev) if err != nil { t.Fatalf("append valid: %v", err) } // Invalid UTF-8 chunk must be rejected. evBad, err := NewTextDeltaEvent("ch", partial, now) if err != nil { t.Fatalf("build invalid event: %v", err) } _, _, err = tail2.Append(evBad) if err == nil { t.Error("invalid UTF-8 must be rejected with error") } // State must not have been mutated by the invalid append. cs := tail2.channelState["ch"] if cs.pendingRunes != 1 { t.Errorf("pendingRunes must stay 1 after invalid append, got %d", cs.pendingRunes) } }) t.Run("reasoningDeltaAlsoCounted", func(t *testing.T) { tail3, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } ev, err := NewReasoningDeltaEvent("ch", "한글", now) if err != nil { t.Fatalf("build event: %v", err) } _, _, err = tail3.Append(ev) if err != nil { t.Fatalf("append: %v", err) } cs := tail3.channelState["ch"] if cs.pendingRunes != 2 { t.Errorf("reasoning rune count want 2, got %d", cs.pendingRunes) } }) } // TestEvidenceTailKeysFragmentsByToolCall verifies that fragment completion for // one tool-call ID does not accidentally trigger completion for another, and // that interleaved fragments on the same channel are tracked independently. // Fragment state is now accumulated automatically by Append. func TestEvidenceTailKeysFragmentsByToolCall(t *testing.T) { toolKinds := []EventKind{EventKindToolCallFragment} req, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } t.Run("interleavedToolIDsIndependent", func(t *testing.T) { // Tool A and Tool B fragments interleaved. Append automatically // registers fragment state. evA1, err := NewToolCallFragmentEvent("ch", "call-A", "funcA", `{"x":1}`, now) if err != nil { t.Fatalf("build event A1: %v", err) } _, _, err = tail.Append(evA1) if err != nil { t.Fatalf("append A1: %v", err) } evB1, err := NewToolCallFragmentEvent("ch", "call-B", "funcB", `{"y":2}`, now) if err != nil { t.Fatalf("build event B1: %v", err) } _, _, err = tail.Append(evB1) if err != nil { t.Fatalf("append B1: %v", err) } // Complete A — should mark A as completed. epoch, signal, ok, err := tail.CompleteFragment("ch", "call-A") if err != nil { t.Fatalf("complete A: %v", err) } if !ok { t.Error("complete A should return true") } if signal == EvidenceTailSignalNone { t.Error("complete A should produce a trigger signal") } if epoch.ID() == 0 { t.Error("complete A should return an epoch") } // Complete B — independent of A. epochB, signalB, okB, err := tail.CompleteFragment("ch", "call-B") if err != nil { t.Fatalf("complete B: %v", err) } if !okB { t.Error("complete B should return true") } if signalB == EvidenceTailSignalNone { t.Error("complete B should produce a trigger signal") } if epochB.ID() == 0 { t.Error("complete B should return an epoch") } // Re-complete A (already completed) — should return false. _, _, ok2, err := tail.CompleteFragment("ch", "call-A") if err != nil { t.Fatalf("complete A again: %v", err) } if ok2 { t.Error("re-completing completed fragment should return false") } }) t.Run("unknownToolIdNoCompletion", func(t *testing.T) { _, _, ok, err := tail.CompleteFragment("ch", "nonexistent") if err != nil { t.Fatalf("complete unknown: %v", err) } if ok { t.Error("unknown tool id must not complete") } }) } // --------------------------------------------------------------------------- // API-2: Prepared release, partial confirm, overflow, recovery transitions // --------------------------------------------------------------------------- // TestEvidenceTailConfirmsOnlyReleasedPrefix verifies that PrepareRelease does // not mutate pending/cursor state, and ConfirmRelease only moves the number // of events matching ReleasedEvents into look-behind. func TestEvidenceTailConfirmsOnlyReleasedPrefix(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} req, err := NewFilterHoldRequirementRolling("ch", textKinds, 2) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Fill to threshold. var lastEpochID uint64 for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "x", now) if err != nil { t.Fatalf("build event: %v", err) } epoch, _, err := tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } if epoch.Triggered() { lastEpochID = epoch.ID() } } cs := tail.channelState["ch"] beforePending := len(cs.pendingEntries) if lastEpochID == 0 { t.Fatal("expected at least one threshold-triggering epoch") } // PrepareRelease must not mutate pending. pr, err := tail.PrepareRelease(lastEpochID) if err != nil { t.Fatalf("prepare release: %v", err) } if len(cs.pendingEntries) != beforePending { t.Error("PrepareRelease must not change pending count") } if cs.pendingRunes != beforePending { t.Error("PrepareRelease must not change pendingRunes") } releaseEvents := pr.ReleaseEvents() if len(releaseEvents) != 3 { t.Errorf("releaseEvents want 3, got %d", len(releaseEvents)) } t.Run("partialConfirm", func(t *testing.T) { conf := ReleaseConfirmation{ReleasedEvents: 2} if err := tail.ConfirmRelease(pr.Token(), conf); err != nil { t.Fatalf("partial confirm: %v", err) } cs = tail.channelState["ch"] if len(cs.pendingEntries) != 1 { t.Errorf("pending after partial confirm want 1, got %d", len(cs.pendingEntries)) } if cs.pendingRunes != 1 { t.Errorf("pendingRunes after partial confirm want 1, got %d", cs.pendingRunes) } if len(cs.committedLookBehind) != 2 { t.Errorf("lookBehind want 2, got %d", len(cs.committedLookBehind)) } }) // Prepare another release. var lastEpochID2 uint64 for i := 0; i < 1; i++ { ev, err := NewTextDeltaEvent("ch", "y", now) if err != nil { t.Fatalf("build event: %v", err) } epoch, _, err := tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } if epoch.Triggered() { lastEpochID2 = epoch.ID() } } if lastEpochID2 == 0 { t.Fatal("expected threshold-triggering epoch for second prepare") } prFull, err := tail.PrepareRelease(lastEpochID2) if err != nil { t.Fatalf("prepare: %v", err) } fullSnapshotSize := len(prFull.ReleaseEvents()) t.Run("fullConfirm", func(t *testing.T) { conf := ReleaseConfirmation{ReleasedEvents: fullSnapshotSize} if err := tail.ConfirmRelease(prFull.Token(), conf); err != nil { t.Fatalf("full confirm: %v", err) } cs = tail.channelState["ch"] if len(cs.pendingEntries) != 0 { t.Errorf("pending after full confirm want 0, got %d", len(cs.pendingEntries)) } }) t.Run("zeroConfirmNoMutation", func(t *testing.T) { // Re-add events for zero-confirm (all were released by full confirm). for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "z", now) if err != nil { t.Fatalf("build event: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } // Prepare a fresh release for zero-confirm test. cs = tail.channelState["ch"] var prZeroEpoch uint64 for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "z", now) if err != nil { t.Fatalf("build event: %v", err) } epoch, _, err := tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } if epoch.Triggered() { prZeroEpoch = epoch.ID() } } if prZeroEpoch == 0 { t.Fatal("expected threshold-triggering epoch for zero-confirm") } prZero, err := tail.PrepareRelease(prZeroEpoch) if err != nil { t.Fatalf("prepare for zero-confirm: %v", err) } // Zero confirm should not change state. prevLookBehind := len(cs.committedLookBehind) err = tail.ConfirmRelease(prZero.Token(), ReleaseConfirmation{ReleasedEvents: 0}) if err != nil { t.Fatalf("zero confirm: %v", err) } if len(cs.committedLookBehind) != prevLookBehind { t.Error("zero confirm must not change look-behind") } }) } // TestEvidenceTailOverflowSignalsNoRelease verifies that exceeding the hard // buffer limit produces an overflow signal without mutating pending state or // producing a release. func TestEvidenceTailOverflowSignalsNoRelease(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} // Use a small maxBuffer to make the test deterministic. req, err := NewFilterHoldRequirementRolling("ch", textKinds, 2) if err != nil { t.Fatalf("build requirement: %v", err) } // Override maxBufferRunes via a custom requirement to keep the test tight. req = FilterHoldRequirement{ channel: "ch", mode: FilterHoldModeRolling, subscribedKinds: textKinds, evidenceRunes: 2, maxBufferRunes: 10, // Minimum valid hard limit. } if err := req.Validate(); err != nil { t.Fatalf("validate custom req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Fill up to 10 runes (maxBufferRunes limit). for i := 0; i < 10; i++ { ev, err := NewTextDeltaEvent("ch", "x", now) if err != nil { t.Fatalf("build event: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } cs := tail.channelState["ch"] if cs.pendingRunes != 10 { t.Errorf("pendingRunes want 10, got %d", cs.pendingRunes) } // 11th rune → overflow beyond maxBufferRunes=10. ev11, err := NewTextDeltaEvent("ch", "x", now) if err != nil { t.Fatalf("build event: %v", err) } _, signal, err := tail.Append(ev11) if err != nil { t.Fatalf("append overflow: %v", err) } if signal != EvidenceTailSignalBufferOverflow { t.Errorf("signal want bufferOverflow, got %s", signal) } // State must not have changed. cs = tail.channelState["ch"] if cs.pendingRunes != 10 { t.Errorf("overflow must not change pendingRunes, want 10 got %d", cs.pendingRunes) } if len(cs.pendingEntries) != 10 { t.Errorf("overflow must not change pending count, want 10 got %d", len(cs.pendingEntries)) } } // TestEvidenceTailRecoveryTransitionsInvalidatePreparedRelease verifies that // ResetForReplace, PrepareContinuation, and DiscardPendingForTerminal all // invalidate prepared tokens and that their preserve/discard scopes differ. func TestEvidenceTailRecoveryTransitionsInvalidatePreparedRelease(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} req, err := NewFilterHoldRequirementRolling("ch", textKinds, 2) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } // --- ResetForReplace ------------------------------------------------------ t.Run("ResetForReplace", func(t *testing.T) { tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "a", now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } // Find the latest epoch ID for prepare release. var resetEpoch uint64 for _, rec := range tail.epochs { if rec.epochID > resetEpoch { resetEpoch = rec.epochID } } if resetEpoch == 0 { t.Fatal("expected at least one epoch for ResetForReplace") } pr, err := tail.PrepareRelease(resetEpoch) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() tail.ResetForReplace() // Pending must be cleared. if tail.channelState["ch"] != nil { t.Error("ResetForReplace must clear channel state") } // Look-behind must be cleared too. // Confirm with invalidated token must fail. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("confirm with invalidated token must fail") } }) // --- PrepareContinuation -------------------------------------------------- t.Run("PrepareContinuationPreservesLookBehind", func(t *testing.T) { tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Add and confirm 2 events. for i := 0; i < 2; i++ { ev, err := NewTextDeltaEvent("ch", "a", now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } // Find the latest epoch ID for prepare release. var contEpoch uint64 for _, rec := range tail.epochs { if rec.epochID > contEpoch { contEpoch = rec.epochID } } if contEpoch == 0 { t.Fatal("expected at least one epoch for PrepareContinuation") } pr, err := tail.PrepareRelease(contEpoch) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() // Confirm 1 event to create look-behind. if err := tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}); err != nil { t.Fatalf("confirm: %v", err) } // Add another event to have pending. ev, err := NewTextDeltaEvent("ch", "b", now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } tail.PrepareContinuation() cs := tail.channelState["ch"] // Pending must be cleared. if len(cs.pendingEntries) != 0 { t.Errorf("pending must be cleared, got %d", len(cs.pendingEntries)) } if cs.pendingRunes != 0 { t.Errorf("pendingRunes must be 0, got %d", cs.pendingRunes) } // Look-behind must be preserved. if len(cs.committedLookBehind) != 1 { t.Errorf("lookBehind want 1 preserved, got %d", len(cs.committedLookBehind)) } }) // --- DiscardPendingForTerminal -------------------------------------------- t.Run("DiscardPendingForTerminalPreservesLookBehind", func(t *testing.T) { tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "a", now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } cs := tail.channelState["ch"] // Find the latest epoch ID for prepare release. var discEpoch uint64 for _, rec := range tail.epochs { if rec.epochID > discEpoch { discEpoch = rec.epochID } } if discEpoch == 0 { t.Fatal("expected at least one epoch for DiscardPendingForTerminal") } pr, err := tail.PrepareRelease(discEpoch) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() // Confirm 2. if err := tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 2}); err != nil { t.Fatalf("confirm: %v", err) } // Add pending event. ev, err := NewTextDeltaEvent("ch", "c", now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } tail.DiscardPendingForTerminal() cs = tail.channelState["ch"] // Pending cleared. if len(cs.pendingEntries) != 0 { t.Errorf("pending must be cleared, got %d", len(cs.pendingEntries)) } // Look-behind preserved. if len(cs.committedLookBehind) != 2 { t.Errorf("lookBehind want 2, got %d", len(cs.committedLookBehind)) } // Prepared token invalidated. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("confirm with invalidated token must fail") } }) t.Run("staleTokenRejected", func(t *testing.T) { tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "a", now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } // Find the latest epoch ID for prepare release. var staleEpoch uint64 for _, rec := range tail.epochs { if rec.epochID > staleEpoch { staleEpoch = rec.epochID } } if staleEpoch == 0 { t.Fatal("expected at least one epoch for staleTokenRejected") } pr, err := tail.PrepareRelease(staleEpoch) if err != nil { t.Fatalf("prepare: %v", err) } // Reset invalidates the token. tail.ResetForReplace() // Confirm with stale token must fail. err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("stale token confirm must fail") } }) t.Run("duplicateConfirmRejected", func(t *testing.T) { tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "a", now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } // Find the latest epoch ID for prepare release. var dupEpoch uint64 for _, rec := range tail.epochs { if rec.epochID > dupEpoch { dupEpoch = rec.epochID } } if dupEpoch == 0 { t.Fatal("expected at least one epoch for duplicateConfirmRejected") } pr, err := tail.PrepareRelease(dupEpoch) if err != nil { t.Fatalf("prepare: %v", err) } // First confirm succeeds. if err := tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 3}); err != nil { t.Fatalf("first confirm: %v", err) } // Duplicate confirm must fail. err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 0}) if err == nil { t.Error("duplicate confirm must fail") } }) } // --------------------------------------------------------------------------- // REVIEW_API-1: Deterministic composition, defensive copy, fail-closed validation // --------------------------------------------------------------------------- // TestFilterHoldRequirementValidatesPublicBoundMatrix verifies that the // constructor with explicit max_buffer_runes override validates the public // bound matrix: default, override, min, max, and reachable threshold. func TestFilterHoldRequirementValidatesPublicBoundMatrix(t *testing.T) { textKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} t.Run("defaultRolling", func(t *testing.T) { _, err := NewFilterHoldRequirementRolling("ch", textKinds, 500) if err != nil { t.Fatalf("default rolling: %v", err) } }) t.Run("overrideWithMaxBuffer", func(t *testing.T) { r, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 200, 2048) if err != nil { t.Fatalf("override rolling: %v", err) } if r.MaxBufferRunes() != 2048 { t.Errorf("buffer want 2048, got %d", r.MaxBufferRunes()) } if r.EvidenceRunes() != 200 { t.Errorf("evidence want 200, got %d", r.EvidenceRunes()) } }) t.Run("minBoundary", func(t *testing.T) { _, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, minEvidenceRunes, minMaxBufferRunes) if err != nil { t.Fatalf("min boundary: %v", err) } }) t.Run("maxBoundary", func(t *testing.T) { _, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, maxEvidenceRunes, maxMaxBufferRunes) if err != nil { t.Fatalf("max boundary: %v", err) } }) t.Run("thresholdEqualToBuffer", func(t *testing.T) { // evidence_runes == max_buffer_runes is allowed (reachable threshold). _, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 100, 100) if err != nil { t.Fatalf("threshold==buffer: %v", err) } }) t.Run("bufferLessThanThresholdRejected", func(t *testing.T) { _, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 200, 100) if err == nil { t.Error("buffer < threshold should be rejected") } }) t.Run("terminalGateOverride", func(t *testing.T) { _, err := NewFilterHoldRequirementTerminalGateWithMaxBuffer("ch", textKinds, EventKindTerminal, 8192) if err != nil { t.Fatalf("terminal gate override: %v", err) } }) t.Run("fragmentGateOverride", func(t *testing.T) { fragKinds := []EventKind{EventKindToolCallFragment} _, err := NewFilterHoldRequirementFragmentGateWithMaxBuffer("ch", fragKinds, EventKindToolCallFragment, 4096) if err != nil { t.Fatalf("fragment gate override: %v", err) } }) } // TestEvidencePlanComposesDeterministicallyAndDefensively verifies that // binding order reversal produces identical results and that returned slices // are defensive copies that cannot mutate the plan. func TestEvidencePlanComposesDeterministicallyAndDefensively(t *testing.T) { textKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} t.Run("bindingOrderReversal", func(t *testing.T) { rollA, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 300, 5000) if err != nil { t.Fatalf("build rollA: %v", err) } rollB, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 400, 3000) if err != nil { t.Fatalf("build rollB: %v", err) } b1, err := NewFilterHoldBinding("f1", rollA, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", rollB, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } planForward, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("forward plan: %v", err) } planReverse, err := NewEvidencePlan([]FilterHoldBinding{b2, b1}) if err != nil { t.Fatalf("reverse plan: %v", err) } reqF, _ := planForward.BlockingRequirement("ch") reqR, _ := planReverse.BlockingRequirement("ch") if reqF.EvidenceRunes() != reqR.EvidenceRunes() { t.Errorf("evidenceRunes mismatch: forward=%d reverse=%d", reqF.EvidenceRunes(), reqR.EvidenceRunes()) } if reqF.MaxBufferRunes() != reqR.MaxBufferRunes() { t.Errorf("maxBufferRunes mismatch: forward=%d reverse=%d", reqF.MaxBufferRunes(), reqR.MaxBufferRunes()) } }) t.Run("returnedSliceMutationIsolation", func(t *testing.T) { roll, err := NewFilterHoldRequirementRolling("ch", textKinds, 500) if err != nil { t.Fatalf("build: %v", err) } bind, err := NewFilterHoldBinding("f", roll, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } kinds := plan.SubscribeKinds("ch") originalLen := len(kinds) originalFirst := kinds[0] // Mutate the returned slice. kinds[0] = "" // Second call should return a fresh copy. kinds2 := plan.SubscribeKinds("ch") if len(kinds2) != originalLen { t.Error("returned slice length changed") } if kinds2[0] != originalFirst { t.Errorf("returned slice was mutated, got %q want %q", kinds2[0], originalFirst) } }) t.Run("incompatibleBoundsRejected", func(t *testing.T) { // Two rolling requirements on same channel where merged buffer < merged threshold. r1, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 500, 600) if err != nil { t.Fatalf("build r1: %v", err) } r2, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 700, 800) if err != nil { t.Fatalf("build r2: %v", err) } b1, err := NewFilterHoldBinding("f1", r1, FilterEnforcementBlocking) if err != nil { t.Fatalf("build b1: %v", err) } b2, err := NewFilterHoldBinding("f2", r2, FilterEnforcementBlocking) if err != nil { t.Fatalf("build b2: %v", err) } _, err = NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err == nil { t.Error("merged buffer (600) < merged threshold (700) should be rejected") } }) } // TestEvidenceTailAppliesSubscriptionsAfterFailClosedValidation verifies that // unsubscribed events do not mutate hold state and that invalid UTF-8 is // rejected before any subscription check (fail-closed). func TestEvidenceTailAppliesSubscriptionsAfterFailClosedValidation(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} req, err := NewFilterHoldRequirementRolling("ch", textKinds, 10) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } t.Run("unsubscribedEventNoStateMutation", func(t *testing.T) { // reasoning_delta is not subscribed. ev, err := NewReasoningDeltaEvent("ch", "hello world!!!", now) if err != nil { t.Fatalf("build event: %v", err) } _, signal, err := tail.Append(ev) if err != nil { t.Fatalf("append unsubscribed: %v", err) } if signal != EvidenceTailSignalNone { t.Errorf("signal want none, got %s", signal) } // No channel state should be created. if tail.channelState["ch"] != nil { t.Error("unsubscribed event must not create channel state") } }) t.Run("observeOnlyInvalidUTF8Rejected", func(t *testing.T) { // Build a plan with both text_delta (subscribed) and reasoning_delta (observe-only). allKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} req2, err := NewFilterHoldRequirementRolling("ch2", allKinds, 10) if err != nil { t.Fatalf("build req2: %v", err) } bind2, err := NewFilterHoldBinding("f2", req2, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind2: %v", err) } plan2, err := NewEvidencePlan([]FilterHoldBinding{bind2}) if err != nil { t.Fatalf("build plan2: %v", err) } tail2, err := NewEvidenceTail(plan2) if err != nil { t.Fatalf("build tail2: %v", err) } // Invalid UTF-8 for a subscribed kind must be rejected (fail-closed). invalidUTF8 := string([]byte{0xFF, 0xFE}) ev, err := NewTextDeltaEvent("ch2", invalidUTF8, now) if err != nil { t.Fatalf("build invalid event: %v", err) } _, _, err = tail2.Append(ev) if err == nil { t.Error("invalid UTF-8 must be rejected") } // State must not have been mutated. cs := tail2.channelState["ch2"] if cs != nil && cs.pendingRunes != 0 { t.Errorf("state must not be mutated, pendingRunes=%d", cs.pendingRunes) } }) t.Run("noHoldInvalidUTF8Rejected", func(t *testing.T) { // None mode with no blocking requirement. noneKinds := []EventKind{EventKindTextDelta} noneReq, err := NewFilterHoldRequirementNone("ch3", noneKinds) if err != nil { t.Fatalf("build none: %v", err) } noneBind, err := NewFilterHoldBinding("f3", noneReq, FilterEnforcementBlocking) if err != nil { t.Fatalf("build none bind: %v", err) } nonePlan, err := NewEvidencePlan([]FilterHoldBinding{noneBind}) if err != nil { t.Fatalf("build none plan: %v", err) } noneTail, err := NewEvidenceTail(nonePlan) if err != nil { t.Fatalf("build none tail: %v", err) } // Invalid UTF-8 in no-blocking path must also be rejected (fail-closed). invalidUTF8 := string([]byte{0xFF, 0xFE}) ev, err := NewTextDeltaEvent("ch3", invalidUTF8, now) if err != nil { t.Fatalf("build invalid event: %v", err) } _, _, err = noneTail.Append(ev) if err == nil { t.Error("invalid UTF-8 must be rejected even in no-blocking path") } }) } // --------------------------------------------------------------------------- // REVIEW_API-2: Epoch-bound snapshot, committed cursor, exact look-behind // --------------------------------------------------------------------------- // TestEvidenceTailBindsPreparedReleaseToEpochSnapshot verifies that the prepared // release is bound to the epoch and that post-prepare appends are excluded // from the snapshot regardless of map iteration order. The evidence_runes // window is large enough to retain all released events in look-behind so the // snapshot-size invariant—not the trim bound—is what is under test here. func TestEvidenceTailBindsPreparedReleaseToEpochSnapshot(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} // evidence_runes=30 matches the total rune count of three 10-rune events // so the threshold triggers on the third append and the look-behind trim // bound retains all three confirmed events. req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 30, 4096) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Fill to threshold. for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", repeatRune('a', 10), now) if err != nil { t.Fatalf("build event %d: %v", i, err) } _, _, _ = tail.Append(ev) if err != nil { t.Fatalf("append %d: %v", i, err) } } cs := tail.channelState["ch"] if cs.pendingRunes != 30 { t.Fatalf("pendingRunes want 30, got %d", cs.pendingRunes) } // Prepare release. var snapshotEpoch uint64 for _, rec := range tail.epochs { if rec.epochID > snapshotEpoch { snapshotEpoch = rec.epochID } } pr, err := tail.PrepareRelease(snapshotEpoch) if err != nil { t.Fatalf("prepare: %v", err) } releaseEvents := pr.ReleaseEvents() if len(releaseEvents) != 3 { t.Errorf("releaseEvents want 3, got %d", len(releaseEvents)) } // Append another event after prepare. evExtra, err := NewTextDeltaEvent("ch", repeatRune('z', 10), now) if err != nil { t.Fatalf("build extra: %v", err) } tail.Append(evExtra) cs = tail.channelState["ch"] if len(cs.pendingEntries) != 4 { t.Fatalf("pendingEntries want 4 after extra append, got %d", len(cs.pendingEntries)) } // Confirm only the 3 events from the snapshot, not the extra. err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 3}) if err != nil { t.Fatalf("confirm snapshot: %v", err) } cs = tail.channelState["ch"] // The extra event should remain in pending. if len(cs.pendingEntries) != 1 { t.Errorf("pending want 1 (extra only), got %d", len(cs.pendingEntries)) } // Look-behind should have exactly 3 events. if len(cs.committedLookBehind) != 3 { t.Errorf("lookBehind want 3, got %d", len(cs.committedLookBehind)) } } // TestEvidenceTailRejectsOverlappingAndOversizedConfirmation verifies that // concurrent prepared tokens are rejected and that oversized confirmation is // refused. func TestEvidenceTailRejectsOverlappingAndOversizedConfirmation(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} req, err := NewFilterHoldRequirementRolling("ch", textKinds, 2) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Fill to threshold. for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "x", now) if err != nil { t.Fatalf("build: %v", err) } _, _, _ = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } // Prepare release. var prepareEpoch uint64 for _, rec := range tail.epochs { if rec.epochID > prepareEpoch { prepareEpoch = rec.epochID } } pr, err := tail.PrepareRelease(prepareEpoch) if err != nil { t.Fatalf("prepare: %v", err) } t.Run("overConfirmationRejected", func(t *testing.T) { // Try to confirm more than the snapshot size (3 events). err := tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 5}) if err == nil { t.Error("over-confirmation must be rejected") } }) t.Run("duplicatePrepareRejected", func(t *testing.T) { // Confirm the outer prepare so preparedByChannel is cleared. err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 3}) if err != nil { t.Fatalf("confirm outer prepare: %v", err) } // Re-fill and prepare again. for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", "y", now) if err != nil { t.Fatalf("build: %v", err) } _, _, _ = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } var epoch2 uint64 for _, rec := range tail.epochs { if rec.epochID > epoch2 { epoch2 = rec.epochID } } _, err = tail.PrepareRelease(epoch2) if err != nil { t.Fatalf("prepare 2: %v", err) } // Try to prepare again with the same epoch while it's already prepared. _, err = tail.PrepareRelease(epoch2) if err == nil { t.Error("duplicate prepare with same epoch must be rejected") } }) } // TestEvidenceTailPreservesCursorAndExactRuneLookBehind verifies that partial // and full confirm updates the committed cursor correctly, that oversized // events are sliced at rune boundaries, and that replace/continuation behave // differently with respect to cursor preservation. func TestEvidenceTailPreservesCursorAndExactRuneLookBehind(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} req, err := NewFilterHoldRequirementRolling("ch", textKinds, 5) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } t.Run("partialConfirmIncrementsCursor", func(t *testing.T) { for i := 0; i < 4; i++ { ev, err := NewTextDeltaEvent("ch", repeatRune('a', 2), now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } var epochID uint64 for _, rec := range tail.epochs { if rec.epochID > epochID { epochID = rec.epochID } } pr, err := tail.PrepareRelease(epochID) if err != nil { t.Fatalf("prepare: %v", err) } if tail.CommittedCursor("ch") != 0 { t.Errorf("initial cursor want 0, got %d", tail.CommittedCursor("ch")) } err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 2}) if err != nil { t.Fatalf("partial confirm: %v", err) } if tail.CommittedCursor("ch") != 2 { t.Errorf("cursor after partial want 2, got %d", tail.CommittedCursor("ch")) } }) t.Run("fullConfirmConsumesEpoch", func(t *testing.T) { for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", repeatRune('b', 2), now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } var epochID uint64 for _, rec := range tail.epochs { if rec.epochID > epochID { epochID = rec.epochID } } pr, err := tail.PrepareRelease(epochID) if err != nil { t.Fatalf("prepare: %v", err) } err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 3}) if err != nil { t.Fatalf("full confirm: %v", err) } if tail.CommittedCursor("ch") != 5 { t.Errorf("cursor after full want 5, got %d", tail.CommittedCursor("ch")) } }) t.Run("exactRuneSuffixFromOversizedEvent", func(t *testing.T) { req2, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch2", textKinds, 2, 15) if err != nil { t.Fatalf("build req2: %v", err) } bind2, err := NewFilterHoldBinding("f2", req2, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind2: %v", err) } plan2, err := NewEvidencePlan([]FilterHoldBinding{bind2}) if err != nil { t.Fatalf("build plan2: %v", err) } tail2, err := NewEvidenceTail(plan2) if err != nil { t.Fatalf("build tail2: %v", err) } for i := 0; i < 5; i++ { ev, err := NewTextDeltaEvent("ch2", repeatRune('c', 2), now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail2.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } for i := 0; i < 1; i++ { ev, err := NewTextDeltaEvent("ch2", "d", now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail2.Append(ev) if err != nil { t.Fatalf("append threshold: %v", err) } } var epochID uint64 for _, rec := range tail2.epochs { if rec.epochID > epochID { epochID = rec.epochID } } pr, err := tail2.PrepareRelease(epochID) if err != nil { t.Fatalf("prepare: %v", err) } err = tail2.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: len(pr.ReleaseEvents())}) if err != nil { t.Fatalf("confirm: %v", err) } cs := tail2.channelState["ch2"] if cs == nil { t.Fatal("channel state missing") } totalRunes := 0 for _, ev := range cs.committedLookBehind { totalRunes += runeCountForEvent(ev) } if totalRunes > 15 { t.Errorf("lookBehind runes %d exceeds maxBufferRunes 15", totalRunes) } }) t.Run("replaceClearsCursor", func(t *testing.T) { req2, err := NewFilterHoldRequirementRolling("ch_r", textKinds, 5) if err != nil { t.Fatalf("build req2: %v", err) } bind2, err := NewFilterHoldBinding("f2", req2, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind2: %v", err) } plan2, err := NewEvidencePlan([]FilterHoldBinding{bind2}) if err != nil { t.Fatalf("build plan2: %v", err) } tail2, err := NewEvidenceTail(plan2) if err != nil { t.Fatalf("build tail2: %v", err) } for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch_r", repeatRune('e', 2), now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail2.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } var epochID uint64 for _, rec := range tail2.epochs { if rec.epochID > epochID { epochID = rec.epochID } } pr, err := tail2.PrepareRelease(epochID) if err != nil { t.Fatalf("prepare: %v", err) } tail2.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 2}) if tail2.CommittedCursor("ch_r") != 2 { t.Fatalf("cursor before replace want 2, got %d", tail2.CommittedCursor("ch_r")) } tail2.ResetForReplace() if tail2.CommittedCursor("ch_r") != 0 { t.Errorf("cursor after replace want 0, got %d", tail2.CommittedCursor("ch_r")) } }) t.Run("continuationPreservesCursor", func(t *testing.T) { req3, err := NewFilterHoldRequirementRolling("ch_c", textKinds, 5) if err != nil { t.Fatalf("build req3: %v", err) } bind3, err := NewFilterHoldBinding("f3", req3, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind3: %v", err) } plan3, err := NewEvidencePlan([]FilterHoldBinding{bind3}) if err != nil { t.Fatalf("build plan3: %v", err) } tail3, err := NewEvidenceTail(plan3) if err != nil { t.Fatalf("build tail3: %v", err) } for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch_c", repeatRune('f', 2), now) if err != nil { t.Fatalf("build: %v", err) } _, _, err = tail3.Append(ev) if err != nil { t.Fatalf("append: %v", err) } } var epochID uint64 for _, rec := range tail3.epochs { if rec.epochID > epochID { epochID = rec.epochID } } pr, err := tail3.PrepareRelease(epochID) if err != nil { t.Fatalf("prepare: %v", err) } tail3.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 1}) if tail3.CommittedCursor("ch_c") != 1 { t.Fatalf("cursor before continuation want 1, got %d", tail3.CommittedCursor("ch_c")) } tail3.PrepareContinuation() if tail3.CommittedCursor("ch_c") != 1 { t.Errorf("cursor after continuation want 1, got %d", tail3.CommittedCursor("ch_c")) } cs := tail3.channelState["ch_c"] if len(cs.pendingEntries) != 0 { t.Errorf("pending after continuation want 0, got %d", len(cs.pendingEntries)) } }) } // TestEvidenceTailFragmentsReleaseOnlyCompletedSafePrefix verifies that when // fragments A and B are interleaved, completing only A produces an epoch for // the safe prefix up to A2 (the first incomplete is B1), and completing B // produces an epoch for all four events. func TestEvidenceTailFragmentsReleaseOnlyCompletedSafePrefix(t *testing.T) { toolKinds := []EventKind{EventKindToolCallFragment} req, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } // Interleave A and B. evA1, _ := NewToolCallFragmentEvent("ch", "call-A", "funcA", `{"a":1}`, now) tail.Append(evA1) evB1, _ := NewToolCallFragmentEvent("ch", "call-B", "funcB", `{"b":2}`, now) tail.Append(evB1) evA2, _ := NewToolCallFragmentEvent("ch", "call-A", "funcA", `{"a":3}`, now) tail.Append(evA2) evB2, _ := NewToolCallFragmentEvent("ch", "call-B", "funcB", `{"b":4}`, now) tail.Append(evB2) if len(tail.channelState["ch"].pendingEntries) != 4 { t.Fatalf("pending want 4, got %d", len(tail.channelState["ch"].pendingEntries)) } // Complete A only — the safe prefix includes entries 0..2 (A1, B1, A2). epochA, signalA, okA, err := tail.CompleteFragment("ch", "call-A") if err != nil { t.Fatalf("complete A: %v", err) } if !okA { t.Error("complete A should return true") } if signalA != EvidenceTailSignalTrigger { t.Errorf("complete A signal want trigger, got %s", signalA) } if epochA.ID() == 0 { t.Fatal("complete A should produce an epoch") } // Safe prefix stops at B1 (entry index 1, first incomplete fragment). prA, err := tail.PrepareRelease(epochA.ID()) if err != nil { t.Fatalf("prepare after A complete: %v", err) } if len(prA.ReleaseEvents()) != 1 { t.Errorf("safe prefix with only A complete want 1 event (A1), got %d", len(prA.ReleaseEvents())) } // Complete B — all fragments are now complete, so all 4 entries form the // safe prefix for the new epoch. epochAB, signalAB, okAB, err := tail.CompleteFragment("ch", "call-B") if err != nil { t.Fatalf("complete B: %v", err) } if !okAB { t.Error("complete B should return true") } if signalAB != EvidenceTailSignalTrigger { t.Errorf("complete B signal want trigger, got %s", signalAB) } if epochAB.ID() == 0 { t.Fatal("complete B should produce an epoch") } // All fragments complete → all 4 events released. prAB, err := tail.PrepareRelease(epochAB.ID()) if err != nil { t.Fatalf("prepare after both complete: %v", err) } if len(prAB.ReleaseEvents()) != 4 { t.Errorf("all complete want 4 events, got %d", len(prAB.ReleaseEvents())) } // PrepareRelease for the first epoch should now fail (its snapshot range // has been invalidated by the second completion). _, err = tail.PrepareRelease(epochA.ID()) if err == nil { t.Error("prepareRelease for consumed epoch should fail") } } // TestEvidenceTailAccumulatesAndBoundsFragmentByToolCall verifies that multiple // fragments with the same toolCallID accumulate and that exceeding the per-ID // bound prevents state creation. func TestEvidenceTailAccumulatesAndBoundsFragmentByToolCall(t *testing.T) { toolKinds := []EventKind{EventKindToolCallFragment} req, err := NewFilterHoldRequirementFragmentGateWithMaxBuffer("ch", toolKinds, EventKindToolCallFragment, 50) if err != nil { t.Fatalf("build requirement: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("build plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("build tail: %v", err) } t.Run("sameIdMultipleFragmentsAccumulate", func(t *testing.T) { for i := 0; i < 3; i++ { ev, err := NewToolCallFragmentEvent("ch", "call-same", "func", fmt.Sprintf(`{"n":%d}`, i), now) if err != nil { t.Fatalf("build fragment %d: %v", i, err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append fragment %d: %v", i, err) } } // Complete the fragment — since all entries are for call-same and // there are no incomplete fragments, the safe prefix includes all 3. epoch, _, ok, err := tail.CompleteFragment("ch", "call-same") if err != nil { t.Fatalf("complete: %v", err) } if !ok { t.Error("complete should return true") } if epoch.ID() == 0 { t.Fatal("complete should produce an epoch") } pr, err := tail.PrepareRelease(epoch.ID()) if err != nil { t.Fatalf("prepare: %v", err) } if len(pr.ReleaseEvents()) != 3 { t.Errorf("releaseEvents want 3, got %d", len(pr.ReleaseEvents())) } }) t.Run("perIdOverflowNoRelease", func(t *testing.T) { req2, err := NewFilterHoldRequirementFragmentGateWithMaxBuffer("ch2", toolKinds, EventKindToolCallFragment, 10) if err != nil { t.Fatalf("build req2: %v", err) } bind2, err := NewFilterHoldBinding("f2", req2, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind2: %v", err) } plan2, err := NewEvidencePlan([]FilterHoldBinding{bind2}) if err != nil { t.Fatalf("build plan2: %v", err) } tail2, err := NewEvidenceTail(plan2) if err != nil { t.Fatalf("build tail2: %v", err) } for i := 0; i < 6; i++ { ev, err := NewToolCallFragmentEvent("ch2", "call-big", "func", "abcd", now) if err != nil { t.Fatalf("build fragment: %v", err) } _, signal, err := tail2.Append(ev) if err != nil { t.Fatalf("append fragment: %v", err) } if i == 5 && signal != EvidenceTailSignalBufferOverflow { t.Errorf("sixth fragment signal want overflow, got %s", signal) } } cs := tail2.channelState["ch2"] if cs != nil { fragment := cs.fragmentState["call-big"] if fragment != nil && fragment.runes > 10 { t.Errorf("fragment runes %d exceeds per-ID bound 10", fragment.runes) } } }) } // TestEvidencePlanComposesMixedRequirementsDeterministically verifies that // mixing rolling, terminal_gate, and fragment_gate bindings on the same // channel produces a deterministic merged requirement regardless of input // order. Mode priority, subscription-kind union, min positive hard bound, // and same-mode trigger compatibility are all checked. Regressions for // reverse ordering and trigger conflicts are included. func TestEvidencePlanComposesMixedRequirementsDeterministically(t *testing.T) { textKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} toolKinds := []EventKind{EventKindToolCallFragment} t.Run("rollingThenFragment", func(t *testing.T) { r, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 300, 5000) if err != nil { t.Fatalf("build rolling: %v", err) } f, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build fragment: %v", err) } b1, err := NewFilterHoldBinding("f1", r, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", f, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } req, _ := plan.BlockingRequirement("ch") if req.Mode() != FilterHoldModeFragmentGate { t.Errorf("strongest mode want fragment_gate, got %s", req.Mode()) } if req.MaxBufferRunes() != 4096 { // min(5000, default 4096) t.Errorf("min hard bound want 4096, got %d", req.MaxBufferRunes()) } kinds := plan.BlockingKinds("ch") if len(kinds) != 3 { t.Errorf("blocking kinds want 3, got %d", len(kinds)) } }) t.Run("fragmentThenRolling", func(t *testing.T) { f, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build fragment: %v", err) } r, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 300, 5000) if err != nil { t.Fatalf("build rolling: %v", err) } b1, err := NewFilterHoldBinding("f1", f, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", r, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } req, _ := plan.BlockingRequirement("ch") if req.Mode() != FilterHoldModeFragmentGate { t.Errorf("strongest mode want fragment_gate, got %s", req.Mode()) } if req.MaxBufferRunes() != 4096 { t.Errorf("min hard bound want 4096, got %d", req.MaxBufferRunes()) } }) t.Run("terminalWithFragmentTriggerSupported", func(t *testing.T) { f, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build fragment: %v", err) } te, err := NewFilterHoldRequirementTerminalGate("ch", textKinds, EventKindTerminal) if err != nil { t.Fatalf("build terminal: %v", err) } b1, err := NewFilterHoldBinding("f1", f, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", te, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("multi-requirement evidence plan failed: %v", err) } if len(plan.BindingsForChannel("ch")) != 2 { t.Errorf("bindings count want 2, got %d", len(plan.BindingsForChannel("ch"))) } }) t.Run("twoTerminalTriggersSupported", func(t *testing.T) { t1, err := NewFilterHoldRequirementTerminalGate("ch", textKinds, EventKindTerminal) if err != nil { t.Fatalf("build terminal 1: %v", err) } t2, err := NewFilterHoldRequirementTerminalGate("ch", textKinds, EventKindProviderError) if err != nil { t.Fatalf("build terminal 2: %v", err) } b1, err := NewFilterHoldBinding("f1", t1, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", t2, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("multi-terminal trigger plan failed: %v", err) } if len(plan.BindingsForChannel("ch")) != 2 { t.Errorf("bindings count want 2, got %d", len(plan.BindingsForChannel("ch"))) } }) t.Run("sameTerminalTriggerAccepted", func(t *testing.T) { t1, err := NewFilterHoldRequirementTerminalGate("ch", textKinds, EventKindTerminal) if err != nil { t.Fatalf("build terminal 1: %v", err) } t2, err := NewFilterHoldRequirementTerminalGateWithMaxBuffer("ch", textKinds, EventKindTerminal, 8192) if err != nil { t.Fatalf("build terminal 2: %v", err) } b1, err := NewFilterHoldBinding("f1", t1, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", t2, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } req, _ := plan.BlockingRequirement("ch") if req.MaxBufferRunes() != 4096 { // min(8192, default) t.Errorf("min hard bound want 4096, got %d", req.MaxBufferRunes()) } }) t.Run("rollingMixedMinMaxBound", func(t *testing.T) { r1, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 200, 3000) if err != nil { t.Fatalf("build rolling 1: %v", err) } r2, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 400, 2000) if err != nil { t.Fatalf("build rolling 2: %v", err) } b1, err := NewFilterHoldBinding("f1", r1, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", r2, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } req, _ := plan.BlockingRequirement("ch") if req.EvidenceRunes() != 400 { t.Errorf("evidenceRunes want max(200,400)=400, got %d", req.EvidenceRunes()) } if req.MaxBufferRunes() != 2000 { t.Errorf("maxBufferRunes want min(3000,2000)=2000, got %d", req.MaxBufferRunes()) } }) } // TestEvidenceTailSeparatesBlockingAndObserveSubscriptions verifies that // observe-only event kinds pass through without contributing to pending state // or creating ready epochs for the blocking path, while blocking kinds // contribute to pending as expected. func TestEvidenceTailSeparatesBlockingAndObserveSubscriptions(t *testing.T) { textKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} toolKinds := []EventKind{EventKindToolCallFragment} t.Run("observeOnlyEventDoesNotCreatePending", func(t *testing.T) { // Non-blocking rolling (observe-only) + blocking fragment_gate. // textDelta/reasoningDelta from the non-blocking binding are // observe-only; toolCallFragment from the blocking binding is // blocking. r, err := NewFilterHoldRequirementRolling("ch", textKinds, 5) if err != nil { t.Fatalf("build rolling: %v", err) } f, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build fragment: %v", err) } b1, err := NewFilterHoldBinding("f1", r, FilterEnforcementObserveOnly) // blocksRelease=false → observe-only if err != nil { t.Fatalf("build binding 1: %v", err) } b2, err := NewFilterHoldBinding("f2", f, FilterEnforcementBlocking) // blocking if err != nil { t.Fatalf("build binding 2: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("NewEvidenceTail: %v", err) } // reasoningDelta is observe-only: should pass through without state. ev, err := NewReasoningDeltaEvent("ch", "hello", now) if err != nil { t.Fatalf("build event: %v", err) } _, signal, err := tail.Append(ev) if err != nil { t.Fatalf("append observe: %v", err) } if signal != EvidenceTailSignalNone { t.Errorf("observe-only signal want none, got %s", signal) } if tail.channelState["ch"] != nil { t.Error("observe-only event must not create channel state") } }) t.Run("blockingEventContributesToPending", func(t *testing.T) { r, err := NewFilterHoldRequirementRolling("ch", textKinds, 5) if err != nil { t.Fatalf("build rolling: %v", err) } bind, err := NewFilterHoldBinding("f", r, FilterEnforcementBlocking) if err != nil { t.Fatalf("build binding: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("NewEvidenceTail: %v", err) } // textDelta is the sole blocking kind: should create channel state // and contribute to pending runes. ev, err := NewTextDeltaEvent("ch", "hello", now) if err != nil { t.Fatalf("build event: %v", err) } _, _, err = tail.Append(ev) if err != nil { t.Fatalf("append blocking: %v", err) } cs := tail.channelState["ch"] if cs == nil { t.Fatal("blocking event must create channel state") } if cs.pendingRunes != 5 { t.Errorf("pendingRunes want 5, got %d", cs.pendingRunes) } }) } // TestEvidenceTailBindsOnlyReadyEpochSnapshots verifies via a public-API table // that non-ready epochs (sub-threshold, unsubscribed, overflow) cannot be // prepared, and that post-snapshot append does not expand the snapshot. func TestEvidenceTailBindsOnlyReadyEpochSnapshots(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} t.Run("subThresholdEpochNoRelease", func(t *testing.T) { // Sub-threshold: evidence_runes=100 but only 10 runes added. req, err := NewFilterHoldRequirementRolling("ch", textKinds, 100) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Append 1 event. No threshold → no epoch. ev, _ := NewTextDeltaEvent("ch", repeatRune('a', 10), now) tail.Append(ev) // No epochs exist. for id := range tail.epochs { t.Errorf("unexpected epoch %d for sub-threshold", id) } }) t.Run("unsubscribedEpochNoRelease", func(t *testing.T) { req, err := NewFilterHoldRequirementRolling("ch", textKinds, 3) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // reasoning_delta not subscribed. Append returns an epoch but the // channel state must not be created and the event must not contribute // to pending runes. ev, _ := NewReasoningDeltaEvent("ch", "unsub", now) _, signal, _ := tail.Append(ev) if signal != EvidenceTailSignalNone { t.Errorf("unsubscribed signal want none, got %s", signal) } if tail.channelState["ch"] != nil { t.Errorf("unsubscribed event must not create channel state, got pendingRunes=%d", tail.channelState["ch"].pendingRunes) } }) t.Run("postSnapshotAppendExcluded", func(t *testing.T) { req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 30, 4096) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Fill to threshold (3 events x 10 runes = 30). for i := 0; i < 3; i++ { ev, _ := NewTextDeltaEvent("ch", repeatRune('a', 10), now) tail.Append(ev) } // Find threshold epoch. var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } if len(pr.ReleaseEvents()) != 3 { t.Fatalf("snapshot want 3 events, got %d", len(pr.ReleaseEvents())) } // Zero-confirm to clear the prepared token. tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 0}) // Append after prepare. evExtra, _ := NewTextDeltaEvent("ch", repeatRune('z', 10), now) tail.Append(evExtra) // Re-prepare same epoch — snapshot should still be 3, not 4. pr2, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("re-prepare: %v", err) } if len(pr2.ReleaseEvents()) != 3 { t.Errorf("re-prepared snapshot want 3, got %d", len(pr2.ReleaseEvents())) } }) t.Run("overflowEpochCannotPrepare", func(t *testing.T) { req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 3, 10) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Fill to 10. for i := 0; i < 10; i++ { ev, _ := NewTextDeltaEvent("ch", "x", now) tail.Append(ev) } // 11th causes overflow. Append returns an overflow signal. ev11, _ := NewTextDeltaEvent("ch", "x", now) beforeOverflow := len(tail.epochs) _, signal, _ := tail.Append(ev11) if signal != EvidenceTailSignalBufferOverflow { t.Errorf("signal want overflow, got %s", signal) } // Overflow must not create any new epoch record. if len(tail.epochs) != beforeOverflow { t.Errorf("overflow must not create epoch record, before=%d after=%d", beforeOverflow, len(tail.epochs)) } }) } // TestEvidenceTailRepreparesOnlyUnconfirmedSnapshotSuffix verifies that after // zero or partial confirm, re-preparing the same epoch returns only the // unconfirmed suffix of the original snapshot. New appends after prepare are // excluded from the re-prepared snapshot. func TestEvidenceTailRepreparesOnlyUnconfirmedSnapshotSuffix(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} t.Run("zeroConfirmReprepareReturnsFull", func(t *testing.T) { req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 30, 4096) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", fmt.Sprintf("ev-%d-pad-%s", i, repeatRune('b', 8)), now) if err != nil { t.Fatalf("new text delta event %d: %v", i, err) } if _, sig, err := tail.Append(ev); err != nil { t.Fatalf("append text event %d: %v", i, err) } else if sig != EvidenceTailSignalNone && sig != EvidenceTailSignalThreshold { t.Errorf("append text event %d signal want none or threshold, got %s", i, sig) } } var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare release: %v", err) } // Zero confirm. if err := tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 0}); err != nil { t.Fatalf("zero confirm: %v", err) } // Re-prepare: should get same 3 events. pr2, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("re-prepare after zero confirm: %v", err) } if len(pr2.ReleaseEvents()) != 3 { t.Errorf("re-prepared want 3 events, got %d", len(pr2.ReleaseEvents())) } // Exact content check with distinct payloads. for i, re := range pr2.ReleaseEvents() { text, err := re.AsTextDelta() if err != nil { t.Fatalf("as text delta %d: %v", i, err) } want := fmt.Sprintf("ev-%d-pad-%s", i, repeatRune('b', 8)) if text != want { t.Errorf("re-prepared[%d] text want %q, got %q", i, want, text) } } }) t.Run("partialConfirmReprepareReturnsSuffix", func(t *testing.T) { req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 30, 4096) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } for i := 0; i < 4; i++ { ev, err := NewTextDeltaEvent("ch", fmt.Sprintf("ev-%d-pad-%s", i, repeatRune('x', 8)), now) if err != nil { t.Fatalf("new text delta event %d: %v", i, err) } if _, sig, err := tail.Append(ev); err != nil { t.Fatalf("append text event %d: %v", i, err) } else if sig != EvidenceTailSignalNone && sig != EvidenceTailSignalThreshold { t.Errorf("append text event %d signal want none or threshold, got %s", i, sig) } } var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare release: %v", err) } if len(pr.ReleaseEvents()) != 4 { t.Fatalf("initial snapshot want 4, got %d", len(pr.ReleaseEvents())) } // Partial confirm 2. err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 2}) if err != nil { t.Fatalf("partial confirm 2: %v", err) } // Re-prepare: should get remaining 2 events. pr2, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("re-prepare after partial: %v", err) } if len(pr2.ReleaseEvents()) != 2 { t.Errorf("re-prepared suffix want 2, got %d", len(pr2.ReleaseEvents())) } // Exact event content verification: remaining events must be entries [2,3]. for i, re := range pr2.ReleaseEvents() { text, err := re.AsTextDelta() if err != nil { t.Fatalf("as text delta %d: %v", i, err) } want := fmt.Sprintf("ev-%d-pad-%s", i+2, repeatRune('x', 8)) if text != want { t.Errorf("re-prepared[%d] text want %q, got %q", i, want, text) } } // Second full confirm must succeed on the re-prepared suffix. err = tail.ConfirmRelease(pr2.Token(), ReleaseConfirmation{ReleasedEvents: 2}) if err != nil { t.Fatalf("second full confirm on re-prepared suffix: %v", err) } // Epoch must now be fully consumed. epochRec, ok := tail.epochs[ep] if !ok { t.Fatal("epoch must exist after second confirm") } if !epochRec.consumed { t.Error("epoch must be consumed after full second confirm") } if epochRec.confirmedCount != 4 { t.Errorf("confirmedCount want 4, got %d", epochRec.confirmedCount) } // Cursor must reflect all 4 confirmed events. if tail.CommittedCursor("ch") != 4 { t.Errorf("committedCursor want 4, got %d", tail.CommittedCursor("ch")) } // Post-snapshot append must remain in pending, not part of consumed epoch. evExtra, err := NewTextDeltaEvent("ch", fmt.Sprintf("extra-%d", 99), now) if err != nil { t.Fatalf("new extra text delta: %v", err) } if _, sig, err := tail.Append(evExtra); err != nil { t.Fatalf("append extra: %v", err) } else if sig != EvidenceTailSignalNone { t.Errorf("append extra signal want none, got %s", sig) } cs := tail.channelState["ch"] if len(cs.pendingEntries) != 1 { t.Errorf("post-append pending want 1, got %d", len(cs.pendingEntries)) } if cs.pendingEntries[0].event.Kind() != EventKindTextDelta { t.Error("post-append entry must be text delta") } }) t.Run("postPrepareAppendExcludedFromReprepare", func(t *testing.T) { req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 30, 4096) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("ch", fmt.Sprintf("ev-%d-pad-%s", i, repeatRune('c', 8)), now) if err != nil { t.Fatalf("new text delta event %d: %v", i, err) } if _, sig, err := tail.Append(ev); err != nil { t.Fatalf("append text event %d: %v", i, err) } else if sig != EvidenceTailSignalNone && sig != EvidenceTailSignalThreshold { t.Errorf("append text event %d signal want none or threshold, got %s", i, sig) } } var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare release: %v", err) } if err := tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 0}); err != nil { t.Fatalf("zero confirm: %v", err) } // Append after prepare. evExtra, err := NewTextDeltaEvent("ch", fmt.Sprintf("extra-%d-pad-%s", 99, repeatRune('d', 8)), now) if err != nil { t.Fatalf("new extra text delta: %v", err) } if _, sig, err := tail.Append(evExtra); err != nil { t.Fatalf("append extra: %v", err) } else if sig != EvidenceTailSignalNone && sig != EvidenceTailSignalThreshold { t.Errorf("append extra signal want none or threshold, got %s", sig) } // Re-prepare: should still be 3, not 4. pr2, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("re-prepare: %v", err) } if len(pr2.ReleaseEvents()) != 3 { t.Errorf("re-prepared want 3 (post-append excluded), got %d", len(pr2.ReleaseEvents())) } // Exact content check: first three payloads in original order. originalPayloads := []string{ fmt.Sprintf("ev-0-pad-%s", repeatRune('c', 8)), fmt.Sprintf("ev-1-pad-%s", repeatRune('c', 8)), fmt.Sprintf("ev-2-pad-%s", repeatRune('c', 8)), } for i, re := range pr2.ReleaseEvents() { text, err := re.AsTextDelta() if err != nil { t.Fatalf("as text delta %d: %v", i, err) } if text != originalPayloads[i] { t.Errorf("re-prepared[%d] want %q, got %q", i, originalPayloads[i], text) } } // Extra payload must NOT appear in re-prepared result. for _, re := range pr2.ReleaseEvents() { text, err := re.AsTextDelta() if err != nil { t.Fatalf("as text delta extra check: %v", err) } if strings.Contains(text, "extra-") { t.Errorf("extra payload must not appear in re-prepared result, got %q", text) } } }) } // TestEvidenceTailFragmentCompletionCreatesSafePrefixEpoch verifies that // CompleteFragment returns (epoch, signal, true, nil) with a contiguous // completed safe-prefix epoch, that interleaved A/B completion works, // that duplicate/unknown completion returns zero epoch, and that // over-confirm (confirming more than the snapshot) is rejected. func TestEvidenceTailFragmentCompletionCreatesSafePrefixEpoch(t *testing.T) { toolKinds := []EventKind{EventKindToolCallFragment} req, err := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } t.Run("interleavedABCompletion", func(t *testing.T) { // Interleave A and B with three fragments: A1, B1, A2. evA1, err := NewToolCallFragmentEvent("ch", "a1", "fn", `{"x":1}`, now) if err != nil { t.Fatalf("new A1 fragment: %v", err) } if _, _, err := tail.Append(evA1); err != nil { t.Fatalf("append A1: %v", err) } evB1, err := NewToolCallFragmentEvent("ch", "b1", "fn", `{"y":2}`, now) if err != nil { t.Fatalf("new B1 fragment: %v", err) } if _, _, err := tail.Append(evB1); err != nil { t.Fatalf("append B1: %v", err) } evA2, err := NewToolCallFragmentEvent("ch", "a2", "fn", `{"x":3}`, now) if err != nil { t.Fatalf("new A2 fragment: %v", err) } if _, _, err := tail.Append(evA2); err != nil { t.Fatalf("append A2: %v", err) } // Complete A (a1+a2): safe prefix = [A1] (B1 at index 1 is incomplete). epA, sigA, okA, err := tail.CompleteFragment("ch", "a1") if err != nil { t.Fatalf("complete A: %v", err) } if !okA { t.Error("complete A should return true") } if sigA != EvidenceTailSignalTrigger { t.Errorf("signal want trigger, got %s", sigA) } if epA.ID() == 0 { t.Fatal("complete A should produce an epoch") } prA, err := tail.PrepareRelease(epA.ID()) if err != nil { t.Fatalf("prepare A release: %v", err) } if len(prA.ReleaseEvents()) != 1 { t.Errorf("A safe prefix want 1 event, got %d", len(prA.ReleaseEvents())) } // Complete B (b1): now a1 and b1 are done, but a2 is still pending. // safe prefix = [A1, B1] (a2 at index 2 is incomplete). epB, sigB, okB, err := tail.CompleteFragment("ch", "b1") if err != nil { t.Fatalf("complete B: %v", err) } if !okB { t.Error("complete B should return true") } if sigB != EvidenceTailSignalTrigger { t.Errorf("signal want trigger, got %s", sigB) } prB, err := tail.PrepareRelease(epB.ID()) if err != nil { t.Fatalf("prepare B release: %v", err) } if len(prB.ReleaseEvents()) != 2 { t.Errorf("AB safe prefix (a2 incomplete) want 2 events, got %d", len(prB.ReleaseEvents())) } // Complete A again? No — a1 is already completed. Complete a2. epA2, _, okA2, err := tail.CompleteFragment("ch", "a2") if err != nil { t.Fatalf("complete A2: %v", err) } if !okA2 { t.Error("complete A2 should return true") } prA2, err := tail.PrepareRelease(epA2.ID()) if err != nil { t.Fatalf("prepare A2 release: %v", err) } if len(prA2.ReleaseEvents()) != 3 { t.Errorf("all complete want 3 events, got %d", len(prA2.ReleaseEvents())) } // Re-preparing A should fail (invalidated by later completion). _, err = tail.PrepareRelease(epA.ID()) if err == nil { t.Error("re-prepare invalidated epoch should fail") } }) t.Run("duplicateCompletionReturnsZero", func(t *testing.T) { ev, err := NewToolCallFragmentEvent("ch", "dup", "fn", `{}`, now) if err != nil { t.Fatalf("new dup fragment: %v", err) } if _, _, err := tail.Append(ev); err != nil { t.Fatalf("append dup: %v", err) } ep1, _, ok1, err := tail.CompleteFragment("ch", "dup") if err != nil { t.Fatalf("complete dup: %v", err) } if !ok1 || ep1.ID() == 0 { t.Fatal("first complete should succeed") } ep2, sig2, ok2, err := tail.CompleteFragment("ch", "dup") if err != nil { t.Fatalf("duplicate complete: %v", err) } if ok2 { t.Error("duplicate complete should return false") } if ep2.ID() != 0 { t.Error("duplicate complete should return zero epoch") } if sig2 != EvidenceTailSignalNone { t.Errorf("duplicate complete signal want none, got %s", sig2) } }) t.Run("unknownToolIdReturnsZero", func(t *testing.T) { ep, sig, ok, err := tail.CompleteFragment("ch", "nonexistent") if err != nil { t.Fatalf("complete nonexistent: %v", err) } if ok { t.Error("unknown ID should return false") } if ep.ID() != 0 { t.Error("unknown ID should return zero epoch") } if sig != EvidenceTailSignalNone { t.Errorf("unknown ID signal want none, got %s", sig) } }) } // TestEvidenceTailFragmentBoundNeverCreatesReleaseEpoch verifies that // exceeding the per-ID bound never creates a releaseable epoch or token, // even on the first overflow event (limit+1). func TestEvidenceTailFragmentBoundNeverCreatesReleaseEpoch(t *testing.T) { toolKinds := []EventKind{EventKindToolCallFragment} t.Run("exactBoundNoRelease", func(t *testing.T) { req, err := NewFilterHoldRequirementFragmentGateWithMaxBuffer("ch", toolKinds, EventKindToolCallFragment, 10) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // 2 fragments of 5 runes each = exactly 10 (the bound). ev1, _ := NewToolCallFragmentEvent("ch", "id", "fn", "aaaaa", now) ev2, _ := NewToolCallFragmentEvent("ch", "id", "fn", "aaaaa", now) tail.Append(ev1) tail.Append(ev2) // Complete: safe prefix includes both entries = 10 runes. ep, _, ok, _ := tail.CompleteFragment("ch", "id") if !ok { t.Fatal("complete at exact bound should return true") } pr, err := tail.PrepareRelease(ep.ID()) if err != nil { t.Fatalf("prepare at exact bound: %v", err) } if len(pr.ReleaseEvents()) != 2 { t.Errorf("release want 2 events, got %d", len(pr.ReleaseEvents())) } }) t.Run("limitPlusOneNeverCreatesEpoch", func(t *testing.T) { // Bound = 10. 3 fragments of 4 runes = 12 > 10. The 3rd should // overflow and not be added to pending state. req, err := NewFilterHoldRequirementFragmentGateWithMaxBuffer("ch", toolKinds, EventKindToolCallFragment, 10) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } ev1, _ := NewToolCallFragmentEvent("ch", "id", "fn", "abcd", now) ev2, _ := NewToolCallFragmentEvent("ch", "id", "fn", "abcd", now) tail.Append(ev1) tail.Append(ev2) // 3rd fragment overflows — signal is overflow, state unchanged. ev3, _ := NewToolCallFragmentEvent("ch", "id", "fn", "abcd", now) _, signal, _ := tail.Append(ev3) if signal != EvidenceTailSignalBufferOverflow { t.Errorf("signal want overflow, got %s", signal) } // Verify the overflow did not add a 3rd pending entry. cs := tail.channelState["ch"] if cs == nil { t.Fatal("channel state missing") } if len(cs.pendingEntries) != 2 { t.Errorf("pendingEntries want 2 (overflow rejected), got %d", len(cs.pendingEntries)) } // The fragment state for "id" still has 2 entries (the overflow was rejected). frag := cs.fragmentState["id"] if frag == nil { t.Fatal("fragment state missing") } if frag.runes != 8 { t.Errorf("fragment runes want 8, got %d", frag.runes) } // No non-invalidated overflow epoch should exist. for _, r := range tail.epochs { if !r.invalidated { t.Errorf("overflow must not create non-invalidated epoch, got ID=%d", r.epochID) } } }) } // TestEvidenceTailBoundsAndCopiesEffectiveLookBehind verifies that the // committed look-behind is trimmed to the channel's effective evidence_runes // (Unicode suffix), that the accessor returns a deep copy isolated from // internal mutation, and that replace/continuation behave correctly with // respect to cursor and look-behind preservation. func TestEvidenceTailBoundsAndCopiesEffectiveLookBehind(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} t.Run("effectiveEvidenceRuneBound", func(t *testing.T) { // evidence_runes=10, maxBufferRunes=4096. Confirm enough events to // exceed 10 runes and verify lookBehind is trimmed to ≤10 runes. req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 10, 4096) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Add 3 events of 5 runes each = 15 runes (exceeds 10). for i := 0; i < 3; i++ { ev, _ := NewTextDeltaEvent("ch", repeatRune('a', 5), now) tail.Append(ev) } var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, _ := tail.PrepareRelease(ep) tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 3}) lb := tail.EffectiveLookBehind("ch") totalRunes := 0 for _, ev := range lb { totalRunes += runeCountForEvent(ev) } if totalRunes > 10 { t.Errorf("lookBehind runes %d exceeds effectiveEvidenceRunes 10", totalRunes) } if len(lb) == 0 { t.Error("lookBehind must not be empty") } }) t.Run("callerMutationIsolation", func(t *testing.T) { req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 5, 4096) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // 5 events of 1 rune each = 5 runes (threshold). for i := 0; i < 5; i++ { ev, _ := NewTextDeltaEvent("ch", "x", now) tail.Append(ev) } var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, _ := tail.PrepareRelease(ep) tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 5}) lb := tail.EffectiveLookBehind("ch") orig := make([]NormalizedEvent, len(lb)) copy(orig, lb) // Mutate returned slice. lb[0] = NormalizedEvent{} // Second call should return a fresh copy. lb2 := tail.EffectiveLookBehind("ch") if len(lb2) != len(orig) { t.Errorf("lookBehind length changed: got %d, want %d", len(lb2), len(orig)) } for i := range lb2 { if lb2[i].Kind() != orig[i].Kind() { t.Errorf("lookBehind[%d] mutated: got %q, want %q", i, lb2[i].Kind(), orig[i].Kind()) } } }) t.Run("replaceClearsLookBehindAndCursor", func(t *testing.T) { req, err := NewFilterHoldRequirementRolling("ch", textKinds, 5) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // 3 events of 1 rune each = 3 runes. Need 5 for threshold. for i := 0; i < 5; i++ { ev, _ := NewTextDeltaEvent("ch", "x", now) tail.Append(ev) } var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, _ := tail.PrepareRelease(ep) tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 2}) if tail.CommittedCursor("ch") != 2 { t.Fatalf("cursor before replace want 2, got %d", tail.CommittedCursor("ch")) } if len(tail.EffectiveLookBehind("ch")) != 2 { t.Fatalf("lookBehind before replace want 2, got %d", len(tail.EffectiveLookBehind("ch"))) } tail.ResetForReplace() if tail.CommittedCursor("ch") != 0 { t.Errorf("cursor after replace want 0, got %d", tail.CommittedCursor("ch")) } if tail.EffectiveLookBehind("ch") != nil { t.Errorf("lookBehind after replace want nil, got %v", tail.EffectiveLookBehind("ch")) } }) t.Run("continuationPreservesCursorAndLookBehind", func(t *testing.T) { req, err := NewFilterHoldRequirementRolling("ch", textKinds, 5) if err != nil { t.Fatalf("build req: %v", err) } bind, err := NewFilterHoldBinding("f", req, FilterEnforcementBlocking) if err != nil { t.Fatalf("build bind: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{bind}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // 3 events of 1 rune each — need 5 for threshold. for i := 0; i < 5; i++ { ev, _ := NewTextDeltaEvent("ch", "x", now) tail.Append(ev) } var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, _ := tail.PrepareRelease(ep) tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 1}) tail.PrepareContinuation() if tail.CommittedCursor("ch") != 1 { t.Errorf("cursor after continuation want 1, got %d", tail.CommittedCursor("ch")) } if len(tail.EffectiveLookBehind("ch")) != 1 { t.Errorf("lookBehind after continuation want 1, got %d", len(tail.EffectiveLookBehind("ch"))) } }) } // TestEvidenceTailRejectsInvalidReleaseLifecycle verifies via a public API that // stale/invalidated/consumed tokens, overlapping prepared tokens, already- // prepared epochs, and out-of-bounds confirmations are all rejected with // deterministic errors (never panic). func TestEvidenceTailRejectsInvalidReleaseLifecycle(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} t.Run("unknownEpochRejected", func(t *testing.T) { req, err := NewFilterHoldRequirementRolling("ch", textKinds, 5) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Non-existent epoch. _, err = tail.PrepareRelease(99999) if err == nil { t.Error("unknown epoch must be rejected") } }) t.Run("consumedEpochRejected", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } // Full confirm → epoch consumed. if err := tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 3}); err != nil { t.Fatalf("confirm: %v", err) } // Re-prepare consumed epoch must fail. _, err = tail.PrepareRelease(ep) if err == nil { t.Error("consumed epoch must be rejected on re-prepare") } }) t.Run("invalidatedEpochRejected", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } // Continuation invalidates the epoch. tail.PrepareContinuation() // Confirm with invalidated token must fail (no panic). err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("invalidated token must be rejected") } }) t.Run("duplicatePrepareRejected", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) // First prepare. _, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare 1: %v", err) } // Second prepare on same channel must fail (overlapping token). _, err = tail.PrepareRelease(ep) if err == nil { t.Error("duplicate prepare must be rejected") } }) t.Run("overConfirmationRejected", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 99}) if err == nil { t.Error("over-confirmation must be rejected") } }) t.Run("staleTokenAfterResetRejected", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() // ResetForReplace invalidates. tail.ResetForReplace() // Confirm with stale token must fail. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("stale token after reset must be rejected") } }) t.Run("staleTokenAfterTerminalDiscardRejected", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() tail.DiscardPendingForTerminal() // Confirm with stale token must fail. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("stale token after terminal discard must be rejected") } }) t.Run("staleTokenAfterContinuationRejected", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() tail.PrepareContinuation() // Confirm with stale token must fail. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("stale token after continuation must be rejected") } }) t.Run("snapshotMismatchRejectsConfirm", func(t *testing.T) { req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 30, 4096) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Fill to threshold (3 events). for i := 0; i < 3; i++ { ev, _ := NewTextDeltaEvent("ch", repeatRune('a', 10), now) tail.Append(ev) } ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } // Zero confirm (does not clear preparedByChannel — token still valid). tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 0}) // Append a new event after prepare. evExtra, _ := NewTextDeltaEvent("ch", repeatRune('z', 10), now) tail.Append(evExtra) // Re-prepare (zero confirm clears preparedByChannel, so prepare succeeds). _, err = tail.PrepareRelease(ep) if err != nil { t.Fatalf("re-prepare: %v", err) } // Confirm with old token. err = tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 0}) if err == nil { t.Error("confirm with stale token after re-prepare must fail") } }) } // TestEvidenceTailRecoveryTransitions verifies that ResetForReplace, // PrepareContinuation, and DiscardPendingForTerminal all invalidate prepared // tokens with deterministic errors (no panic), clear prepared/preparedByChannel // state, and that subsequent append/prepare cycles work correctly after each // transition. func TestEvidenceTailRecoveryTransitions(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} t.Run("ResetForReplaceInvalidatesAllTokens", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() tail.ResetForReplace() // Token must be invalid. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("reset must invalidate token") } // After reset, state is cleared. New events and new prepare cycle should work. for i := 0; i < 5; i++ { ev, _ := NewTextDeltaEvent("ch", "x", now) tail.Append(ev) } newEp := findLatestEpoch(t, tail) pr2, err := tail.PrepareRelease(newEp) if err != nil { t.Fatalf("new prepare after reset: %v", err) } // New token must work. err = tail.ConfirmRelease(pr2.Token(), ReleaseConfirmation{ReleasedEvents: 3}) if err != nil { t.Fatalf("new confirm after reset: %v", err) } }) t.Run("PrepareContinuationInvalidatesTokensPreservesLookBehind", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() // Confirm partial to create look-behind. tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) // Append pending event. ev, _ := NewTextDeltaEvent("ch", "y", now) tail.Append(ev) tail.PrepareContinuation() // Token must be invalid. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("continuation must invalidate token") } // Look-behind must be preserved. lb := tail.EffectiveLookBehind("ch") if len(lb) != 1 { t.Errorf("lookBehind want 1 after continuation, got %d", len(lb)) } // Cursor must be preserved. if tail.CommittedCursor("ch") != 1 { t.Errorf("cursor want 1 after continuation, got %d", tail.CommittedCursor("ch")) } // Pending must be cleared. cs := tail.channelState["ch"] if len(cs.pendingEntries) != 0 { t.Errorf("pending must be cleared, got %d", len(cs.pendingEntries)) } }) t.Run("DiscardPendingForTerminalInvalidatesTokensPreservesLookBehind", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() // Confirm partial. tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) // Append pending. ev, _ := NewTextDeltaEvent("ch", "z", now) tail.Append(ev) tail.DiscardPendingForTerminal() // Token must be invalid. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Error("terminal discard must invalidate token") } // Look-behind preserved. lb := tail.EffectiveLookBehind("ch") if len(lb) != 1 { t.Errorf("lookBehind want 1, got %d", len(lb)) } // Pending cleared. cs := tail.channelState["ch"] if len(cs.pendingEntries) != 0 { t.Errorf("pending must be cleared, got %d", len(cs.pendingEntries)) } }) t.Run("recoveryTransitionsNoPanicOnConfirm", func(t *testing.T) { // Verify that outstanding tokens after each recovery transition // always return an error (not panic) regardless of token value. transitions := []func(t *EvidenceTail){ func(t *EvidenceTail) { t.ResetForReplace() }, func(t *EvidenceTail) { t.PrepareContinuation() }, func(t *EvidenceTail) { t.DiscardPendingForTerminal() }, } for _, tr := range transitions { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } token := pr.Token() tr(tail) // Must not panic; must return error. err = tail.ConfirmRelease(token, ReleaseConfirmation{ReleasedEvents: 1}) if err == nil { t.Errorf("transition must invalidate token") } } }) t.Run("invalidatedEpochHasClearedRecordFields", func(t *testing.T) { tail := buildThresholdTail(t, textKinds, 5) ep := findLatestEpoch(t, tail) _, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare: %v", err) } tail.PrepareContinuation() // Find the record and verify fields are cleared. for _, rec := range tail.epochs { if rec.epochID == ep { if !rec.invalidated { t.Error("epoch must be invalidated") } if rec.prepared { t.Error("prepared must be false after invalidation") } if rec.token != "" { t.Error("token must be empty after invalidation") } if rec.snapshotSize != 0 { t.Errorf("snapshotSize must be 0 after invalidation, got %d", rec.snapshotSize) } } } }) } // TestEvidenceTailFragmentCompletionPreservesOtherChannels verifies that // CompleteFragment invalidates only the same channel's superseded epochs and // prepared tokens, while leaving unrelated channel outstanding tokens intact. // A rolling/text channel with a prepared release token must still be able to // confirm normally after a fragment/tool channel CompleteFragment invalidates // its own channel's state. func TestEvidenceTailFragmentCompletionPreservesOtherChannels(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} toolKinds := []EventKind{EventKindToolCallFragment} // Build a plan with two channels: "text" (rolling) and "tool" (fragment_gate). textReq, err := NewFilterHoldRequirementRollingWithMaxBuffer("text", textKinds, 15, 4096) if err != nil { t.Fatalf("build text req: %v", err) } toolReq, err := NewFilterHoldRequirementFragmentGate("tool", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build tool req: %v", err) } tool2Req, err := NewFilterHoldRequirementFragmentGate("tool2", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build tool2 req: %v", err) } interleaveReq, err := NewFilterHoldRequirementFragmentGate("interleave", toolKinds, EventKindToolCallFragment) if err != nil { t.Fatalf("build interleave req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{ mustBind(t, textReq), mustBind(t, toolReq), mustBind(t, tool2Req), mustBind(t, interleaveReq), }) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } t.Run("fragmentCompletionDoesNotBreakRollingTokenConfirm", func(t *testing.T) { // Prepare a rolling/text token first. for i := 0; i < 3; i++ { ev, err := NewTextDeltaEvent("text", fmt.Sprintf("text-%d", i), now) if err != nil { t.Fatalf("new text delta event %d: %v", i, err) } if _, _, err := tail.Append(ev); err != nil { t.Fatalf("append text event %d: %v", i, err) } } var textEp uint64 for _, r := range tail.epochs { if r.channel == "text" && r.epochID > textEp { textEp = r.epochID } } textPR, err := tail.PrepareRelease(textEp) if err != nil { t.Fatalf("prepare text release: %v", err) } if len(textPR.ReleaseEvents()) != 3 { t.Fatalf("text snapshot want 3 events, got %d", len(textPR.ReleaseEvents())) } // Now do fragment operations on the "tool" channel. evA1, err := NewToolCallFragmentEvent("tool", "a1", "fn", `{"x":1}`, now) if err != nil { t.Fatalf("new tool fragment A1: %v", err) } if _, _, err := tail.Append(evA1); err != nil { t.Fatalf("append tool A1: %v", err) } evA2, err := NewToolCallFragmentEvent("tool", "a1", "fn", `{"x":2}`, now) if err != nil { t.Fatalf("new tool fragment A2: %v", err) } if _, _, err := tail.Append(evA2); err != nil { t.Fatalf("append tool A2: %v", err) } // Complete fragment "a1" on tool channel. _, _, toolOk, err := tail.CompleteFragment("tool", "a1") if err != nil { t.Fatalf("complete fragment tool: %v", err) } if !toolOk { t.Fatal("fragment completion should return true") } // Verify the text channel token is still valid. err = tail.ConfirmRelease(textPR.Token(), ReleaseConfirmation{ReleasedEvents: 3}) if err != nil { t.Fatalf("confirm text release after tool fragment completion: %v", err) } // Text cursor must be 3. if tail.CommittedCursor("text") != 3 { t.Errorf("text cursor want 3, got %d", tail.CommittedCursor("text")) } // Verify the tool fragment epoch was created by CompleteFragment. toolEpochCreated := false for _, r := range tail.epochs { if r.channel == "tool" && r.epochID > 0 { toolEpochCreated = true break } } if !toolEpochCreated { t.Error("tool fragment completion should create an epoch") } }) t.Run("sameChannelOldTokenRejected", func(t *testing.T) { // Prepare a tool channel token, then complete fragment should invalidate it. evT1, err := NewToolCallFragmentEvent("tool2", "t1", "fn", `"a"`, now) if err != nil { t.Fatalf("new tool fragment T1: %v", err) } if _, _, err := tail.Append(evT1); err != nil { t.Fatalf("append tool2 T1: %v", err) } evT2, err := NewToolCallFragmentEvent("tool2", "t1", "fn", `"b"`, now) if err != nil { t.Fatalf("new tool fragment T2: %v", err) } if _, _, err := tail.Append(evT2); err != nil { t.Fatalf("append tool2 T2: %v", err) } // Complete t1 -> creates safe prefix epoch for tool2. _, _, ok, err := tail.CompleteFragment("tool2", "t1") if err != nil { t.Fatalf("complete tool2 t1: %v", err) } if !ok { t.Fatal("first complete should succeed") } // Prepare the safe prefix token. var toolEp uint64 for _, r := range tail.epochs { if r.channel == "tool2" && r.epochID > toolEp { toolEp = r.epochID } } toolPR, err := tail.PrepareRelease(toolEp) if err != nil { t.Fatalf("prepare tool2 release: %v", err) } oldToken := toolPR.Token() if oldToken == "" { t.Fatal("prepared token must not be empty") } // Add another fragment and complete it -> new safe prefix -> invalidates old. evT3, err := NewToolCallFragmentEvent("tool2", "t2", "fn", `"c"`, now) if err != nil { t.Fatalf("new tool call fragment: %v", err) } if _, _, err := tail.Append(evT3); err != nil { t.Fatalf("append tool2 fragment: %v", err) } _, _, _, err = tail.CompleteFragment("tool2", "t2") if err != nil { t.Fatalf("complete fragment tool2 t2: %v", err) } // Old token must now be stale. err = tail.ConfirmRelease(oldToken, ReleaseConfirmation{ReleasedEvents: 2}) if err == nil { t.Error("old tool token after superseding completion must be rejected") } }) t.Run("interleavedFragmentABExpandsSafePrefixAfterConfirm", func(t *testing.T) { // A1, B1, A2 interleaved with toolCallID "a"/"b"/"a". // Complete A → safe=[A1] (B1 incomplete at index 1). Partial confirm 1. // Complete B → all fragments done, safe=[B1,A2] (remaining pending). Full confirm 2. evA1, err := NewToolCallFragmentEvent("interleave", "a", "fn", `"a1"`, now) if err != nil { t.Fatalf("new A1 fragment: %v", err) } if _, sig, err := tail.Append(evA1); err != nil { t.Fatalf("append A1: %v", err) } else if sig != EvidenceTailSignalNone { t.Errorf("append A1 signal want none, got %s", sig) } evB1, err := NewToolCallFragmentEvent("interleave", "b", "fn", `"b1"`, now) if err != nil { t.Fatalf("new B1 fragment: %v", err) } if _, _, err := tail.Append(evB1); err != nil { t.Fatalf("append B1: %v", err) } evA2, err := NewToolCallFragmentEvent("interleave", "a", "fn", `"a2"`, now) if err != nil { t.Fatalf("new A2 fragment: %v", err) } if _, _, err := tail.Append(evA2); err != nil { t.Fatalf("append A2: %v", err) } // Complete A → safe prefix [A1] (B1 incomplete at index 1). epA, _, ok, err := tail.CompleteFragment("interleave", "a") if err != nil { t.Fatalf("complete A: %v", err) } if !ok { t.Fatal("complete A should succeed") } prA, err := tail.PrepareRelease(epA.ID()) if err != nil { t.Fatalf("prepare A release: %v", err) } if len(prA.ReleaseEvents()) != 1 { t.Fatalf("A safe prefix want 1 event, got %d", len(prA.ReleaseEvents())) } // Verify exact content of A1. reA1, err := prA.ReleaseEvents()[0].AsToolCallFragment() if err != nil { t.Fatalf("as tool call fragment [0]: %v", err) } if reA1.ID != "a" || reA1.Name != "fn" || reA1.Arguments != `"a1"` { t.Errorf("A1 fragment want id=a name=fn args=\"a1\", got id=%q name=%q args=%q", reA1.ID, reA1.Name, reA1.Arguments) } err = tail.ConfirmRelease(prA.Token(), ReleaseConfirmation{ReleasedEvents: 1}) if err != nil { t.Fatalf("partial confirm A: %v", err) } // Complete B → all fragments completed, safe prefix = remaining [B1,A2]. epB, _, ok, err := tail.CompleteFragment("interleave", "b") if err != nil { t.Fatalf("complete B: %v", err) } if !ok { t.Fatal("complete B should succeed") } prB, err := tail.PrepareRelease(epB.ID()) if err != nil { t.Fatalf("prepare B release: %v", err) } if len(prB.ReleaseEvents()) != 2 { t.Fatalf("remaining safe prefix want 2 events [B1,A2], got %d", len(prB.ReleaseEvents())) } // Exact public content and order: [B1, A2]. for i := 0; i < 2; i++ { re, err := prB.ReleaseEvents()[i].AsToolCallFragment() if err != nil { t.Fatalf("as tool call fragment[%d]: %v", i, err) } var wantID, wantArgs string if i == 0 { wantID, wantArgs = "b", `"b1"` } else { wantID, wantArgs = "a", `"a2"` } if re.ID != wantID { t.Errorf("release[%d] ID want %q, got %q", i, wantID, re.ID) } if re.Arguments != wantArgs { t.Errorf("release[%d] arguments want %q, got %q", i, wantArgs, re.Arguments) } } // Genuine partial confirm: ReleasedEvents=1 on a 2-event snapshot. err = tail.ConfirmRelease(prB.Token(), ReleaseConfirmation{ReleasedEvents: 1}) if err != nil { t.Fatalf("partial confirm 1 of 2: %v", err) } // Re-prepare: should return only the unconfirmed suffix [A2]. prB2, err := tail.PrepareRelease(epB.ID()) if err != nil { t.Fatalf("re-prepare after partial: %v", err) } if len(prB2.ReleaseEvents()) != 1 { t.Fatalf("re-prepared suffix want 1 event [A2], got %d", len(prB2.ReleaseEvents())) } reSuffix, err := prB2.ReleaseEvents()[0].AsToolCallFragment() if err != nil { t.Fatalf("as tool call fragment suffix: %v", err) } if reSuffix.ID != "a" || reSuffix.Arguments != `"a2"` { t.Errorf("suffix event want ID=a args=\"a2\", got ID=%q args=%q", reSuffix.ID, reSuffix.Arguments) } // Fragment "a" is already completed from the first CompleteFragment // call, so no further completion is needed. Full confirm the remaining // re-prepared suffix. err = tail.ConfirmRelease(prB2.Token(), ReleaseConfirmation{ReleasedEvents: 1}) if err != nil { t.Fatalf("confirm remaining suffix: %v", err) } if tail.CommittedCursor("interleave") != 3 { t.Errorf("cursor want 3 (1+1+1), got %d", tail.CommittedCursor("interleave")) } // Both fragment states "a" and "b" must be cleaned up after the final // confirm. epB.snapStartSeq advanced to 2 by the partial confirm, so // the re-prepared suffix [A2] covers A2 at absolute seq 2; the // cleanup loop's confirmedEndSeq = 2 + 1 = 3 drops every entry with // seq < 3, which removes both A and B. cs := tail.channelState["interleave"] if _, exists := cs.fragmentState["a"]; exists { t.Error("fragmentState[\"a\"] must be removed after A2 confirmed") } if _, exists := cs.fragmentState["b"]; exists { t.Error("fragmentState[\"b\"] must be removed after B1 confirmed") } // No pending entries remain. if len(cs.pendingEntries) != 0 { t.Errorf("pendingEntries want 0, got %d", len(cs.pendingEntries)) } // Re-preparing the first A epoch must fail (consumed). _, err = tail.PrepareRelease(epA.ID()) if err == nil { t.Error("re-prepare consumed epoch must fail") } }) } // buildReplacementTail builds a two-channel rolling tail whose evidence window // is wide enough to hold a full replacement payload without look-behind // trimming, so replacement content assertions observe exact payload identity. func buildReplacementTail(t *testing.T) *EvidenceTail { t.Helper() textKinds := []EventKind{EventKindTextDelta} reqA, err := NewFilterHoldRequirementRollingWithMaxBuffer("a", textKinds, 12, 4096) if err != nil { t.Fatalf("build channel a req: %v", err) } reqB, err := NewFilterHoldRequirementRollingWithMaxBuffer("b", textKinds, 12, 4096) if err != nil { t.Fatalf("build channel b req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, reqA), mustBind(t, reqB)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } return tail } // appendReplacementText appends one text_delta event and returns the epoch it // produced. The caller drives readiness through the 12-rune rolling threshold. func appendReplacementText(t *testing.T, tail *EvidenceTail, channel, text string) EvidenceEpoch { t.Helper() ev, err := NewTextDeltaEvent(channel, text, now) if err != nil { t.Fatalf("new text delta %q: %v", text, err) } epoch, _, err := tail.Append(ev) if err != nil { t.Fatalf("append %q: %v", text, err) } return epoch } func mustReplacementPayload(t *testing.T, channel string, texts ...string) []NormalizedEvent { t.Helper() out := make([]NormalizedEvent, len(texts)) for i, text := range texts { ev, err := NewTextDeltaEvent(channel, text, now) if err != nil { t.Fatalf("new replacement event %q: %v", text, err) } out[i] = ev } return out } func lookBehindTexts(t *testing.T, events []NormalizedEvent) []string { t.Helper() texts := make([]string, 0, len(events)) for _, ev := range events { text, err := ev.AsTextDelta() if err != nil { t.Fatalf("look-behind AsTextDelta: %v", err) } texts = append(texts, text) } return texts } func pendingTexts(t *testing.T, tail *EvidenceTail, channel string) []string { t.Helper() cs := tail.channelState[channel] if cs == nil { return nil } texts := make([]string, 0, len(cs.pendingEntries)) for _, entry := range cs.pendingEntries { text, err := entry.event.AsTextDelta() if err != nil { t.Fatalf("pending AsTextDelta: %v", err) } texts = append(texts, text) } return texts } // TestEvidenceTailReplacementSettlement verifies that the epoch-bound // replacement seam settles only the target epoch's captured range. Validation // fails closed before any state mutation, zero/partial/full sink progress puts // exactly the accepted replacement prefix into look-behind and cursor, and the // same channel's post-snapshot entries plus the unrelated channel's pending // entries and outstanding prepared token all survive untouched. func TestEvidenceTailReplacementSettlement(t *testing.T) { t.Run("fails_closed_before_any_mutation", func(t *testing.T) { tail := buildReplacementTail(t) epoch := appendReplacementText(t, tail, "a", "a-original-1") if epoch.ID() == 0 { t.Fatal("append should produce a ready epoch at the rolling threshold") } terminal, err := NewTerminalEvent("a", now) if err != nil { t.Fatalf("new terminal event: %v", err) } for _, tc := range []struct { name string epochID uint64 events []NormalizedEvent }{ {"unknown_epoch", epoch.ID() + 999, mustReplacementPayload(t, "a", "repl-1")}, {"empty_payload", epoch.ID(), nil}, {"foreign_channel_event", epoch.ID(), mustReplacementPayload(t, "b", "repl-1")}, {"non_releasable_kind", epoch.ID(), []NormalizedEvent{terminal}}, } { if _, err := tail.prepareReplacement(tc.epochID, tc.events); err == nil { t.Fatalf("%s: prepareReplacement error = nil, want fail-closed", tc.name) } } // Nothing may have been consumed, released, or reserved. if got := pendingTexts(t, tail, "a"); !reflect.DeepEqual(got, []string{"a-original-1"}) { t.Fatalf("pending on a = %v, want [a-original-1]", got) } if got := tail.CommittedCursor("a"); got != 0 { t.Fatalf("committed cursor on a = %d, want 0", got) } if len(tail.replacements) != 0 { t.Fatalf("prepared replacements = %d, want 0", len(tail.replacements)) } if _, err := tail.PrepareRelease(epoch.ID()); err != nil { t.Fatalf("PrepareRelease after fail-closed prepare: %v", err) } }) for _, tc := range []struct { name string releasedEvents int wantLookBehind []string wantCursorAfter int }{ {"zero_progress", 0, []string{}, 0}, {"partial_progress", 1, []string{"repl-1"}, 1}, {"full_progress", 2, []string{"repl-1", "repl-2"}, 2}, } { t.Run(tc.name, func(t *testing.T) { tail := buildReplacementTail(t) // Unrelated channel: one pending entry and an outstanding prepared // release token that must stay confirmable across the replacement. epochB := appendReplacementText(t, tail, "b", "b-original-1") preparedB, err := tail.PrepareRelease(epochB.ID()) if err != nil { t.Fatalf("PrepareRelease on b: %v", err) } // Target channel: two originals captured by the target epoch, then a // post-snapshot entry that a later overlapping epoch also captures. appendReplacementText(t, tail, "a", "a-original-1") target := appendReplacementText(t, tail, "a", "a-original-2") laterEpoch := appendReplacementText(t, tail, "a", "a-post") prepared, err := tail.prepareReplacement( target.ID(), mustReplacementPayload(t, "a", "repl-1", "repl-2"), ) if err != nil { t.Fatalf("prepareReplacement: %v", err) } if len(prepared.ReleaseEvents()) != 2 { t.Fatalf("prepared replacement events = %d, want 2", len(prepared.ReleaseEvents())) } // Prepare alone must not settle anything. if got := pendingTexts(t, tail, "a"); len(got) != 3 { t.Fatalf("pending on a after prepare = %v, want 3 entries", got) } if err := tail.confirmReplacement( prepared.Token(), ReleaseConfirmation{ReleasedEvents: tc.releasedEvents}, ); err != nil { t.Fatalf("confirmReplacement: %v", err) } // The whole target range is consumed even at zero progress, and the // post-snapshot entry survives. if got := pendingTexts(t, tail, "a"); !reflect.DeepEqual(got, []string{"a-post"}) { t.Fatalf("pending on a = %v, want [a-post]", got) } if got := tail.CommittedCursor("a"); got != tc.wantCursorAfter { t.Fatalf("committed cursor on a = %d, want %d", got, tc.wantCursorAfter) } if got := lookBehindTexts(t, tail.EffectiveLookBehind("a")); !reflect.DeepEqual(got, tc.wantLookBehind) { t.Fatalf("look-behind on a = %v, want %v", got, tc.wantLookBehind) } // The token is single-use and the target epoch is settled. if err := tail.confirmReplacement( prepared.Token(), ReleaseConfirmation{ReleasedEvents: tc.releasedEvents}, ); err == nil { t.Fatal("re-confirming a replacement token must fail") } if _, err := tail.PrepareRelease(target.ID()); err == nil { t.Fatal("preparing the settled target epoch must fail") } // The later same-channel epoch overlapped the consumed range. if _, err := tail.PrepareRelease(laterEpoch.ID()); err == nil { t.Fatal("preparing an overlapping stale epoch must fail") } // The unrelated channel kept its pending entry, its prepared token, // and its own cursor accounting. if got := pendingTexts(t, tail, "b"); !reflect.DeepEqual(got, []string{"b-original-1"}) { t.Fatalf("pending on b = %v, want [b-original-1]", got) } if err := tail.ConfirmRelease(preparedB.Token(), ReleaseConfirmation{ReleasedEvents: 1}); err != nil { t.Fatalf("ConfirmRelease on b after replacement: %v", err) } if got := tail.CommittedCursor("b"); got != 1 { t.Fatalf("committed cursor on b = %d, want 1", got) } // The target channel still accepts and releases new evidence. next := appendReplacementText(t, tail, "a", "a-tail") preparedNext, err := tail.PrepareRelease(next.ID()) if err != nil { t.Fatalf("PrepareRelease on a after replacement: %v", err) } if len(preparedNext.ReleaseEvents()) != 2 { t.Fatalf("post-replacement snapshot = %d events, want 2 [a-post,a-tail]", len(preparedNext.ReleaseEvents())) } if err := tail.ConfirmRelease(preparedNext.Token(), ReleaseConfirmation{ReleasedEvents: 2}); err != nil { t.Fatalf("ConfirmRelease on a after replacement: %v", err) } if got := tail.CommittedCursor("a"); got != tc.wantCursorAfter+2 { t.Fatalf("committed cursor on a = %d, want %d", got, tc.wantCursorAfter+2) } }) } t.Run("stale_consumed_non_empty_range_rejected", func(t *testing.T) { tail := buildReplacementTail(t) // Old epoch: two originals captured on channel "a". appendReplacementText(t, tail, "a", "a-old-1") oldEpoch := appendReplacementText(t, tail, "a", "a-old-2") // Newer overlapping epoch on the same channel, capturing a-post. newerEpoch := appendReplacementText(t, tail, "a", "a-newer") // Fully confirm the newer epoch so the original range is consumed. preparedNewer, err := tail.PrepareRelease(newerEpoch.ID()) if err != nil { t.Fatalf("PrepareRelease on newer: %v", err) } if err := tail.ConfirmRelease(preparedNewer.Token(), ReleaseConfirmation{ReleasedEvents: 2}); err != nil { t.Fatalf("ConfirmRelease newer: %v", err) } // The old epoch's captured range is now consumed by the newer release. cs := tail.channelState["a"] oldRec := tail.epochs[oldEpoch.ID()] if got := cs.rangeEntryCount(oldRec.snapStartSeq, oldRec.snapEndSeq); got != 0 { t.Fatalf("old range remaining entries = %d, want 0", got) } // prepareReplacement must fail-closed with snapshot mismatch. _, err = tail.prepareReplacement( oldEpoch.ID(), mustReplacementPayload(t, "a", "repl-stale-1", "repl-stale-2"), ) if err != errPreparedSnapshotMismatch { t.Fatalf("prepareReplacement stale = %v, want errPreparedSnapshotMismatch", err) } // prepare must not have created a replacement record. if len(tail.replacements) != 0 { t.Fatalf("prepared replacements = %d, want 0", len(tail.replacements)) } // Pending must still contain only the post-snapshot entry. if got := pendingTexts(t, tail, "a"); !reflect.DeepEqual(got, []string{"a-newer"}) { t.Fatalf("pending on a = %v, want [a-newer]", got) } if len(tail.replacements) != 0 { t.Fatalf("prepared replacements = %d, want 0", len(tail.replacements)) } }) t.Run("empty_range_terminal_trigger_accepts_replacement", func(t *testing.T) { textKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} req, err := NewFilterHoldRequirementTerminalGate("ch-empty", textKinds, EventKindTerminal) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Append a terminal trigger with no prior pending entries. The // trigger creates an epoch with snapStartSeq == snapEndSeq == 0, // i.e. an exact [n,n) empty captured range. termEv, err := NewTerminalEvent("ch-empty", now) if err != nil { t.Fatalf("new terminal event: %v", err) } ep, _, err := tail.Append(termEv) if err != nil { t.Fatalf("append terminal: %v", err) } epRec := tail.epochs[ep.ID()] if epRec.snapEndSeq-epRec.snapStartSeq != 0 { t.Fatalf("terminal trigger epoch range width = %d, want 0", epRec.snapEndSeq-epRec.snapStartSeq) } // A non-empty replacement payload must be accepted for the empty range. prepared, err := tail.prepareReplacement( ep.ID(), mustReplacementPayload(t, "ch-empty", "repl-1", "repl-2"), ) if err != nil { t.Fatalf("prepareReplacement empty range: %v", err) } if len(prepared.ReleaseEvents()) != 2 { t.Fatalf("prepared events = %d, want 2", len(prepared.ReleaseEvents())) } // Confirm with zero progress: range is consumed but nothing enters // look-behind since the sink accepted none. if err := tail.confirmReplacement(prepared.Token(), ReleaseConfirmation{ReleasedEvents: 0}); err != nil { t.Fatalf("confirmReplacement empty range: %v", err) } // The epoch must be consumed. if !tail.epochs[ep.ID()].consumed { t.Fatal("epoch must be consumed after confirm") } if len(tail.replacements) != 0 { t.Fatalf("replacements = %d, want 0", len(tail.replacements)) } }) t.Run("empty_range_after_fully_released_prefix_accepts_replacement", func(t *testing.T) { // Regression: firstPendingSequence() must return cs.nextSequence (not 0) // when pending is empty so that a terminal trigger after fully releasing // a prefix creates an epoch with the current monotonic [n,n) boundary. tail := buildReplacementTail(t) // Step 1: Append two text events on channel "a" to cross the 12-rune // rolling threshold and create a ready epoch. The accumulated text // advances nextSequence past 0. appendReplacementText(t, tail, "a", "a-prefix-12") original := appendReplacementText(t, tail, "a", "a-prefix-34") cs := tail.channelState["a"] if cs.nextSequence != 2 { t.Fatalf("nextSequence after two appends = %d, want 2", cs.nextSequence) } // Step 2: Fully release the original epoch so pending becomes empty. prepared, err := tail.PrepareRelease(original.ID()) if err != nil { t.Fatalf("PrepareRelease original: %v", err) } if err := tail.ConfirmRelease(prepared.Token(), ReleaseConfirmation{ReleasedEvents: 2}); err != nil { t.Fatalf("ConfirmRelease original: %v", err) } // Verify pending is empty and nextSequence > 0. cs = tail.channelState["a"] if len(cs.pendingEntries) != 0 { t.Fatalf("pending after full release = %d, want 0", len(cs.pendingEntries)) } if cs.nextSequence != 2 { t.Fatalf("nextSequence after release = %d, want 2", cs.nextSequence) } // Step 3: Append a terminal trigger. With the fix, firstPendingSequence() // must return 2 (cs.nextSequence), so the epoch range becomes [2,2). termEv, err := NewTerminalEvent("a", now) if err != nil { t.Fatalf("new terminal event: %v", err) } ep, _, err := tail.Append(termEv) if err != nil { t.Fatalf("append terminal after full release: %v", err) } epRec := tail.epochs[ep.ID()] rangeWidth := epRec.snapEndSeq - epRec.snapStartSeq if rangeWidth != 0 { t.Fatalf("terminal epoch range width = %d, want 0", rangeWidth) } if epRec.snapStartSeq != 2 { t.Fatalf("terminal epoch snapStartSeq = %d, want 2 (current monotonic)", epRec.snapStartSeq) } if epRec.snapEndSeq != 2 { t.Fatalf("terminal epoch snapEndSeq = %d, want 2 (current monotonic)", epRec.snapEndSeq) } // Step 4: Replacement must be accepted for this non-zero empty epoch. preparedRepl, err := tail.prepareReplacement( ep.ID(), mustReplacementPayload(t, "a", "repl-1", "repl-2"), ) if err != nil { t.Fatalf("prepareReplacement non-zero empty epoch: %v", err) } if len(preparedRepl.ReleaseEvents()) != 2 { t.Fatalf("prepared events = %d, want 2", len(preparedRepl.ReleaseEvents())) } // Confirm with zero progress: the empty range is consumed but nothing // enters look-behind. if err := tail.confirmReplacement(preparedRepl.Token(), ReleaseConfirmation{ReleasedEvents: 0}); err != nil { t.Fatalf("confirmReplacement non-zero empty: %v", err) } // The epoch must be consumed. if !tail.epochs[ep.ID()].consumed { t.Fatal("epoch must be consumed after confirm") } if len(tail.replacements) != 0 { t.Fatalf("replacements = %d, want 0", len(tail.replacements)) } // Pending must still be empty (terminal event excluded from pending). cs = tail.channelState["a"] if len(cs.pendingEntries) != 0 { t.Fatalf("pending after terminal+replacement = %d, want 0", len(cs.pendingEntries)) } // Committed cursor must reflect the original release, not the empty replacement. if got := tail.CommittedCursor("a"); got != 2 { t.Fatalf("committed cursor on a = %d, want 2", got) } // nextSequence must remain unchanged after the terminal trigger. cs = tail.channelState["a"] if cs.nextSequence != 2 { t.Fatalf("nextSequence after terminal = %d, want 2", cs.nextSequence) } }) } // terminal and provider-error triggers create an exact safe-prefix ready // epoch that excludes the control event itself, and that the signal is // EvidenceTailSignalTrigger. Outstanding tokens after the trigger are // properly invalidated. func TestEvidenceTailTerminalGateConfiguredTrigger(t *testing.T) { allKinds := []EventKind{EventKindTextDelta, EventKindReasoningDelta} t.Run("configuredTerminalTriggerExcludesControlEvent", func(t *testing.T) { req, err := NewFilterHoldRequirementTerminalGate("ch", allKinds, EventKindTerminal) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Add 3 text events (data, not trigger). for i := 0; i < 3; i++ { ev, _ := NewTextDeltaEvent("ch", fmt.Sprintf("data-%d", i), now) tail.Append(ev) } // Append terminal trigger. evTerm, _ := NewTerminalEvent("ch", now) ep, signal, err := tail.Append(evTerm) if err != nil { t.Fatalf("append terminal: %v", err) } if signal != EvidenceTailSignalTrigger { t.Errorf("signal want trigger, got %s", signal) } if ep.ID() == 0 { t.Fatal("terminal trigger must produce an epoch") } if !ep.Triggered() { t.Error("epoch must be marked triggered") } // Prepare release: snapshot must contain only the 3 data events, // NOT the terminal event. pr, err := tail.PrepareRelease(ep.ID()) if err != nil { t.Fatalf("prepare: %v", err) } if len(pr.ReleaseEvents()) != 3 { t.Errorf("snapshot want 3 data events, got %d", len(pr.ReleaseEvents())) } // Full confirm. if err := tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 3}); err != nil { t.Fatalf("confirm: %v", err) } // Pending must be empty (terminal event excluded from pending). cs := tail.channelState["ch"] if len(cs.pendingEntries) != 0 { t.Errorf("pending must be 0 after terminal release, got %d", len(cs.pendingEntries)) } }) t.Run("configuredProviderErrorTriggerExcludesControlEvent", func(t *testing.T) { req, err := NewFilterHoldRequirementTerminalGate("ch", allKinds, EventKindProviderError) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Add 2 data events. for i := 0; i < 2; i++ { ev, _ := NewTextDeltaEvent("ch", fmt.Sprintf("data-%d", i), now) tail.Append(ev) } // Append provider-error trigger. evPE, _ := NewProviderErrorEvent("ch", mustDesc(t, "provider_error", "err", "provider_failure", ""), mustCauses(t, "stage", "err", "consumer", "filter", "rule"), now) ep, signal, err := tail.Append(evPE) if err != nil { t.Fatalf("append provider-error: %v", err) } if signal != EvidenceTailSignalTrigger { t.Errorf("signal want trigger, got %s", signal) } if ep.ID() == 0 { t.Fatal("provider-error trigger must produce an epoch") } pr, err := tail.PrepareRelease(ep.ID()) if err != nil { t.Fatalf("prepare: %v", err) } // Only 2 data events, not the provider-error control event. if len(pr.ReleaseEvents()) != 2 { t.Errorf("snapshot want 2 data events, got %d", len(pr.ReleaseEvents())) } // Verify none of the release events is a provider-error kind. for _, re := range pr.ReleaseEvents() { if re.Kind() == EventKindProviderError { t.Error("control event must not appear in release snapshot") } } }) t.Run("terminalTriggerWithNoPendingData", func(t *testing.T) { req, err := NewFilterHoldRequirementTerminalGate("ch", allKinds, EventKindTerminal) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // No data events, just terminal trigger. evTerm, _ := NewTerminalEvent("ch", now) ep, signal, err := tail.Append(evTerm) if err != nil { t.Fatalf("append terminal: %v", err) } if signal != EvidenceTailSignalTrigger { t.Errorf("signal want trigger, got %s", signal) } if ep.ID() == 0 { t.Fatal("terminal trigger must produce epoch even with no data") } // Prepare should return empty snapshot. _, err = tail.PrepareRelease(ep.ID()) if err == nil { t.Error("empty snapshot should be rejected on prepare") } }) t.Run("subsequentAppendAfterTerminal", func(t *testing.T) { req, err := NewFilterHoldRequirementTerminalGate("ch", allKinds, EventKindTerminal) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Add data, terminal trigger, confirm. for i := 0; i < 2; i++ { ev, _ := NewTextDeltaEvent("ch", "data", now) tail.Append(ev) } evTerm, _ := NewTerminalEvent("ch", now) ep, _, _ := tail.Append(evTerm) pr, _ := tail.PrepareRelease(ep.ID()) tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 2}) // Now append more data and a new terminal trigger. for i := 0; i < 2; i++ { ev, _ := NewTextDeltaEvent("ch", fmt.Sprintf("more-%d", i), now) tail.Append(ev) } evTerm2, _ := NewTerminalEvent("ch", now) ep2, signal2, err := tail.Append(evTerm2) if err != nil { t.Fatalf("append second terminal: %v", err) } if signal2 != EvidenceTailSignalTrigger { t.Errorf("second signal want trigger, got %s", signal2) } pr2, _ := tail.PrepareRelease(ep2.ID()) if err != nil { t.Fatalf("prepare second: %v", err) } if len(pr2.ReleaseEvents()) != 2 { t.Errorf("second snapshot want 2 data events, got %d", len(pr2.ReleaseEvents())) } }) } // --- test helper functions --- func mustDesc(t *testing.T, errorType, code, message, param string) ExternalDescriptor { t.Helper() d, err := NewExternalDescriptor(errorType, code, message, param) if err != nil { t.Fatalf("external descriptor: %v", err) } return d } func mustCauses(t *testing.T, stage, code, consumer, filter, ruleID string) FailureCauseChain { t.Helper() c, err := NewFailureCause(stage, code, consumer, filter, ruleID) if err != nil { t.Fatalf("failure cause: %v", err) } chain, err := NewFailureCauseChain([]FailureCause{c}) if err != nil { t.Fatalf("failure cause chain: %v", err) } return chain } func mustBind(t *testing.T, req FilterHoldRequirement) FilterHoldBinding { t.Helper() b, err := NewFilterHoldBinding("f-"+req.Channel(), req, FilterEnforcementBlocking) if err != nil { t.Fatalf("bind: %v", err) } return b } func buildThresholdTail(t *testing.T, kinds []EventKind, threshold int) *EvidenceTail { t.Helper() req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", kinds, threshold, 4096) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Fill to threshold. for i := 0; i < 3; i++ { ev, _ := NewTextDeltaEvent("ch", repeatRune('a', threshold/3+1), now) tail.Append(ev) } return tail } func findLatestEpoch(t *testing.T, tail *EvidenceTail) uint64 { t.Helper() var ep uint64 for _, rec := range tail.epochs { if rec.epochID > ep { ep = rec.epochID } } if ep == 0 { t.Fatal("no epoch found") } return ep } func repeatRune(r rune, n int) string { b := make([]rune, n) for i := range b { b[i] = r } return string(b) } // TestEvidenceTailPreservesExactKoreanSuffixContent verifies that // 2/200/500-rune oversized Korean events keep exact trailing content through // release/confirm. The table drives each bound with a single Korean event // larger than the bound, a release+confirm cycle, and an ExactLookBehind // assertion on the trailing bound runes. func TestEvidenceTailPreservesExactKoreanSuffixContent(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} type koreanCase struct { name string bound int runes int } cases := []koreanCase{ {"twoRuneBound", 2, 50}, {"twoHundredRuneBound", 200, 500}, {"fiveHundredRuneBound", 500, 1500}, } // distinctHangul produces n position-distinct valid Hangul syllable runes // so that a wrong suffix offset always produces a different rune, not a // same-rune false match. The largest case remains within the 11,172-code // point Hangul syllables block. distinctHangul := func(n int) string { b := make([]rune, n) for i := 0; i < n; i++ { b[i] = rune(0xAC00 + i%11172) } return string(b) } for _, tc := range cases { tc := tc t.Run(tc.name, func(t *testing.T) { req, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch-kr", textKinds, tc.bound, 2*tc.runes) if err != nil { t.Fatalf("build req: %v", err) } plan, err := NewEvidencePlan([]FilterHoldBinding{mustBind(t, req)}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } // Single Korean event with tc.runes unique Hangul syllables. koreanText := distinctHangul(tc.runes) ev, err := NewTextDeltaEvent("ch-kr", koreanText, now) if err != nil { t.Fatalf("new Korean text delta: %v", err) } if _, sig, err := tail.Append(ev); err != nil { t.Fatalf("append Korean event: %v", err) } else if sig != EvidenceTailSignalNone && sig != EvidenceTailSignalThreshold { t.Errorf("append Korean signal want none or threshold, got %s", sig) } var ep uint64 for _, r := range tail.epochs { if r.epochID > ep { ep = r.epochID } } pr, err := tail.PrepareRelease(ep) if err != nil { t.Fatalf("prepare release: %v", err) } if len(pr.ReleaseEvents()) != 1 { t.Fatalf("release want 1 event, got %d", len(pr.ReleaseEvents())) } releaseText, err := pr.ReleaseEvents()[0].AsTextDelta() if err != nil { t.Fatalf("as text delta (release): %v", err) } if len([]rune(releaseText)) != tc.runes { t.Errorf("release text want %d runes, got %d", tc.runes, len([]rune(releaseText))) } // Full confirm. if err := tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 1}); err != nil { t.Fatalf("confirm: %v", err) } // EffectiveLookBehind must be a deep copy of at most tc.bound runes. lb := tail.EffectiveLookBehind("ch-kr") if len(lb) != 1 { t.Fatalf("lookBehind want 1 event, got %d", len(lb)) } lbText, err := lb[0].AsTextDelta() if err != nil { t.Fatalf("as text delta (lookBehind): %v", err) } if len([]rune(lbText)) != tc.bound { t.Errorf("lookBehind trailing want %d runes, got %d", tc.bound, len([]rune(lbText))) } // Exact trailing content: the last tc.bound runes must be the // suffix of koreanText. With distinct-rune fixture, a wrong // whole-rune offset in the trim logic flips to a different rune. allR := []rune(koreanText) wantSuffix := string(allR[len(allR)-tc.bound:]) if lbText != wantSuffix { t.Errorf("lookBehind text want exact suffix of %d chars, got %q", tc.bound, lbText) } // Exact release text match: release must equal the full koreanText. if releaseText != koreanText { t.Errorf("release text want exact %q, got %q", koreanText, releaseText) } // Mutation-isolated requery: mutate the first returned event // in-place and verify that re-querying EffectiveLookBehind still // returns the untouched exact suffix, proving the accessor // returns a deep copy. mutatedOriginal := lbText lb[0] = NormalizedEvent{} lbRequery := tail.EffectiveLookBehind("ch-kr") if len(lbRequery) != 1 { t.Fatalf("requery lookBehind want 1 event, got %d", len(lbRequery)) } lbRequeryText, err := lbRequery[0].AsTextDelta() if err != nil { t.Fatalf("as text delta (requery): %v", err) } if lbRequeryText != mutatedOriginal { t.Errorf("defensive copy requery want exact %q after mutation, got %q", mutatedOriginal, lbRequeryText) } // Cursor must be 1. if tail.CommittedCursor("ch-kr") != 1 { t.Errorf("cursor want 1, got %d", tail.CommittedCursor("ch-kr")) } }) } } type mockTailFilter struct { id string req FilterHoldRequirement } func (m *mockTailFilter) ID() string { return m.id } func (m *mockTailFilter) Applies(c FilterContext) bool { return true } func (m *mockTailFilter) HoldRequirement(c FilterContext) FilterHoldRequirement { return m.req } func (m *mockTailFilter) Evaluate(ctx context.Context, fc FilterContext, batch EvidenceBatch) (FilterDecision, error) { ev, _ := NewSanitizedEvidence(EventKindTextDelta, "ch", "rule", "code", [32]byte{1}, 0, 0, FilterOutcomeKindEvaluated, now) return NewFilterDecision(FilterDecisionKindPass, "consumer.id", m.id, "rule.id", ev, nil) } func TestMultiRequirementPolicyHookScheduleAndLifecycle(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} toolKinds := []EventKind{EventKindToolCallFragment} t.Run("allFourModesScheduleUnionAndBounds", func(t *testing.T) { roll, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 100, 2000) if err != nil { t.Fatalf("roll: %v", err) } term, err := NewFilterHoldRequirementTerminalGateWithMaxBuffer("ch", textKinds, EventKindTerminal, 1500) if err != nil { t.Fatalf("term: %v", err) } frag, err := NewFilterHoldRequirementFragmentGateWithMaxBuffer("ch", toolKinds, EventKindToolCallFragment, 1800) if err != nil { t.Fatalf("frag: %v", err) } none, err := NewFilterHoldRequirementNone("ch", textKinds) if err != nil { t.Fatalf("none: %v", err) } b1, _ := NewFilterHoldBinding("f1_roll", roll, FilterEnforcementBlocking) b2, _ := NewFilterHoldBinding("f2_term", term, FilterEnforcementBlocking) b3, _ := NewFilterHoldBinding("f3_frag", frag, FilterEnforcementBlocking) b4, _ := NewFilterHoldBinding("f4_none", none, FilterEnforcementObserveOnly) plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2, b3, b4}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } bindings := plan.BindingsForChannel("ch") if len(bindings) != 4 { t.Fatalf("bindings count want 4, got %d", len(bindings)) } subscribeKinds := plan.SubscribeKinds("ch") if len(subscribeKinds) != 2 { t.Errorf("subscribeKinds len want 2 (textDelta, toolCallFragment), got %d", len(subscribeKinds)) } repReq, ok := plan.BlockingRequirementFor("ch") if !ok { t.Fatal("expected blocking requirement for ch") } if repReq.MaxBufferRunes() != 1500 { t.Errorf("min hard bound want 1500, got %d", repReq.MaxBufferRunes()) } rollReq, ok := plan.plan.allBindings["f1_roll"] if !ok || rollReq.Requirement().EvidenceRunes() != 100 { t.Errorf("f1_roll evidence runes want 100") } }) t.Run("incompatibleAggregateRejectedWhenRollingExceedsHardBound", func(t *testing.T) { roll, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 2500, 3000) if err != nil { t.Fatalf("roll: %v", err) } term, err := NewFilterHoldRequirementTerminalGateWithMaxBuffer("ch", textKinds, EventKindTerminal, 2000) if err != nil { t.Fatalf("term: %v", err) } b1, _ := NewFilterHoldBinding("f1_roll", roll, FilterEnforcementBlocking) b2, _ := NewFilterHoldBinding("f2_term", term, FilterEnforcementBlocking) _, err = NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err == nil { t.Error("rolling evidence > hard bound should be rejected as incompatible aggregate") } }) t.Run("earlyRollingEpochWithDeferredTerminalGateDoesNotReleaseEarly", func(t *testing.T) { roll, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 10, 4096) if err != nil { t.Fatalf("roll: %v", err) } term, err := NewFilterHoldRequirementTerminalGate("ch", textKinds, EventKindTerminal) if err != nil { t.Fatalf("term: %v", err) } fRoll := &mockTailFilter{id: "f1_roll", req: roll} fTerm := &mockTailFilter{id: "f2_term", req: term} fRegRoll, _ := NewFilterRegistration(fRoll, "cap1", true, FilterEnforcementBlocking, time.Second, 10) fRegTerm, _ := NewFilterRegistration(fTerm, "cap2", true, FilterEnforcementBlocking, time.Second, 20) snap, err := NewFilterRegistrySnapshot("gen1", []FilterRegistration{fRegRoll, fRegTerm}, nil) if err != nil { t.Fatalf("snapshot: %v", err) } reqCtx, err := NewRequestFilterContext("gen1", "att1", "env", "ep", "fam", "s1", CommitStateTransportUncommitted, false, false, "corr") if err != nil { t.Fatalf("reqCtx: %v", err) } reqSnap, err := snap.BeginRequest(reqCtx) if err != nil { t.Fatalf("reqSnap: %v", err) } target, _ := NewAttemptTarget("mg", "m", "p", "path", []string{"cap1", "cap2"}) resolved, err := reqSnap.ResolveAttempt(target) if err != nil { t.Fatalf("resolveAttempt: %v", err) } plan, err := NewEvidencePlanFromResolvedFilters(resolved) if err != nil { t.Fatalf("plan from resolved: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } ev, _ := NewTextDeltaEvent("ch", "012345678901234", now) epoch, signal, err := tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } if signal != EvidenceTailSignalThreshold { t.Fatalf("signal want threshold, got %s", signal) } epochFilters, err := plan.BindEpochFilters(epoch, resolved) if err != nil { t.Fatalf("bind epoch filters: %v", err) } if len(epochFilters) != 2 { t.Fatalf("epochFilters count want 2, got %d", len(epochFilters)) } if !epochFilters[0].EvaluatedForEpoch() { t.Errorf("f1_roll should be evaluated for epoch") } if epochFilters[1].EvaluatedForEpoch() { t.Errorf("f2_term should NOT be evaluated for epoch (deferred)") } outcomeTerm := epochFilters[1].NormalizeOutcome() if outcomeTerm.Kind() != FilterOutcomeKindDeferredByRequirement { t.Errorf("f2_term outcome want deferred_by_requirement, got %s", outcomeTerm.Kind()) } termEv, _ := NewTerminalEvent("ch", now) termEpoch, termSignal, err := tail.Append(termEv) if err != nil { t.Fatalf("append terminal: %v", err) } if termSignal != EvidenceTailSignalTrigger { t.Fatalf("terminal signal want trigger, got %s", termSignal) } prepared, err := tail.PrepareRelease(termEpoch.ID()) if err != nil { t.Fatalf("prepare release: %v", err) } if len(prepared.ReleaseEvents()) != 1 { t.Errorf("release events count want 1, got %d", len(prepared.ReleaseEvents())) } err = tail.ConfirmRelease(prepared.Token(), ReleaseConfirmation{ReleasedEvents: 1}) if err != nil { t.Fatalf("confirm release: %v", err) } if tail.CommittedCursor("ch") != 1 { t.Errorf("committed cursor want 1, got %d", tail.CommittedCursor("ch")) } }) } func TestMultiRequirementPolicyHookApplicabilityMatrix(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} toolKinds := []EventKind{EventKindToolCallFragment} t.Run("terminalPendingPresentAndReady", func(t *testing.T) { roll, _ := NewFilterHoldRequirementRolling("ch", textKinds, 100) term, _ := NewFilterHoldRequirementTerminalGate("ch", textKinds, EventKindTerminal) bRoll, _ := NewFilterHoldBinding("f_roll", roll, FilterEnforcementBlocking) bTerm, _ := NewFilterHoldBinding("f_term", term, FilterEnforcementBlocking) plan, _ := NewEvidencePlan([]FilterHoldBinding{bRoll, bTerm}) tail, _ := NewEvidenceTail(plan) evText, _ := NewTextDeltaEvent("ch", "012345678901234", now) _, sig, err := tail.Append(evText) if err != nil { t.Fatalf("append text: %v", err) } if sig != EvidenceTailSignalNone { t.Fatalf("expected signal none for text append, got %s", sig) } evTerm, _ := NewTerminalEvent("ch", now) epoch, sig, err := tail.Append(evTerm) if err != nil { t.Fatalf("append terminal: %v", err) } if sig != EvidenceTailSignalTrigger { t.Fatalf("expected signal trigger for terminal append, got %s", sig) } appRoll, ok := epoch.ApplicabilityFor("f_roll") if !ok || !appRoll.SubscribedEventPresent() || !appRoll.TriggerReady() { t.Errorf("f_roll app want (true, true), got (%v, %v)", appRoll.SubscribedEventPresent(), appRoll.TriggerReady()) } appTerm, ok := epoch.ApplicabilityFor("f_term") if !ok || !appTerm.SubscribedEventPresent() || !appTerm.TriggerReady() { t.Errorf("f_term app want (true, true), got (%v, %v)", appTerm.SubscribedEventPresent(), appTerm.TriggerReady()) } }) t.Run("blockingAndObserveOverlapBufferedOnce", func(t *testing.T) { roll, _ := NewFilterHoldRequirementRolling("ch", textKinds, 100) bBlock, _ := NewFilterHoldBinding("f_block", roll, FilterEnforcementBlocking) bObs, _ := NewFilterHoldBinding("f_obs", roll, FilterEnforcementObserveOnly) plan, _ := NewEvidencePlan([]FilterHoldBinding{bBlock, bObs}) tail, _ := NewEvidenceTail(plan) evText, _ := NewTextDeltaEvent("ch", "hello world", now) epoch, sig, err := tail.Append(evText) if err != nil { t.Fatalf("append text: %v", err) } if sig != EvidenceTailSignalNone { t.Fatalf("expected signal none, got %s", sig) } appBlock, ok := epoch.ApplicabilityFor("f_block") if !ok || !appBlock.SubscribedEventPresent() { t.Errorf("f_block subscribed want true, got %v", appBlock.SubscribedEventPresent()) } appObs, ok := epoch.ApplicabilityFor("f_obs") if !ok || !appObs.SubscribedEventPresent() { t.Errorf("f_obs subscribed want true, got %v", appObs.SubscribedEventPresent()) } cs := tail.channelState["ch"] if cs == nil || cs.pendingRunes != 11 { t.Errorf("pending runes want 11 (buffered once), got %v", cs) } }) t.Run("textRollingDoesNotCountToolRunes", func(t *testing.T) { rollText, _ := NewFilterHoldRequirementRolling("ch", textKinds, 20) fragTool, _ := NewFilterHoldRequirementFragmentGate("ch", toolKinds, EventKindToolCallFragment) bText, _ := NewFilterHoldBinding("f_text", rollText, FilterEnforcementBlocking) bTool, _ := NewFilterHoldBinding("f_tool", fragTool, FilterEnforcementBlocking) plan, _ := NewEvidencePlan([]FilterHoldBinding{bText, bTool}) tail, _ := NewEvidenceTail(plan) evFrag, _ := NewToolCallFragmentEvent("ch", "call-1", "func", `{"query":"test_search_query_tool_call"}`, now) epoch, sig, err := tail.Append(evFrag) if err != nil { t.Fatalf("append tool fragment: %v", err) } if sig != EvidenceTailSignalNone { t.Fatalf("tool fragment should not trigger text rolling, got signal %s", sig) } appText, ok := epoch.ApplicabilityFor("f_text") if !ok { t.Fatalf("missing f_text applicability") } if appText.SubscribedEventPresent() || appText.TriggerReady() { t.Errorf("f_text app want (false, false), got (%v, %v)", appText.SubscribedEventPresent(), appText.TriggerReady()) } }) t.Run("hardBoundFinalReadiness", func(t *testing.T) { roll, _ := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 50, 100) b1, _ := NewFilterHoldBinding("f1", roll, FilterEnforcementBlocking) b2, _ := NewFilterHoldBinding("f2", roll, FilterEnforcementObserveOnly) plan, _ := NewEvidencePlan([]FilterHoldBinding{b1, b2}) tail, _ := NewEvidenceTail(plan) largeText := make([]byte, 120) for i := range largeText { largeText[i] = 'a' } evOverflow, _ := NewTextDeltaEvent("ch", string(largeText), now) epoch, sig, err := tail.Append(evOverflow) if err != nil { t.Fatalf("append overflow: %v", err) } if sig != EvidenceTailSignalBufferOverflow { t.Fatalf("expected signal overflow, got %s", sig) } app1, _ := epoch.ApplicabilityFor("f1") if !app1.TriggerReady() { t.Errorf("blocking filter f1 triggerReady want true on hard bound") } app2, _ := epoch.ApplicabilityFor("f2") if app2.TriggerReady() { t.Errorf("observe filter f2 triggerReady want false on hard bound") } }) t.Run("providerErrorNoneReadiness", func(t *testing.T) { errKinds := []EventKind{EventKindProviderError} noneReq, _ := NewFilterHoldRequirementNone("ch", errKinds) bNone, _ := NewFilterHoldBinding("f_none", noneReq, FilterEnforcementObserveOnly) plan, _ := NewEvidencePlan([]FilterHoldBinding{bNone}) tail, _ := NewEvidenceTail(plan) evErr, _ := NewProviderErrorEvent("ch", mustDesc(t, "err_500", "500", "internal", ""), mustCauses(t, "stage", "err", "prov", "filter", "rule"), now) epoch, sig, err := tail.Append(evErr) if err != nil { t.Fatalf("append provider error: %v", err) } if sig != EvidenceTailSignalReady { t.Fatalf("expected signal ready for provider error pass-through, got %s", sig) } app, _ := epoch.ApplicabilityFor("f_none") if !app.SubscribedEventPresent() || !app.TriggerReady() { t.Errorf("f_none app want (true, true), got (%v, %v)", app.SubscribedEventPresent(), app.TriggerReady()) } }) } func TestMultiRequirementPolicyHookStorageEnvelope(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} t.Run("storageEnvelopeAccessorsAndMinMaxCombination", func(t *testing.T) { roll, _ := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 80, 2000) term, _ := NewFilterHoldRequirementTerminalGateWithMaxBuffer("ch", textKinds, EventKindTerminal, 1500) b1, _ := NewFilterHoldBinding("f1_roll", roll, FilterEnforcementBlocking) b2, _ := NewFilterHoldBinding("f2_term", term, FilterEnforcementBlocking) plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("NewEvidencePlan: %v", err) } if plan.MaxBufferRunes("ch") != 1500 { t.Errorf("MaxBufferRunes want 1500 (min hard bound), got %d", plan.MaxBufferRunes("ch")) } if plan.RollingEvidenceRunes("ch") != 80 { t.Errorf("RollingEvidenceRunes want 80 (max rolling evidence), got %d", plan.RollingEvidenceRunes("ch")) } }) t.Run("incompatibleAggregateRejectedWhenRollingExceedsHardBound", func(t *testing.T) { rroll, _ := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 1200, 2000) term, _ := NewFilterHoldRequirementTerminalGateWithMaxBuffer("ch", textKinds, EventKindTerminal, 1000) b1, _ := NewFilterHoldBinding("f1_roll", rroll, FilterEnforcementBlocking) b2, _ := NewFilterHoldBinding("f2_term", term, FilterEnforcementBlocking) _, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err == nil { t.Error("expected rejection when rolling evidence (1200) > min hard bound (1000)") } }) t.Run("runtimeStorageEnvelopeLifecycle", func(t *testing.T) { rroll, err := NewFilterHoldRequirementRollingWithMaxBuffer("ch", textKinds, 80, 2000) if err != nil { t.Fatalf("rolling req: %v", err) } term, err := NewFilterHoldRequirementTerminalGateWithMaxBuffer("ch", textKinds, EventKindTerminal, 1500) if err != nil { t.Fatalf("terminal req: %v", err) } b1, _ := NewFilterHoldBinding("f1_roll", rroll, FilterEnforcementBlocking) b2, _ := NewFilterHoldBinding("f2_term", term, FilterEnforcementBlocking) plan, err := NewEvidencePlan([]FilterHoldBinding{b1, b2}) if err != nil { t.Fatalf("plan: %v", err) } tail, err := NewEvidenceTail(plan) if err != nil { t.Fatalf("tail: %v", err) } sourceRunes := make([]rune, 100) for i := range sourceRunes { sourceRunes[i] = rune(0x4E00 + i) } text := string(sourceRunes) prefix80 := string(sourceRunes[:80]) expected := string(sourceRunes[20:]) if prefix80 == expected { t.Fatal("fixture must distinguish leading and trailing 80 runes") } representative, err := plan.BlockingRequirement("ch") if err != nil || representative.Mode() != FilterHoldModeTerminalGate { t.Fatalf("representative mode want terminal_gate, got %s (err=%v)", representative.Mode(), err) } ev, _ := NewTextDeltaEvent("ch", text, now) _, sig, err := tail.Append(ev) if err != nil { t.Fatalf("append text: %v", err) } cs := tail.channelState["ch"] if cs == nil { t.Fatal("expected channel state for ch after append") } if cs.maxBufferRunes != 1500 { t.Errorf("runtime maxBufferRunes want 1500, got %d", cs.maxBufferRunes) } if cs.effectiveEvidenceRunes != 80 { t.Errorf("runtime effectiveEvidenceRunes want 80, got %d", cs.effectiveEvidenceRunes) } evTerm, _ := NewTerminalEvent("ch", now) ep, sig, err := tail.Append(evTerm) if err != nil { t.Fatalf("append terminal: %v", err) } if sig != EvidenceTailSignalTrigger { t.Errorf("terminal signal want trigger, got %s", sig) } pr, err := tail.PrepareRelease(ep.ID()) if err != nil { t.Fatalf("prepare: %v", err) } if pr.Token() == "" { t.Error("prepared token must be non-empty") } if n := len(pr.ReleaseEvents()); n != 1 { t.Fatalf("prepared snapshot want 1 event, got %d", n) } if err := tail.ConfirmRelease(pr.Token(), ReleaseConfirmation{ReleasedEvents: 1}); err != nil { t.Fatalf("confirm: %v", err) } if got := tail.CommittedCursor("ch"); got != 1 { t.Errorf("committed cursor want 1, got %d", got) } if got := len(tail.channelState["ch"].pendingEntries); got != 0 { t.Errorf("pending entries want 0, got %d", got) } lb := tail.EffectiveLookBehind("ch") if len(lb) != 1 { t.Fatalf("look-behind want 1 event, got %d", len(lb)) } released, err := lb[0].AsTextDelta() if err != nil { t.Fatalf("look-behind event not text: %v", err) } if n := len([]rune(released)); n != 80 { t.Errorf("look-behind text want 80 runes, got %d", n) } runes := []rune(text) trailing80 := string(runes[len(runes)-80:]) if released != trailing80 { t.Errorf("look-behind text want trailing 80 runes %q, got %q", trailing80, released) } }) } func TestMultiRequirementPolicyHookExactSetBinding(t *testing.T) { textKinds := []EventKind{EventKindTextDelta} r1, _ := NewFilterHoldRequirementRolling("ch1", textKinds, 10) r2, _ := NewFilterHoldRequirementRolling("ch2", textKinds, 10) b1, _ := NewFilterHoldBinding("f1", r1, FilterEnforcementBlocking) b2, _ := NewFilterHoldBinding("f2", r2, FilterEnforcementBlocking) plan, _ := NewEvidencePlan([]FilterHoldBinding{b1, b2}) f1 := &mockTailFilter{id: "f1", req: r1} f2 := &mockTailFilter{id: "f2", req: r2} fReg1, _ := NewFilterRegistration(f1, "cap1", true, FilterEnforcementBlocking, time.Second, 10) fReg2, _ := NewFilterRegistration(f2, "cap2", true, FilterEnforcementBlocking, time.Second, 20) snap, _ := NewFilterRegistrySnapshot("gen1", []FilterRegistration{fReg1, fReg2}, nil) reqCtx, _ := NewRequestFilterContext("gen1", "att1", "env", "ep", "fam", "s1", CommitStateTransportUncommitted, false, false, "corr") reqSnap, _ := snap.BeginRequest(reqCtx) target, _ := NewAttemptTarget("mg", "m", "p", "path", []string{"cap1", "cap2"}) resolved, _ := reqSnap.ResolveAttempt(target) tail, _ := NewEvidenceTail(plan) ev, _ := NewTextDeltaEvent("ch1", "012345678901234", now) epoch, _, err := tail.Append(ev) if err != nil { t.Fatalf("append: %v", err) } // Pre-build resolvedExtra for table-driven matrix use. f3 := &mockTailFilter{id: "f3", req: r1} fReg3, _ := NewFilterRegistration(f3, "cap3", true, FilterEnforcementBlocking, time.Second, 30) snapExtra, _ := NewFilterRegistrySnapshot("gen1", []FilterRegistration{fReg1, fReg2, fReg3}, nil) reqSnapExtra, _ := snapExtra.BeginRequest(reqCtx) targetExtra, _ := NewAttemptTarget("mg", "m", "p", "path", []string{"cap1", "cap2", "cap3"}) resolvedExtra, _ := reqSnapExtra.ResolveAttempt(targetExtra) type matrixCase struct { name string epoch EvidenceEpoch filters []ResolvedFilter wantErr string successIDs []string notApplicable string } appF1, _ := NewFilterApplicability("f1", true, false) appF2, _ := NewFilterApplicability("f2", false, false) epochWithOnlyF1, _ := NewEvidenceEpochWithApplicabilities(99, "ch1", FilterHoldModeNone, false, "test_only_f1", []FilterApplicability{appF1}) appF3, _ := NewFilterApplicability("f3", true, false) epochWithF1F2F3, _ := NewEvidenceEpochWithApplicabilities(98, "ch1", FilterHoldModeNone, false, "test_f1_f2_f3", []FilterApplicability{appF1, appF2, appF3}) matrix := []matrixCase{ { name: "exactSetSuccess", epoch: epoch, filters: resolved, wantErr: "", successIDs: []string{"f1", "f2"}, notApplicable: "f2", }, { name: "resolvedEmpty", epoch: epoch, filters: nil, wantErr: "streamgate: resolved filter count mismatch: got 0, expected 2", }, { name: "resolvedSubset", epoch: epoch, filters: resolved[:1], wantErr: "streamgate: resolved filter count mismatch: got 1, expected 2", }, { name: "resolvedExtra", epoch: epoch, filters: resolvedExtra, wantErr: "streamgate: resolved filter count mismatch: got 3, expected 2", }, { name: "resolvedDuplicate", epoch: epoch, filters: []ResolvedFilter{resolved[0], resolved[0]}, wantErr: "streamgate: duplicate resolved filter id in bind epoch: f1", }, { name: "applicabilityMissing", epoch: epochWithOnlyF1, filters: resolved, wantErr: "streamgate: applicability count mismatch: got 1, expected 2", }, { name: "applicabilityExtra", epoch: epochWithF1F2F3, filters: resolved, wantErr: "streamgate: applicability count mismatch: got 3, expected 2", }, } for _, tc := range matrix { tc := tc t.Run(tc.name, func(t *testing.T) { bound, err := plan.BindEpochFilters(tc.epoch, tc.filters) if tc.wantErr != "" { if err == nil { t.Errorf("expected error containing %q, got nil", tc.wantErr) } else if !strings.Contains(err.Error(), tc.wantErr) { t.Errorf("error want to contain %q, got %v", tc.wantErr, err) } return } if err != nil { t.Fatalf("BindEpochFilters %s: %v", tc.name, err) } if len(bound) != len(tc.successIDs) { t.Fatalf("bound count want %d, got %d", len(tc.successIDs), len(bound)) } for i, id := range tc.successIDs { if bound[i].ID() != id { t.Errorf("bound[%d] ID want %s, got %s", i, id, bound[i].ID()) } } if tc.notApplicable != "" { for _, ef := range bound { if ef.ID() == tc.notApplicable && ef.EvaluatedForEpoch() { t.Errorf("%s should be explicit not-applicable in this epoch", tc.notApplicable) } } } }) } t.Run("applicabilityDuplicate", func(t *testing.T) { dupApps := []FilterApplicability{appF1, appF1} _, err := NewEvidenceEpochWithApplicabilities(97, "ch1", FilterHoldModeNone, false, "test_dup_app", dupApps) if err == nil { t.Error("expected constructor error for duplicate applicability filter id") } else if !strings.Contains(err.Error(), "streamgate: duplicate epoch applicability filter id: f1") { t.Errorf("constructor error want to contain %q, got %v", "duplicate epoch applicability filter id: f1", err) } }) } func TestMultiRequirementPolicyHookTypedPolicySeam(t *testing.T) { req, _ := NewFilterHoldRequirementRolling("ch", []EventKind{EventKindTextDelta}, 10) t.Run("validEnforcements", func(t *testing.T) { bBlock, err := NewFilterHoldBinding("f1", req, FilterEnforcementBlocking) if err != nil || !bBlock.BlocksRelease() || bBlock.Enforcement() != FilterEnforcementBlocking { t.Errorf("blocking binding invalid: %v, blocksRelease=%v", err, bBlock.BlocksRelease()) } bObs, err := NewFilterHoldBinding("f2", req, FilterEnforcementObserveOnly) if err != nil || bObs.BlocksRelease() || bObs.Enforcement() != FilterEnforcementObserveOnly { t.Errorf("observe binding invalid: %v, blocksRelease=%v", err, bObs.BlocksRelease()) } noneReq, _ := NewFilterHoldRequirementNone("ch", []EventKind{EventKindTextDelta}) bNone, err := NewFilterHoldBinding("f3", noneReq, FilterEnforcementBlocking) if err != nil || bNone.BlocksRelease() { t.Errorf("none mode binding should have blocksRelease=false even if enforcement=blocking: %v, blocksRelease=%v", err, bNone.BlocksRelease()) } }) t.Run("invalidEnforcementRejected", func(t *testing.T) { _, err := NewFilterHoldBinding("f1", req, FilterEnforcement("invalid_mode")) if err == nil { t.Error("expected error for invalid FilterEnforcement, got nil") } }) }