iop/packages/go/streamgate/evidence_tail.go
toki 7634ca8962 feat(stream-evidence-gate-core): complete stream evidence gate core implementation
- OpenAI request rebuilder with tool validation and provider tunnel
- Edge config runtime refresh for stream evidence gate
- Filter observation contract and runtime with sink/correlation
- Stream gate dispatcher, release sink, and vertical slice
- Recovery coordinator for evidence tail
- Parallel evaluation and commit boundary
- E2E test script for OpenAI vLLM
- Archive completed task groups to archive/2026/07
2026-07-28 04:08:12 +09:00

2438 lines
84 KiB
Go

// Package streamgate provides transport-agnostic event contract types for the
// stream evidence gate. It owns the normalized event lifecycle, immutable
// evidence/filter decision types, and terminal/failure/release payloads.
//
// This package must not import apps/, proto/, or packages/go/config/.
// It depends only on the Go standard library.
package streamgate
import (
"errors"
"fmt"
"sort"
"strings"
"unicode/utf8"
)
// FilterHoldMode identifies how a filter holds pending evidence before
// allowing release. It is a closed set of modes; each mode defines its own
// trigger condition and field requirements.
type FilterHoldMode string
// knownFilterHoldModes is the closed set of valid FilterHoldMode values.
var knownFilterHoldModes = map[FilterHoldMode]struct{}{
FilterHoldModeNone: {},
FilterHoldModeRolling: {},
FilterHoldModeTerminalGate: {},
FilterHoldModeFragmentGate: {},
}
const (
// FilterHoldModeNone indicates no hold. Events pass through immediately.
FilterHoldModeNone FilterHoldMode = "none"
// FilterHoldModeRolling indicates the filter holds until the rolling
// rune evidence window meets a configured threshold. Time is not a
// release condition.
FilterHoldModeRolling FilterHoldMode = "rolling_window"
// FilterHoldModeTerminalGate indicates the filter holds all events until
// a terminal event is received. Terminal or provider-error events serve
// as the trigger.
FilterHoldModeTerminalGate FilterHoldMode = "terminal_gate"
// FilterHoldModeFragmentGate indicates the filter holds tool-call
// fragments until a CompleteFragment signal arrives for the keyed ID.
FilterHoldModeFragmentGate FilterHoldMode = "fragment_gate"
)
// Validate returns nil when the mode is a known lifecycle value.
func (m FilterHoldMode) Validate() error {
switch m {
case FilterHoldModeNone, FilterHoldModeRolling, FilterHoldModeTerminalGate, FilterHoldModeFragmentGate:
return nil
}
return errors.New("streamgate: unknown filter hold mode: " + string(m))
}
// IsValidMode returns true when m is a known hold mode value.
func IsValidMode(m FilterHoldMode) bool {
_, ok := knownFilterHoldModes[m]
return ok
}
// defaultEvidenceRunes is the default Unicode rune evidence window used when
// no policy override is provided.
const defaultEvidenceRunes = 500
// minEvidenceRunes is the absolute minimum allowed rune window.
const minEvidenceRunes = 1
// maxEvidenceRunes is the absolute maximum allowed rune window.
const maxEvidenceRunes = 65536
// defaultMaxBufferRunes is the default hard buffer limit in runes.
const defaultMaxBufferRunes = 4096
// minMaxBufferRunes is the absolute minimum allowed buffer limit.
const minMaxBufferRunes = 10
// maxMaxBufferRunes is the absolute maximum allowed buffer limit.
const maxMaxBufferRunes = 1048576
// FilterHoldRequirement is an immutable specification of how a filter holds
// pending evidence before releasing it. The fields match the SDD naming and
// semantics exactly. A requirement is mode-specific: rolling requires a
// positive rune threshold, terminal gate requires a terminal trigger kind,
// fragment gate requires a fragment trigger kind, and none is a no-op.
type FilterHoldRequirement struct {
channel string
mode FilterHoldMode
subscribedKinds []EventKind
evidenceRunes int
triggerKind EventKind
maxBufferRunes int
}
// NewFilterHoldRequirementNone creates a FilterHoldRequirement of mode none.
// None mode has no rune threshold, no trigger, and releases immediately.
// Subscribed kinds are required but unused.
func NewFilterHoldRequirementNone(channel string, subscribedKinds []EventKind) (FilterHoldRequirement, error) {
if err := validateHoldChannel("none", channel); err != nil {
return FilterHoldRequirement{}, err
}
kinds, err := validateSubscribedKinds("none", subscribedKinds)
if err != nil {
return FilterHoldRequirement{}, err
}
req := FilterHoldRequirement{
channel: channel,
mode: FilterHoldModeNone,
subscribedKinds: kinds,
evidenceRunes: 0,
triggerKind: "",
maxBufferRunes: 0,
}
if err := req.Validate(); err != nil {
return FilterHoldRequirement{}, err
}
return req, nil
}
// NewFilterHoldRequirementRollingWithMaxBuffer creates a FilterHoldRequirement
// of mode rolling_window with an explicit max_buffer_runes override. The
// rune threshold must be within the allowed range. The max_buffer_runes must
// be >= evidence_runes to ensure the buffer can hold the evidence window.
func NewFilterHoldRequirementRollingWithMaxBuffer(channel string, subscribedKinds []EventKind, evidenceRunes int, maxBufferRunes int) (FilterHoldRequirement, error) {
if err := validateHoldChannel("rolling_window", channel); err != nil {
return FilterHoldRequirement{}, err
}
if evidenceRunes < minEvidenceRunes || evidenceRunes > maxEvidenceRunes {
return FilterHoldRequirement{}, fmt.Errorf("streamgate: rolling_window evidence_runes must be between %d and %d, got %d", minEvidenceRunes, maxEvidenceRunes, evidenceRunes)
}
if maxBufferRunes < evidenceRunes {
return FilterHoldRequirement{}, fmt.Errorf("streamgate: rolling_window max_buffer_runes must be >= evidence_runes, got evidence=%d buffer=%d", evidenceRunes, maxBufferRunes)
}
if maxBufferRunes < minMaxBufferRunes || maxBufferRunes > maxMaxBufferRunes {
return FilterHoldRequirement{}, fmt.Errorf("streamgate: rolling_window max_buffer_runes must be between %d and %d, got %d", minMaxBufferRunes, maxMaxBufferRunes, maxBufferRunes)
}
kinds, err := validateSubscribedKinds("rolling_window", subscribedKinds)
if err != nil {
return FilterHoldRequirement{}, err
}
req := FilterHoldRequirement{
channel: channel,
mode: FilterHoldModeRolling,
subscribedKinds: kinds,
evidenceRunes: evidenceRunes,
triggerKind: "",
maxBufferRunes: maxBufferRunes,
}
if err := req.Validate(); err != nil {
return FilterHoldRequirement{}, err
}
return req, nil
}
// NewFilterHoldRequirementRolling creates a FilterHoldRequirement of mode
// rolling_window. The rune threshold must be within the allowed range. The
// subscribed kinds define which event types contribute to the rolling count.
func NewFilterHoldRequirementRolling(channel string, subscribedKinds []EventKind, evidenceRunes int) (FilterHoldRequirement, error) {
return NewFilterHoldRequirementRollingWithMaxBuffer(channel, subscribedKinds, evidenceRunes, defaultMaxBufferRunes)
}
// NewFilterHoldRequirementTerminalGateWithMaxBuffer creates a FilterHoldRequirement
// of mode terminal_gate with an explicit max_buffer_runes override.
func NewFilterHoldRequirementTerminalGateWithMaxBuffer(channel string, subscribedKinds []EventKind, triggerKind EventKind, maxBufferRunes int) (FilterHoldRequirement, error) {
if err := validateHoldChannel("terminal_gate", channel); err != nil {
return FilterHoldRequirement{}, err
}
if triggerKind != EventKindTerminal && triggerKind != EventKindProviderError {
return FilterHoldRequirement{}, fmt.Errorf("streamgate: terminal_gate trigger must be terminal or provider_error, got %q", triggerKind)
}
if maxBufferRunes < minMaxBufferRunes || maxBufferRunes > maxMaxBufferRunes {
return FilterHoldRequirement{}, fmt.Errorf("streamgate: terminal_gate max_buffer_runes must be between %d and %d, got %d", minMaxBufferRunes, maxMaxBufferRunes, maxBufferRunes)
}
kinds, err := validateSubscribedKinds("terminal_gate", subscribedKinds)
if err != nil {
return FilterHoldRequirement{}, err
}
req := FilterHoldRequirement{
channel: channel,
mode: FilterHoldModeTerminalGate,
subscribedKinds: kinds,
evidenceRunes: 0,
triggerKind: triggerKind,
maxBufferRunes: maxBufferRunes,
}
if err := req.Validate(); err != nil {
return FilterHoldRequirement{}, err
}
return req, nil
}
// NewFilterHoldRequirementTerminalGate creates a FilterHoldRequirement of
// mode terminal_gate. The trigger kind must be terminal or provider_error.
func NewFilterHoldRequirementTerminalGate(channel string, subscribedKinds []EventKind, triggerKind EventKind) (FilterHoldRequirement, error) {
return NewFilterHoldRequirementTerminalGateWithMaxBuffer(channel, subscribedKinds, triggerKind, defaultMaxBufferRunes)
}
// NewFilterHoldRequirementFragmentGateWithMaxBuffer creates a FilterHoldRequirement
// of mode fragment_gate with an explicit max_buffer_runes override.
func NewFilterHoldRequirementFragmentGateWithMaxBuffer(channel string, subscribedKinds []EventKind, triggerKind EventKind, maxBufferRunes int) (FilterHoldRequirement, error) {
if err := validateHoldChannel("fragment_gate", channel); err != nil {
return FilterHoldRequirement{}, err
}
if triggerKind != EventKindToolCallFragment {
return FilterHoldRequirement{}, fmt.Errorf("streamgate: fragment_gate trigger must be tool_call_fragment, got %q", triggerKind)
}
if maxBufferRunes < minMaxBufferRunes || maxBufferRunes > maxMaxBufferRunes {
return FilterHoldRequirement{}, fmt.Errorf("streamgate: fragment_gate max_buffer_runes must be between %d and %d, got %d", minMaxBufferRunes, maxMaxBufferRunes, maxBufferRunes)
}
kinds, err := validateSubscribedKinds("fragment_gate", subscribedKinds)
if err != nil {
return FilterHoldRequirement{}, err
}
req := FilterHoldRequirement{
channel: channel,
mode: FilterHoldModeFragmentGate,
subscribedKinds: kinds,
evidenceRunes: 0,
triggerKind: triggerKind,
maxBufferRunes: maxBufferRunes,
}
if err := req.Validate(); err != nil {
return FilterHoldRequirement{}, err
}
return req, nil
}
// NewFilterHoldRequirementFragmentGate creates a FilterHoldRequirement of
// mode fragment_gate. The trigger kind must be tool_call_fragment.
func NewFilterHoldRequirementFragmentGate(channel string, subscribedKinds []EventKind, triggerKind EventKind) (FilterHoldRequirement, error) {
return NewFilterHoldRequirementFragmentGateWithMaxBuffer(channel, subscribedKinds, triggerKind, defaultMaxBufferRunes)
}
func validateHoldChannel(mode, channel string) error {
if channel == "" {
return errors.New("streamgate: hold requirement channel is required")
}
return nil
}
func validateSubscribedKinds(mode string, kinds []EventKind) ([]EventKind, error) {
if len(kinds) == 0 {
return nil, fmt.Errorf("streamgate: %s hold requirement subscribed kinds must not be empty", mode)
}
seen := make(map[EventKind]struct{}, len(kinds))
for _, k := range kinds {
if err := k.Validate(); err != nil {
return nil, fmt.Errorf("streamgate: %s hold requirement subscribed kind: %v", mode, err)
}
if _, dup := seen[k]; dup {
return nil, fmt.Errorf("streamgate: %s hold requirement duplicate subscribed kind: %s", mode, k)
}
seen[k] = struct{}{}
}
out := make([]EventKind, len(kinds))
copy(out, kinds)
sort.Slice(out, func(i, j int) bool { return out[i] < out[j] })
return out, nil
}
// Validate returns nil when the requirement is in a consistent state.
// Mode-specific field rules are enforced.
func (r FilterHoldRequirement) Validate() error {
if r.channel == "" {
return errors.New("streamgate: hold requirement channel is required")
}
if err := r.mode.Validate(); err != nil {
return err
}
if len(r.subscribedKinds) == 0 {
return errors.New("streamgate: hold requirement subscribed kinds must not be empty")
}
switch r.mode {
case FilterHoldModeNone:
if r.evidenceRunes != 0 {
return errors.New("streamgate: none mode must not have evidence_runes")
}
if r.triggerKind != "" {
return errors.New("streamgate: none mode must not have a trigger")
}
if r.maxBufferRunes != 0 {
return errors.New("streamgate: none mode must not have max_buffer_runes")
}
case FilterHoldModeRolling:
if r.evidenceRunes < minEvidenceRunes || r.evidenceRunes > maxEvidenceRunes {
return errors.New("streamgate: rolling_window evidence_runes out of range")
}
if r.triggerKind != "" {
return errors.New("streamgate: rolling_window mode must not have a trigger")
}
if r.maxBufferRunes < minMaxBufferRunes || r.maxBufferRunes > maxMaxBufferRunes {
return errors.New("streamgate: rolling_window max_buffer_runes out of range")
}
case FilterHoldModeTerminalGate:
if r.evidenceRunes != 0 {
return errors.New("streamgate: terminal_gate mode must not have evidence_runes")
}
if r.triggerKind != EventKindTerminal && r.triggerKind != EventKindProviderError {
return errors.New("streamgate: terminal_gate trigger must be terminal or provider_error")
}
if r.maxBufferRunes < minMaxBufferRunes || r.maxBufferRunes > maxMaxBufferRunes {
return errors.New("streamgate: terminal_gate max_buffer_runes out of range")
}
case FilterHoldModeFragmentGate:
if r.evidenceRunes != 0 {
return errors.New("streamgate: fragment_gate mode must not have evidence_runes")
}
if r.triggerKind != EventKindToolCallFragment {
return errors.New("streamgate: fragment_gate trigger must be tool_call_fragment")
}
if r.maxBufferRunes < minMaxBufferRunes || r.maxBufferRunes > maxMaxBufferRunes {
return errors.New("streamgate: fragment_gate max_buffer_runes out of range")
}
}
return nil
}
// Channel returns the requirement channel key.
func (r FilterHoldRequirement) Channel() string { return r.channel }
// Mode returns the hold mode.
func (r FilterHoldRequirement) Mode() FilterHoldMode { return r.mode }
// SubscribedKinds returns a defensive copy of the subscribed event kinds.
func (r FilterHoldRequirement) SubscribedKinds() []EventKind {
if r.subscribedKinds == nil {
return nil
}
out := make([]EventKind, len(r.subscribedKinds))
copy(out, r.subscribedKinds)
return out
}
// EvidenceRunes returns the rolling rune threshold. Returns 0 for non-rolling
// modes.
func (r FilterHoldRequirement) EvidenceRunes() int { return r.evidenceRunes }
// TriggerKind returns the trigger event kind. Returns "" for modes that do
// not use a trigger.
func (r FilterHoldRequirement) TriggerKind() EventKind { return r.triggerKind }
// MaxBufferRunes returns the hard buffer limit in runes.
func (r FilterHoldRequirement) MaxBufferRunes() int { return r.maxBufferRunes }
// IsBlocking returns true when this requirement's mode blocks channel release.
// None and observe-only modes are non-blocking; rolling, terminal_gate, and
// fragment_gate are blocking.
func (r FilterHoldRequirement) IsBlocking() bool {
switch r.mode {
case FilterHoldModeRolling, FilterHoldModeTerminalGate, FilterHoldModeFragmentGate:
return true
default:
return false
}
}
// FilterHoldBinding pairs a stable filter ID with its requirement and an
// explicit FilterEnforcement resolved by the caller. The binding is
// immutable to callers.
type FilterHoldBinding struct {
filterID StableToken
requirement FilterHoldRequirement
enforcement FilterEnforcement
blocksRelease bool
}
// NewFilterHoldBinding creates a FilterHoldBinding with validation. The
// returned value is immutable to callers.
func NewFilterHoldBinding(filterID string, req FilterHoldRequirement, enforcement FilterEnforcement) (FilterHoldBinding, error) {
f, err := NewStableTokenRequired("filterID", filterID)
if err != nil {
return FilterHoldBinding{}, err
}
if err := req.Validate(); err != nil {
return FilterHoldBinding{}, err
}
if err := enforcement.Validate(); err != nil {
return FilterHoldBinding{}, err
}
blocksRelease := (enforcement == FilterEnforcementBlocking && req.Mode() != FilterHoldModeNone)
return FilterHoldBinding{
filterID: f,
requirement: req,
enforcement: enforcement,
blocksRelease: blocksRelease,
}, nil
}
// NewFilterHoldBindingFromResolvedFilter constructs a FilterHoldBinding from a ResolvedFilter.
func NewFilterHoldBindingFromResolvedFilter(rf ResolvedFilter) (FilterHoldBinding, error) {
return NewFilterHoldBinding(rf.FilterID(), rf.HoldRequirement(), rf.Enforcement())
}
// FilterID returns the stable filter id token.
func (b FilterHoldBinding) FilterID() string { return b.filterID.value }
// Requirement returns the hold requirement.
func (b FilterHoldBinding) Requirement() FilterHoldRequirement { return b.requirement }
// Enforcement returns the effective enforcement.
func (b FilterHoldBinding) Enforcement() FilterEnforcement { return b.enforcement }
// BlocksRelease returns whether this binding blocks channel release.
func (b FilterHoldBinding) BlocksRelease() bool { return b.blocksRelease }
// FilterApplicability captures typed readiness and subscription status for a single filter in an epoch.
type FilterApplicability struct {
filterID StableToken
subscribedEventPresent bool
triggerReady bool
}
// NewFilterApplicability creates a validated FilterApplicability.
func NewFilterApplicability(filterID string, subscribedEventPresent, triggerReady bool) (FilterApplicability, error) {
st, err := NewStableTokenRequired("filterID", filterID)
if err != nil {
return FilterApplicability{}, fmt.Errorf("streamgate: filter applicability id: %w", err)
}
return FilterApplicability{
filterID: st,
subscribedEventPresent: subscribedEventPresent,
triggerReady: triggerReady,
}, nil
}
// FilterID returns the stable filter id.
func (a FilterApplicability) FilterID() string { return a.filterID.value }
// SubscribedEventPresent returns whether subscribed events were present for this filter in the epoch.
func (a FilterApplicability) SubscribedEventPresent() bool { return a.subscribedEventPresent }
// TriggerReady returns whether the filter trigger condition is satisfied in the epoch.
func (a FilterApplicability) TriggerReady() bool { return a.triggerReady }
// evidencePlan holds the compiled per-channel hold plan derived from bindings.
type evidencePlan struct {
bindingsByChannel map[string][]FilterHoldBinding
allBindings map[string]FilterHoldBinding
observeOnlyKinds map[string][]EventKind
allKinds map[string][]EventKind
blockingKinds map[string][]EventKind
maxBufferRunes map[string]int
rollingEvidenceRunes map[string]int
}
// compileEvidencePlan produces a compiled plan from a slice of bindings.
func compileEvidencePlan(bindings []FilterHoldBinding) (evidencePlan, error) {
plan := evidencePlan{
bindingsByChannel: make(map[string][]FilterHoldBinding),
allBindings: make(map[string]FilterHoldBinding),
observeOnlyKinds: make(map[string][]EventKind),
allKinds: make(map[string][]EventKind),
blockingKinds: make(map[string][]EventKind),
maxBufferRunes: make(map[string]int),
rollingEvidenceRunes: make(map[string]int),
}
for _, b := range bindings {
fID := b.FilterID()
if fID == "" {
return evidencePlan{}, errors.New("streamgate: binding filter id is required")
}
if _, dup := plan.allBindings[fID]; dup {
return evidencePlan{}, fmt.Errorf("streamgate: duplicate binding filter id: %s", fID)
}
req := b.requirement
if err := req.Validate(); err != nil {
return evidencePlan{}, fmt.Errorf("streamgate: binding requirement invalid: %w", err)
}
ch := req.Channel()
plan.allBindings[fID] = b
plan.bindingsByChannel[ch] = append(plan.bindingsByChannel[ch], b)
plan.allKinds[ch] = mergeKinds(plan.allKinds[ch], req.SubscribedKinds())
if b.BlocksRelease() && req.IsBlocking() {
plan.blockingKinds[ch] = mergeKinds(plan.blockingKinds[ch], req.SubscribedKinds())
plan.maxBufferRunes[ch] = minPositiveBound(plan.maxBufferRunes[ch], req.MaxBufferRunes())
if req.Mode() == FilterHoldModeRolling {
plan.rollingEvidenceRunes[ch] = maxInt(plan.rollingEvidenceRunes[ch], req.EvidenceRunes())
}
} else {
plan.observeOnlyKinds[ch] = mergeKinds(plan.observeOnlyKinds[ch], req.SubscribedKinds())
}
}
// Sort bindings per channel by filter id ascending
for ch := range plan.bindingsByChannel {
sort.SliceStable(plan.bindingsByChannel[ch], func(i, j int) bool {
return plan.bindingsByChannel[ch][i].FilterID() < plan.bindingsByChannel[ch][j].FilterID()
})
}
// Validate bounds per channel: max rolling evidence must be <= shared hard bound
for ch := range plan.bindingsByChannel {
maxBuf := plan.maxBufferRunes[ch]
rollEv := plan.rollingEvidenceRunes[ch]
if rollEv > 0 && maxBuf > 0 && rollEv > maxBuf {
return evidencePlan{}, fmt.Errorf("streamgate: max rolling evidence %d > shared hard bound %d for channel %s", rollEv, maxBuf, ch)
}
}
return plan, nil
}
// isKindSubscribed returns true when kind is in the subscribed kinds slice.
func isKindSubscribed(kind EventKind, subscribed []EventKind) bool {
for _, k := range subscribed {
if k == kind {
return true
}
}
return false
}
// appendUniqueKind appends kind to existing only if not already present.
func appendUniqueKind(existing []EventKind, kind EventKind) []EventKind {
for _, k := range existing {
if k == kind {
return existing
}
}
return append(existing, kind)
}
// strongerMode returns true when a is strictly stronger than b in the
// hold priority: terminal_gate > fragment_gate > rolling_window.
func strongerMode(a, b FilterHoldMode) bool {
priority := map[FilterHoldMode]int{
FilterHoldModeRolling: 1,
FilterHoldModeFragmentGate: 2,
FilterHoldModeTerminalGate: 3,
}
return priority[a] > priority[b]
}
// minPositiveBound returns the smaller of two positive values; if either is
// zero (e.g. none mode has no hard bound) it returns the other.
func minPositiveBound(a, b int) int {
if a == 0 {
return b
}
if b == 0 {
return a
}
if a < b {
return a
}
return b
}
// maxInt returns the larger of two ints.
func maxInt(a, b int) int {
if a > b {
return a
}
return b
}
func mergeKinds(a, b []EventKind) []EventKind {
seen := make(map[EventKind]struct{}, len(a)+len(b))
var out []EventKind
for _, k := range a {
if _, ok := seen[k]; !ok {
seen[k] = struct{}{}
out = append(out, k)
}
}
for _, k := range b {
if _, ok := seen[k]; !ok {
seen[k] = struct{}{}
out = append(out, k)
}
}
sort.Slice(out, func(i, j int) bool { return out[i] < out[j] })
return out
}
// EvidencePlan is the compiled channel plan for evidence holding. It is
// immutable to callers once created. The plan only includes blocking
// requirements; observe-only and none-mode bindings never create a channel
// hold on their own.
type EvidencePlan struct {
plan evidencePlan
}
// NewEvidencePlan compiles a channel hold plan from the provided bindings.
// Only blocking bindings affect the hold decision. Returns an error if any
// binding's requirement fails validation.
func NewEvidencePlan(bindings []FilterHoldBinding) (EvidencePlan, error) {
if len(bindings) == 0 {
return EvidencePlan{plan: evidencePlan{
bindingsByChannel: make(map[string][]FilterHoldBinding),
allBindings: make(map[string]FilterHoldBinding),
observeOnlyKinds: make(map[string][]EventKind),
allKinds: make(map[string][]EventKind),
blockingKinds: make(map[string][]EventKind),
maxBufferRunes: make(map[string]int),
rollingEvidenceRunes: make(map[string]int),
}}, nil
}
compiled, err := compileEvidencePlan(bindings)
if err != nil {
return EvidencePlan{}, err
}
return EvidencePlan{plan: compiled}, nil
}
// NewEvidencePlanFromResolvedFilters compiles an EvidencePlan from ResolvedFilters.
func NewEvidencePlanFromResolvedFilters(resolved []ResolvedFilter) (EvidencePlan, error) {
if len(resolved) == 0 {
return NewEvidencePlan(nil)
}
seen := make(map[string]struct{}, len(resolved))
bindings := make([]FilterHoldBinding, len(resolved))
for i, rf := range resolved {
fID := rf.FilterID()
if fID == "" {
return EvidencePlan{}, errors.New("streamgate: resolved filter id is required")
}
if _, dup := seen[fID]; dup {
return EvidencePlan{}, fmt.Errorf("streamgate: duplicate resolved filter id: %s", fID)
}
seen[fID] = struct{}{}
b, err := NewFilterHoldBindingFromResolvedFilter(rf)
if err != nil {
return EvidencePlan{}, fmt.Errorf("streamgate: build binding from resolved filter: %w", err)
}
bindings[i] = b
}
return NewEvidencePlan(bindings)
}
// BindingsForChannel returns a defensive copy of all bindings for the channel in stable filter id order.
func (p EvidencePlan) BindingsForChannel(channel string) []FilterHoldBinding {
bindings := p.plan.bindingsByChannel[channel]
if len(bindings) == 0 {
return nil
}
out := make([]FilterHoldBinding, len(bindings))
copy(out, bindings)
return out
}
// Bindings returns a defensive copy of all bindings in stable filter id order.
func (p EvidencePlan) Bindings() []FilterHoldBinding {
var all []FilterHoldBinding
for _, b := range p.plan.allBindings {
all = append(all, b)
}
sort.SliceStable(all, func(i, j int) bool {
return all[i].FilterID() < all[j].FilterID()
})
return all
}
// HasModeForChannel returns true when channel has at least one binding with the given mode.
func (p EvidencePlan) HasModeForChannel(channel string, mode FilterHoldMode) bool {
bindings := p.plan.bindingsByChannel[channel]
for _, b := range bindings {
if b.Requirement().Mode() == mode {
return true
}
}
return false
}
// IsTerminalTriggerForChannel returns true when eventKind is a configured terminal trigger for any binding on channel.
func (p EvidencePlan) IsTerminalTriggerForChannel(channel string, kind EventKind) bool {
bindings := p.plan.bindingsByChannel[channel]
for _, b := range bindings {
if b.BlocksRelease() && b.Requirement().Mode() == FilterHoldModeTerminalGate {
if b.Requirement().TriggerKind() == kind {
return true
}
}
}
return false
}
// BlockingRequirementFor returns a representative storage envelope requirement for the channel.
func (p EvidencePlan) BlockingRequirementFor(channel string) (FilterHoldRequirement, bool) {
bindings := p.plan.bindingsByChannel[channel]
if len(bindings) == 0 {
return FilterHoldRequirement{}, false
}
var strongestMode FilterHoldMode = FilterHoldModeNone
var triggerKind EventKind
for _, b := range bindings {
if !b.BlocksRelease() || !b.Requirement().IsBlocking() {
continue
}
m := b.Requirement().Mode()
if strongerMode(m, strongestMode) {
strongestMode = m
triggerKind = b.Requirement().TriggerKind()
}
}
if strongestMode == FilterHoldModeNone {
return FilterHoldRequirement{}, false
}
maxBuf := p.plan.maxBufferRunes[channel]
rollEv := p.plan.rollingEvidenceRunes[channel]
blockingKinds := p.plan.blockingKinds[channel]
var req FilterHoldRequirement
var err error
switch strongestMode {
case FilterHoldModeRolling:
if rollEv <= 0 {
rollEv = defaultEvidenceRunes
}
if maxBuf <= 0 {
maxBuf = defaultMaxBufferRunes
}
if maxBuf < rollEv {
maxBuf = rollEv
}
req, err = NewFilterHoldRequirementRollingWithMaxBuffer(channel, blockingKinds, rollEv, maxBuf)
case FilterHoldModeTerminalGate:
if maxBuf <= 0 {
maxBuf = defaultMaxBufferRunes
}
req, err = NewFilterHoldRequirementTerminalGateWithMaxBuffer(channel, blockingKinds, triggerKind, maxBuf)
case FilterHoldModeFragmentGate:
if maxBuf <= 0 {
maxBuf = defaultMaxBufferRunes
}
req, err = NewFilterHoldRequirementFragmentGateWithMaxBuffer(channel, blockingKinds, triggerKind, maxBuf)
}
if err != nil {
return FilterHoldRequirement{}, false
}
return req, true
}
// HasBlockingRequirement returns true when the given channel has at least one
// blocking requirement in the plan.
func (p EvidencePlan) HasBlockingRequirement(channel string) bool {
bindings := p.plan.bindingsByChannel[channel]
for _, b := range bindings {
if b.BlocksRelease() && b.Requirement().IsBlocking() {
return true
}
}
return false
}
// BlockingRequirement returns the representative blocking requirement for the
// channel, or an error if no blocking requirement exists.
func (p EvidencePlan) BlockingRequirement(channel string) (FilterHoldRequirement, error) {
req, ok := p.BlockingRequirementFor(channel)
if !ok {
return FilterHoldRequirement{}, errors.New("streamgate: no blocking requirement for channel: " + channel)
}
return req, nil
}
// SubscribeKinds returns a defensive copy of the merged subscribed kinds for the channel.
func (p EvidencePlan) SubscribeKinds(channel string) []EventKind {
kinds := p.plan.allKinds[channel]
if kinds == nil {
return nil
}
out := make([]EventKind, len(kinds))
copy(out, kinds)
return out
}
// BlockingSubscribeKinds returns a defensive copy of the blocking subscribed
// kinds for the channel. These are the kinds that contribute to pending state.
// Returns nil when the channel has no blocking requirement.
func (p EvidencePlan) BlockingSubscribeKinds(channel string) []EventKind {
kinds := p.plan.blockingKinds[channel]
if kinds == nil {
return nil
}
out := make([]EventKind, len(kinds))
copy(out, kinds)
return out
}
// BlockingKinds returns the union of event kinds that contribute to the
// channel's blocking pending state. Returns nil when the channel has no
// blocking requirement.
func (p EvidencePlan) BlockingKinds(channel string) []EventKind {
kinds := p.plan.blockingKinds[channel]
if kinds == nil {
return nil
}
out := make([]EventKind, len(kinds))
copy(out, kinds)
return out
}
// ObserveKinds returns the union of event kinds that contribute to
// observation-only subscriptions but never to blocking pending state.
// Returns nil when the channel has no observe-only subscription.
func (p EvidencePlan) ObserveKinds(channel string) []EventKind {
kinds := p.plan.observeOnlyKinds[channel]
if kinds == nil {
return nil
}
out := make([]EventKind, len(kinds))
copy(out, kinds)
return out
}
// AllBindings returns a defensive copy of all bindings in stable filter ID order.
func (p EvidencePlan) AllBindings() []FilterHoldBinding {
if p.plan.allBindings == nil {
return nil
}
out := make([]FilterHoldBinding, 0, len(p.plan.allBindings))
for _, b := range p.plan.allBindings {
out = append(out, b)
}
sort.SliceStable(out, func(i, j int) bool {
return out[i].FilterID() < out[j].FilterID()
})
return out
}
func (p EvidencePlan) bindingIDs() []string {
bindings := p.AllBindings()
ids := make([]string, len(bindings))
for i, b := range bindings {
ids[i] = b.FilterID()
}
return ids
}
// MaxBufferRunes returns the compiled max buffer runes for the channel.
func (p EvidencePlan) MaxBufferRunes(channel string) int {
return p.plan.maxBufferRunes[channel]
}
// RollingEvidenceRunes returns the compiled rolling evidence runes for the channel.
func (p EvidencePlan) RollingEvidenceRunes(channel string) int {
return p.plan.rollingEvidenceRunes[channel]
}
// BindEpochFilters evaluates epoch applicability and binds resolved filters to EpochFilters in stable filter ID order.
func (p EvidencePlan) BindEpochFilters(epoch EvidenceEpoch, resolvedFilters []ResolvedFilter) ([]EpochFilter, error) {
if epoch.ID() == 0 {
return nil, errors.New("streamgate: epoch id must be positive")
}
planIDs := p.bindingIDs()
if len(planIDs) == 0 && len(resolvedFilters) == 0 {
return nil, nil
}
resolvedByID, err := indexResolvedFilters(resolvedFilters)
if err != nil {
return nil, err
}
appsByID, err := indexApplicabilities(epoch.Applicabilities())
if err != nil {
return nil, err
}
if err := requireExactFilterIDs(planIDs, resolvedByID, appsByID); err != nil {
return nil, err
}
out := make([]EpochFilter, 0, len(planIDs))
for _, id := range planIDs {
rf := resolvedByID[id]
app := appsByID[id]
ef, err := rf.BindEpoch(epoch.ID(), app)
if err != nil {
return nil, fmt.Errorf("streamgate: bind epoch filter %s: %w", id, err)
}
out = append(out, ef)
}
return out, nil
}
func indexResolvedFilters(filters []ResolvedFilter) (map[string]ResolvedFilter, error) {
out := make(map[string]ResolvedFilter, len(filters))
for _, rf := range filters {
fID := rf.FilterID()
if fID == "" {
return nil, errors.New("streamgate: resolved filter id is required")
}
if _, dup := out[fID]; dup {
return nil, fmt.Errorf("streamgate: duplicate resolved filter id in bind epoch: %s", fID)
}
out[fID] = rf
}
return out, nil
}
func indexApplicabilities(apps []FilterApplicability) (map[string]FilterApplicability, error) {
out := make(map[string]FilterApplicability, len(apps))
for _, app := range apps {
fID := app.FilterID()
if fID == "" {
return nil, errors.New("streamgate: applicability filter id is required")
}
if _, dup := out[fID]; dup {
return nil, fmt.Errorf("streamgate: duplicate applicability filter id in bind epoch: %s", fID)
}
out[fID] = app
}
return out, nil
}
func requireExactFilterIDs(planIDs []string, resolvedByID map[string]ResolvedFilter, appsByID map[string]FilterApplicability) error {
if len(resolvedByID) != len(planIDs) {
return fmt.Errorf("streamgate: resolved filter count mismatch: got %d, expected %d", len(resolvedByID), len(planIDs))
}
if len(appsByID) != len(planIDs) {
return fmt.Errorf("streamgate: applicability count mismatch: got %d, expected %d", len(appsByID), len(planIDs))
}
for _, id := range planIDs {
if _, ok := resolvedByID[id]; !ok {
return fmt.Errorf("streamgate: missing resolved filter for id %s", id)
}
if _, ok := appsByID[id]; !ok {
return fmt.Errorf("streamgate: missing applicability for id %s", id)
}
}
return nil
}
// EvidenceEpoch is an immutable snapshot of evidence state at the moment of
// a state transition (e.g. threshold reached, trigger fired). It carries
// only safe, externally-facing identification and filter applicabilities.
type EvidenceEpoch struct {
id uint64
channel string
mode FilterHoldMode
triggered bool
reason string
applicabilities map[string]FilterApplicability
}
// NewEvidenceEpoch creates an EvidenceEpoch with the given parameters.
func NewEvidenceEpoch(id uint64, channel string, mode FilterHoldMode, triggered bool, reason string) (EvidenceEpoch, error) {
return NewEvidenceEpochWithApplicabilities(id, channel, mode, triggered, reason, nil)
}
// NewEvidenceEpochWithApplicabilities creates an EvidenceEpoch carrying filter applicabilities.
func NewEvidenceEpochWithApplicabilities(
id uint64, channel string, mode FilterHoldMode, triggered bool, reason string,
apps []FilterApplicability,
) (EvidenceEpoch, error) {
if channel == "" {
return EvidenceEpoch{}, errors.New("streamgate: epoch channel is required")
}
if err := mode.Validate(); err != nil {
return EvidenceEpoch{}, err
}
appMap := make(map[string]FilterApplicability, len(apps))
for _, app := range apps {
fID := app.FilterID()
if fID == "" {
return EvidenceEpoch{}, errors.New("streamgate: epoch applicability filter id is required")
}
if _, dup := appMap[fID]; dup {
return EvidenceEpoch{}, fmt.Errorf("streamgate: duplicate epoch applicability filter id: %s", fID)
}
appMap[fID] = app
}
return EvidenceEpoch{
id: id,
channel: channel,
mode: mode,
triggered: triggered,
reason: reason,
applicabilities: appMap,
}, nil
}
// ID returns the epoch identifier.
func (e EvidenceEpoch) ID() uint64 { return e.id }
// Channel returns the epoch channel.
func (e EvidenceEpoch) Channel() string { return e.channel }
// Mode returns the epoch hold mode.
func (e EvidenceEpoch) Mode() FilterHoldMode { return e.mode }
// Triggered returns whether the epoch was a trigger event.
func (e EvidenceEpoch) Triggered() bool { return e.triggered }
// Reason returns the epoch reason string.
func (e EvidenceEpoch) Reason() string { return e.reason }
// ApplicabilityFor returns the applicability for the filter ID, if present.
func (e EvidenceEpoch) ApplicabilityFor(filterID string) (FilterApplicability, bool) {
if e.applicabilities == nil {
return FilterApplicability{}, false
}
app, ok := e.applicabilities[filterID]
return app, ok
}
// Applicabilities returns a defensive copy of all applicabilities in stable filter ID order.
func (e EvidenceEpoch) Applicabilities() []FilterApplicability {
if e.applicabilities == nil {
return nil
}
out := make([]FilterApplicability, 0, len(e.applicabilities))
for _, app := range e.applicabilities {
out = append(out, app)
}
sort.Slice(out, func(i, j int) bool { return out[i].filterID.value < out[j].filterID.value })
return out
}
// EvidenceTailSignal identifies the kind of signal emitted by evidence tail
// operations.
type EvidenceTailSignal string
const (
// EvidenceTailSignalNone indicates no signal was produced.
EvidenceTailSignalNone EvidenceTailSignal = "none"
// EvidenceTailSignalThreshold indicates the rolling rune threshold was
// reached and a batch is ready.
EvidenceTailSignalThreshold EvidenceTailSignal = "threshold"
// EvidenceTailSignalTrigger indicates the trigger event (terminal or
// fragment) was received.
EvidenceTailSignalTrigger EvidenceTailSignal = "trigger"
// EvidenceTailSignalBufferOverflow indicates the hard buffer limit was
// exceeded. No release is produced.
EvidenceTailSignalBufferOverflow EvidenceTailSignal = "overflow"
// EvidenceTailSignalReady indicates the state is ready for release
// (used when all accumulated events are eligible).
EvidenceTailSignalReady EvidenceTailSignal = "ready"
)
// EvidenceTailSignalOverflow is a typed signal for buffer overflow. It carries
// only the channel, bound, and cause code; no raw content.
type EvidenceTailSignalOverflow struct {
channel string
bound int
code string
}
// NewEvidenceTailSignalOverflow creates a typed overflow signal with the
// given parameters.
func NewEvidenceTailSignalOverflow(channel string, bound int, code string) (EvidenceTailSignalOverflow, error) {
if channel == "" {
return EvidenceTailSignalOverflow{}, errors.New("streamgate: overflow signal channel is required")
}
if bound <= 0 {
return EvidenceTailSignalOverflow{}, errors.New("streamgate: overflow signal bound must be positive")
}
if code == "" {
return EvidenceTailSignalOverflow{}, errors.New("streamgate: overflow signal cause code is required")
}
return EvidenceTailSignalOverflow{
channel: channel,
bound: bound,
code: code,
}, nil
}
// Channel returns the overflow channel.
func (s EvidenceTailSignalOverflow) Channel() string { return s.channel }
// Bound returns the overflow rune bound.
func (s EvidenceTailSignalOverflow) Bound() int { return s.bound }
// Code returns the overflow cause code.
func (s EvidenceTailSignalOverflow) Code() string { return s.code }
// EvidenceTail manages per-channel pending/look-behind state and tool-call
// fragment state for evidence holding. It is the core state machine that
// accumulates normalized events, validates UTF-8, counts Unicode runes, and
// produces epochs when thresholds or triggers are met.
type EvidenceTail struct {
plan EvidencePlan
epochCounter uint64
// channelState is keyed by channel name.
channelState map[string]*channelState
// epochs tracks epoch records by ID for validation during PrepareRelease
// and ConfirmRelease. Each epoch is bound to a specific channel and can
// only be prepared once.
epochs map[uint64]*epochRecord
// preparedByChannel tracks which channel has a pending prepared release
// for each epoch, preventing overlapping prepared tokens. Prepared
// replacements share this slot so a replacement and a release can never be
// outstanding on the same channel at once.
preparedByChannel map[string]string
// replacements tracks in-flight prepared replacements by token. Each record
// binds the replacement payload to the exact epoch, channel, and pending
// sequence range it settles.
replacements map[string]*replacementRecord
// tokenNonce is a monotonic counter that ensures every generated release
// token is globally unique. It prevents stale token reuse when the same
// epoch is re-prepared after zero or partial confirm.
tokenNonce uint64
}
// channelState holds the runtime state for a single channel.
type channelState struct {
// pendingEntries is the ordered list of pending entry records for this channel.
// Each entry tracks the event, its rune count, and associated fragment ID (if any).
pendingEntries []pendingEntry
// nextSequence is a monotonic counter that assigns an immutable identity
// to each pending entry. It is never decremented or reused, even after
// confirm removes entries from the front. This allows epoch records to
// anchor their snapshot to a fixed sequence range rather than a mutable
// slice index.
nextSequence int
// committedLookBehind is the bounded list of events already confirmed
// released downstream. Used for cross-boundary repeat detection.
committedLookBehind []NormalizedEvent
// committedCursor is the total number of events confirmed released
// downstream. This is a monotonic counter that increases with each
// successful ConfirmRelease and resets on replace.
committedCursor int
// pendingRunes is the cumulative Unicode rune count of pendingEntries.
pendingRunes int
// effectiveEvidenceRunes is the rolling evidence window used for
// look-behind trimming on this channel. For non-rolling modes it equals
// maxBufferRunes.
effectiveEvidenceRunes int
// fragmentState tracks incomplete tool-call fragments keyed by toolCallID.
fragmentState map[string]*fragmentState
// maxBufferRunes is the hard buffer limit for this channel.
maxBufferRunes int
}
// pendingEntry represents a single pending event with its metadata.
type pendingEntry struct {
event NormalizedEvent
runes int
kind EventKind
toolCallID string // empty if not a fragment
sequence int // monotonic sequence number
}
// fragmentState tracks an incomplete tool-call fragment.
// entries is the ordered list of pending entry sequences for this fragment ID.
// runes is the total rune count of all entries.
type fragmentState struct {
toolCallID string
entries []int // sequence numbers of pending entries
runes int
completed bool // true when CompleteFragment has been called for this ID
}
// epochRecord tracks an epoch's state through its lifecycle. Each epoch
// is bound to an exact pending-entry sequence range captured at creation
// time; remainingSnapshot returns only entries within that range that
// have not yet been confirmed.
//
// A record is valid for PrepareRelease only when:
// - consumed is false
// - invalidated is false
// - prepared and token carry a pending token whose channel matches
//
// Once invalidated (by recovery/continuation/terminal discard or a newer
// completion epoch), all token/prepared state is cleared and ConfirmRelease
// will reject any stale outstanding token with a deterministic error.
type epochRecord struct {
epochID uint64
channel string
mode FilterHoldMode // the hold mode at epoch creation
consumed bool // true after full confirm or replace
prepared bool // true after successful PrepareRelease
token string // the prepared release token, if any
snapshotSize int // number of events captured in the prepared snapshot
snapStartSeq int // first pending-entry monotonic sequence in the snapshot
snapEndSeq int // last pending-entry monotonic sequence (exclusive) in the snapshot
confirmedCount int // total events confirmed for this epoch across partial confirms
invalidated bool // true after ResetForReplace/PrepareContinuation/DiscardPendingForTerminal
}
// replacementRecord tracks a prepared replacement through its lifecycle. It
// anchors the substituted payload to one epoch, that epoch's channel, and the
// immutable pending-entry sequence range the epoch captured, so confirm can
// settle exactly that range and nothing else.
type replacementRecord struct {
token string
epochID uint64
channel string
snapStartSeq int // first pending-entry sequence being replaced
snapEndSeq int // last pending-entry sequence (exclusive) being replaced
events []NormalizedEvent // normalized replacement payload snapshot
}
// NewEvidenceTail creates a new EvidenceTail from the given compiled plan.
// It initializes per-channel state on first use.
func NewEvidenceTail(plan EvidencePlan) (*EvidenceTail, error) {
if err := plan.validate(); err != nil {
return nil, err
}
return &EvidenceTail{
plan: plan,
epochCounter: 0,
channelState: make(map[string]*channelState),
epochs: make(map[uint64]*epochRecord),
preparedByChannel: make(map[string]string),
replacements: make(map[string]*replacementRecord),
tokenNonce: 0,
}, nil
}
// validate checks the plan is consistent.
func (p EvidencePlan) validate() error {
for fID, b := range p.plan.allBindings {
if err := b.Requirement().Validate(); err != nil {
return fmt.Errorf("streamgate: compiled plan filter %s: %v", fID, err)
}
}
return nil
}
var (
// errStalePreparedRelease is returned by ConfirmRelease when the token
// has already been consumed, invalidated, or does not match the current
// prepared state for the epoch's channel.
errStalePreparedRelease = errors.New("streamgate: stale prepared release token")
// errPreparedSnapshotMismatch is returned when the current pending
// prefix does not match the sequence range captured at PrepareRelease.
errPreparedSnapshotMismatch = errors.New("streamgate: prepared snapshot prefix mismatch")
)
// nextObservationEpoch generates an epoch ID for a non-blocking observation
// event. No epoch record is created; the epoch cannot be prepared or
// confirmed. This is used for events that pass through (unsubscribed,
// observe-only, sub-threshold, overflow) and must never acquire release
// capability.
func (t *EvidenceTail) nextObservationEpoch(channel string, mode FilterHoldMode, triggered bool, reason string) (EvidenceEpoch, error) {
return t.nextObservationEpochWithApps(channel, mode, triggered, reason, nil)
}
func (t *EvidenceTail) nextObservationEpochWithApps(channel string, mode FilterHoldMode, triggered bool, reason string, apps []FilterApplicability) (EvidenceEpoch, error) {
t.epochCounter++
epoch, err := NewEvidenceEpochWithApplicabilities(t.epochCounter, channel, mode, triggered, reason, apps)
if err != nil {
return EvidenceEpoch{}, err
}
return epoch, nil
}
func (t *EvidenceTail) nextReadyEpochRecord(channel string, mode FilterHoldMode, triggered bool, reason string, firstSeq, lastSeqExclusive int) (EvidenceEpoch, *epochRecord, error) {
return t.nextReadyEpochRecordWithApps(channel, mode, triggered, reason, firstSeq, lastSeqExclusive, nil)
}
func (t *EvidenceTail) nextReadyEpochRecordWithApps(channel string, mode FilterHoldMode, triggered bool, reason string, firstSeq, lastSeqExclusive int, apps []FilterApplicability) (EvidenceEpoch, *epochRecord, error) {
t.epochCounter++
epoch, err := NewEvidenceEpochWithApplicabilities(t.epochCounter, channel, mode, triggered, reason, apps)
if err != nil {
return EvidenceEpoch{}, nil, err
}
record := &epochRecord{
epochID: t.epochCounter,
channel: channel,
mode: mode,
consumed: false,
prepared: false,
token: "",
snapshotSize: 0,
snapStartSeq: firstSeq,
snapEndSeq: lastSeqExclusive,
}
t.epochs[t.epochCounter] = record
return epoch, record, nil
}
// computeApplicabilities derives typed readiness and subscription status for each filter in the plan.
type applicabilityInput struct {
event *NormalizedEvent
isTerminalTrigger bool
isFragmentComplete bool
isHardBound bool
toolCallID string
eventBuffered bool
}
type pendingStats struct {
present bool
runes int
}
func subscribedPendingStats(cs *channelState, subscribedKinds []EventKind) pendingStats {
var stats pendingStats
if cs == nil || len(cs.pendingEntries) == 0 {
return stats
}
for _, entry := range cs.pendingEntries {
if isKindSubscribed(entry.kind, subscribedKinds) {
stats.present = true
stats.runes += entry.runes
}
}
return stats
}
func deriveTriggerReady(
b FilterHoldBinding,
req FilterHoldRequirement,
stats pendingStats,
cs *channelState,
input applicabilityInput,
) bool {
if input.isHardBound {
if b.BlocksRelease() && req.IsBlocking() {
return true
}
return false
}
if input.isTerminalTrigger {
switch req.Mode() {
case FilterHoldModeTerminalGate:
return input.event != nil && (input.event.Kind() == req.TriggerKind() || input.event.Kind() == EventKindTerminal || input.event.Kind() == EventKindProviderError)
case FilterHoldModeRolling, FilterHoldModeFragmentGate, FilterHoldModeNone:
return true
}
return false
}
if input.isFragmentComplete {
switch req.Mode() {
case FilterHoldModeFragmentGate:
return cs != nil && cs.isFragmentCompleted(input.toolCallID)
case FilterHoldModeRolling:
return stats.runes >= req.EvidenceRunes()
}
return false
}
switch req.Mode() {
case FilterHoldModeRolling:
return stats.runes >= req.EvidenceRunes()
case FilterHoldModeNone:
return input.event != nil && input.event.Kind() == EventKindProviderError && isKindSubscribed(EventKindProviderError, req.SubscribedKinds())
}
return false
}
func (t *EvidenceTail) computeApplicabilities(
channel string,
event *NormalizedEvent,
isTerminalTrigger bool,
isFragmentComplete bool,
isHardBound bool,
toolCallID string,
eventBuffered bool,
) ([]FilterApplicability, error) {
allBindings := t.plan.AllBindings()
if len(allBindings) == 0 {
return nil, nil
}
cs := t.channelState[channel]
input := applicabilityInput{
event: event,
isTerminalTrigger: isTerminalTrigger,
isFragmentComplete: isFragmentComplete,
isHardBound: isHardBound,
toolCallID: toolCallID,
eventBuffered: eventBuffered,
}
apps := make([]FilterApplicability, 0, len(allBindings))
for _, b := range allBindings {
fID := b.FilterID()
req := b.Requirement()
if req.Channel() != channel {
app, err := NewFilterApplicability(fID, false, false)
if err != nil {
return nil, err
}
apps = append(apps, app)
continue
}
stats := subscribedPendingStats(cs, req.SubscribedKinds())
if input.event != nil && !input.eventBuffered && isKindSubscribed(input.event.Kind(), req.SubscribedKinds()) {
stats.present = true
stats.runes += runeCountForEvent(*input.event)
}
subscribed := stats.present
ready := deriveTriggerReady(b, req, stats, cs, input)
app, err := NewFilterApplicability(fID, subscribed, ready)
if err != nil {
return nil, err
}
apps = append(apps, app)
}
return apps, nil
}
func (cs *channelState) unconfirmedSnapshot(record *epochRecord) []ReleaseEvent {
var releaseEvents []ReleaseEvent
for _, entry := range cs.pendingEntries {
if entry.sequence >= record.snapStartSeq && entry.sequence < record.snapEndSeq {
re, err := normalizedToReleaseEvent(entry.event)
if err == nil {
releaseEvents = append(releaseEvents, re)
}
}
}
return releaseEvents
}
func (t *EvidenceTail) validatePreparedOwnership(record *epochRecord, token string) error {
if record.invalidated {
return errStalePreparedRelease
}
if t.preparedByChannel[record.channel] != token {
return errStalePreparedRelease
}
return nil
}
func (cs *channelState) pendingPrefixMatchesRange(n int, snapStartSeq int) bool {
if n > len(cs.pendingEntries) {
return false
}
for i := 0; i < n; i++ {
if cs.pendingEntries[i].sequence != snapStartSeq+i {
return false
}
}
return true
}
// rangeEntryCount returns how many pending entries still fall inside the
// half-open sequence range captured by an epoch snapshot.
func (cs *channelState) rangeEntryCount(startSeq, endSeqExclusive int) int {
count := 0
for _, entry := range cs.pendingEntries {
if entry.sequence >= startSeq && entry.sequence < endSeqExclusive {
count++
}
}
return count
}
// dropFragmentEntriesBelow removes fragment bookkeeping for pending entries
// whose sequence is already settled. Absolute sequences are preserved for the
// remaining entries so subsequent partial/full settlements stay correct.
func (cs *channelState) dropFragmentEntriesBelow(settledEndSeq int) {
for id, fragment := range cs.fragmentState {
var remaining []int
for _, seq := range fragment.entries {
if seq >= settledEndSeq {
remaining = append(remaining, seq)
}
}
if len(remaining) == 0 {
delete(cs.fragmentState, id)
} else {
fragment.entries = remaining
}
}
}
func (t *EvidenceTail) getOrCreateChannel(channel string, maxBuffer int, evidenceRunes int) (*channelState, error) {
cs, ok := t.channelState[channel]
if !ok {
effective := evidenceRunes
if effective <= 0 {
effective = maxBuffer
}
cs = &channelState{
pendingEntries: make([]pendingEntry, 0),
fragmentState: make(map[string]*fragmentState),
maxBufferRunes: maxBuffer,
effectiveEvidenceRunes: effective,
}
t.channelState[channel] = cs
}
return cs, nil
}
func (t *EvidenceTail) Append(event NormalizedEvent) (EvidenceEpoch, EvidenceTailSignal, error) {
if err := event.Validate(); err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, fmt.Errorf("streamgate: append event validate: %v", err)
}
channel := event.Channel()
req, hasBlocking := t.plan.BlockingRequirementFor(channel)
if err := validateEventUTF8(event); err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
if !hasBlocking {
apps, err := t.computeApplicabilities(channel, &event, false, false, false, "", false)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
epoch, err := t.nextObservationEpochWithApps(channel, FilterHoldModeNone, false, "pass_through_no_requirement", apps)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
return epoch, EvidenceTailSignalReady, nil
}
if event.Kind() == EventKindTerminal || event.Kind() == EventKindProviderError || t.plan.IsTerminalTriggerForChannel(channel, event.Kind()) {
cs, err := t.getOrCreateChannel(channel, t.plan.MaxBufferRunes(channel), t.plan.RollingEvidenceRunes(channel))
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
safeEnd := cs.nextSequence
firstSeq := cs.firstPendingSequence()
apps, err := t.computeApplicabilities(channel, &event, true, false, false, "", false)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
epoch, _, err := t.nextReadyEpochRecordWithApps(
channel, req.Mode(), true, "configured_trigger",
firstSeq, safeEnd, apps,
)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
return epoch, EvidenceTailSignalTrigger, nil
}
blockingKinds := t.plan.BlockingKinds(channel)
observeKinds := t.plan.ObserveKinds(channel)
if !isKindSubscribed(event.Kind(), blockingKinds) {
reason := "event_not_subscribed"
if isKindSubscribed(event.Kind(), observeKinds) {
reason = "observe_only_pass_through"
}
apps, err := t.computeApplicabilities(channel, &event, false, false, false, "", false)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
epoch, err := t.nextObservationEpochWithApps(channel, req.Mode(), false, reason, apps)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
return epoch, EvidenceTailSignalNone, nil
}
cs, err := t.getOrCreateChannel(channel, t.plan.MaxBufferRunes(channel), t.plan.RollingEvidenceRunes(channel))
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
runesAdded := runeCountForEvent(event)
newPendingRunes := cs.pendingRunes + runesAdded
if newPendingRunes > cs.maxBufferRunes {
apps, err := t.computeApplicabilities(channel, &event, false, false, true, "", false)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
epoch, err := t.nextObservationEpochWithApps(channel, req.Mode(), false, "overflow_beyond_buffer_limit", apps)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
return epoch, EvidenceTailSignalBufferOverflow, nil
}
seq := cs.nextSequence
cs.nextSequence++
entry := pendingEntry{
event: cloneNormalizedEvent(event),
runes: runesAdded,
kind: event.Kind(),
sequence: seq,
}
if event.Kind() == EventKindToolCallFragment {
tc, err := event.AsToolCallFragment()
if err == nil {
entry.toolCallID = tc.ID
fragment := cs.fragmentState[tc.ID]
if fragment == nil {
fragment = &fragmentState{
toolCallID: tc.ID,
entries: []int{},
runes: 0,
completed: false,
}
cs.fragmentState[tc.ID] = fragment
}
newFragmentRunes := fragment.runes + runesAdded
if newFragmentRunes > cs.maxBufferRunes {
return EvidenceEpoch{}, EvidenceTailSignalBufferOverflow, nil
}
fragment.entries = append(fragment.entries, entry.sequence)
fragment.runes = newFragmentRunes
}
}
cs.pendingEntries = append(cs.pendingEntries, entry)
cs.pendingRunes = newPendingRunes
apps, err := t.computeApplicabilities(channel, &event, false, false, false, "", true)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
anyReady := false
hasRollingReady := false
for _, app := range apps {
if app.TriggerReady() {
anyReady = true
for _, b := range t.plan.BindingsForChannel(channel) {
if b.FilterID() == app.FilterID() && b.Requirement().Mode() == FilterHoldModeRolling {
hasRollingReady = true
}
}
}
}
if anyReady {
safeEnd := cs.nextSequence
firstSeq := cs.firstPendingSequence()
epoch, _, err := t.nextReadyEpochRecordWithApps(
channel, req.Mode(), true, "threshold_or_trigger_reached",
firstSeq, safeEnd, apps,
)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
sig := EvidenceTailSignalTrigger
if hasRollingReady {
sig = EvidenceTailSignalThreshold
}
return epoch, sig, nil
}
epoch, err := t.nextObservationEpochWithApps(channel, req.Mode(), false, "sub_threshold_accumulating", apps)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, err
}
return epoch, EvidenceTailSignalNone, nil
}
func (cs *channelState) firstPendingSequence() int {
if len(cs.pendingEntries) == 0 {
return cs.nextSequence
}
return cs.pendingEntries[0].sequence
}
func (t *EvidenceTail) CompleteFragment(
channel, toolCallID string,
) (EvidenceEpoch, EvidenceTailSignal, bool, error) {
if channel == "" {
return EvidenceEpoch{}, EvidenceTailSignalNone, false, errors.New("streamgate: complete fragment channel is required")
}
if toolCallID == "" {
return EvidenceEpoch{}, EvidenceTailSignalNone, false, errors.New("streamgate: complete fragment tool call id is required")
}
if !t.plan.HasModeForChannel(channel, FilterHoldModeFragmentGate) {
return EvidenceEpoch{}, EvidenceTailSignalNone, false, fmt.Errorf("streamgate: channel %s is not fragment_gate mode", channel)
}
cs, ok := t.channelState[channel]
if !ok {
return EvidenceEpoch{}, EvidenceTailSignalNone, false, nil
}
fragment, exists := cs.fragmentState[toolCallID]
if !exists || fragment.completed {
return EvidenceEpoch{}, EvidenceTailSignalNone, false, nil
}
fragment.completed = true
safeLen := cs.fragmentSafePrefixLength()
if safeLen == 0 {
return EvidenceEpoch{}, EvidenceTailSignalNone, true, nil
}
t.markChannelEpochsInvalidated(channel)
delete(t.preparedByChannel, channel)
t.dropChannelReplacements(channel)
firstSeq := cs.firstPendingSequence()
safeEndSeq := firstSeq + safeLen
apps, err := t.computeApplicabilities(channel, nil, false, true, false, toolCallID, false)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, true, err
}
epoch, _, err := t.nextReadyEpochRecordWithApps(
channel, FilterHoldModeFragmentGate, true, "fragment_completion_safe_prefix",
firstSeq, safeEndSeq, apps,
)
if err != nil {
return EvidenceEpoch{}, EvidenceTailSignalNone, true, err
}
return epoch, EvidenceTailSignalTrigger, true, nil
}
// isFragmentCompleted returns true when the fragment for the given ID is completed.
func (cs *channelState) isFragmentCompleted(toolCallID string) bool {
fragment, exists := cs.fragmentState[toolCallID]
return exists && fragment.completed
}
// fragmentSafePrefixLength returns the index in pendingEntries up to which
// all fragments are completed and the prefix is contiguous. Returns 0 if no
// fragments are completed.
func (cs *channelState) fragmentSafePrefixLength() int {
if len(cs.pendingEntries) == 0 {
return 0
}
// Find the first incomplete fragment ID.
firstIncompleteID := ""
for _, entry := range cs.pendingEntries {
if entry.kind == EventKindToolCallFragment && !cs.isFragmentCompleted(entry.toolCallID) {
firstIncompleteID = entry.toolCallID
break
}
}
if firstIncompleteID == "" {
// All fragments are completed.
return len(cs.pendingEntries)
}
// Find the index of the first incomplete fragment.
for i, entry := range cs.pendingEntries {
if entry.toolCallID == firstIncompleteID {
return i
}
}
return 0
}
// PreparedRelease is an opaque token representing a prepared release that
// has not yet been confirmed. It is immutable and cannot be tampered with.
type PreparedRelease struct {
token string
releaseEvents []ReleaseEvent
}
// Token returns the prepared release token string.
func (pr PreparedRelease) Token() string { return pr.token }
// ReleaseEvents returns a defensive copy of the prepared release events.
func (pr PreparedRelease) ReleaseEvents() []ReleaseEvent {
if pr.releaseEvents == nil {
return nil
}
out := make([]ReleaseEvent, len(pr.releaseEvents))
copy(out, pr.releaseEvents)
return out
}
// ReleaseConfirmation carries the number of events successfully released
// during a confirm operation.
type ReleaseConfirmation struct {
ReleasedEvents int
}
// PrepareRelease creates a prepared release for the given epoch ID.
// Only epoch records that were created via a ready transition (threshold
// reached, trigger fired, fragment completion) carry release capability.
// Observation-only epochs (unsubscribed, overflow, observe-only, sub-threshold)
// have no record and cannot be prepared.
//
// PrepareRelease validates the epoch record exists and is not consumed or
// invalidated, generates a release snapshot from the epoch's fixed sequence
// range using unconfirmedSnapshot (which respects partial confirms), and
// creates an opaque token. It does NOT modify pending/look-behind/cursor
// state.
//
// PrepareRelease rejects: unknown/observation epochs, already consumed/
// invalidated epochs, overlapping prepared tokens for the same channel,
// and already-prepared epochs.
// HasUnconfirmedEvents returns true if the specified epoch has unconfirmed release payload events.
func (t *EvidenceTail) HasUnconfirmedEvents(epochID uint64) bool {
if t == nil || t.epochs == nil {
return false
}
record, ok := t.epochs[epochID]
if !ok || record.consumed || record.invalidated {
return false
}
channelCS := t.channelState[record.channel]
if channelCS == nil {
return false
}
return len(channelCS.unconfirmedSnapshot(record)) > 0
}
func (t *EvidenceTail) PrepareRelease(epochID uint64) (PreparedRelease, error) {
// Validate epoch record exists and is not consumed or invalidated.
record, ok := t.epochs[epochID]
if !ok {
return PreparedRelease{}, fmt.Errorf("streamgate: prepare release unknown or observation-only epoch %d", epochID)
}
if record.consumed || record.invalidated {
return PreparedRelease{}, fmt.Errorf("streamgate: prepare release epoch %d already consumed or invalidated", epochID)
}
// Overlapping prepared token on the same channel is rejected.
if t.preparedByChannel[record.channel] != "" {
return PreparedRelease{}, fmt.Errorf("streamgate: prepare release channel %s already has prepared token %s", record.channel, t.preparedByChannel[record.channel])
}
if record.prepared && record.token != "" {
return PreparedRelease{}, fmt.Errorf("streamgate: prepare release epoch %d already prepared with token %s", epochID, record.token)
}
channel := record.channel
// Get the release snapshot from the epoch's fixed sequence range.
// This uses monotonic sequence identity so it is safe after partial confirm.
channelCS := t.channelState[channel]
if channelCS == nil {
return PreparedRelease{}, errors.New("streamgate: prepare release no channel state for epoch channel")
}
releaseEvents := channelCS.unconfirmedSnapshot(record)
if len(releaseEvents) == 0 {
return PreparedRelease{}, fmt.Errorf("streamgate: prepare release epoch %d has no unconfirmed events", epochID)
}
// Generate an opaque, globally unique token with a monotonic nonce.
// This prevents stale token reuse when the same epoch is re-prepared
// after zero or partial confirm (the snapshot content could be identical).
t.tokenNonce++
token := fmt.Sprintf("prepared-%d-%d-%d", epochID, len(releaseEvents), t.tokenNonce)
// Store snapshot size for downstream ConfirmRelease validation.
record.snapshotSize = len(releaseEvents)
record.prepared = true
record.token = token
t.preparedByChannel[channel] = token
return PreparedRelease{
token: token,
releaseEvents: releaseEvents,
}, nil
}
// ConfirmRelease applies the confirmed number of released events to the
// committed look-behind and cursor. Only this method modifies pending/
// look-behind/cursor state. Zero, partial, and full confirmations are
// supported. Stale tokens and duplicate confirms are rejected.
//
// ConfirmRelease validates: token format and existence, that the token
// matches the epoch's prepared token, that ReleasedEvents does not exceed
// the prepared count, and that the epoch has not been invalidated.
//
// ConfirmRelease validates:
// - The token matches the epoch record and matches the current prepared
// token for the epoch's channel.
// - The epoch is not invalidated.
// - The pending prefix starting at the epoch's snapStartSeq matches
// the expected monotonic sequence.
// - ReleasedEvents does not exceed the prepared count.
func (t *EvidenceTail) ConfirmRelease(token string, confirmation ReleaseConfirmation) error {
if token == "" {
return errors.New("streamgate: confirm release token is required")
}
if !strings.HasPrefix(token, "prepared-") {
return errors.New("streamgate: confirm release invalid token format")
}
// Find the epoch record associated with this token.
var targetEpoch *epochRecord
for _, record := range t.epochs {
if record.token == token {
targetEpoch = record
break
}
}
if targetEpoch == nil {
return errors.New("streamgate: confirm release token not found")
}
if targetEpoch.consumed {
return fmt.Errorf("streamgate: confirm release epoch %d already consumed", targetEpoch.epochID)
}
// Validate ownership: token must match current prepared state and
// epoch must not be invalidated.
if err := t.validatePreparedOwnership(targetEpoch, token); err != nil {
return err
}
channel := targetEpoch.channel
cs, ok := t.channelState[channel]
if !ok {
return errors.New("streamgate: confirm release channel state not found")
}
n := confirmation.ReleasedEvents
if n < 0 {
return errors.New("streamgate: confirm release released events must be non-negative")
}
if n > targetEpoch.snapshotSize {
return fmt.Errorf("streamgate: confirm release %d exceeds prepared snapshot %d", n, targetEpoch.snapshotSize)
}
// Validate that the pending prefix matches the epoch's sequence range.
// This guards against a prepared token being confirmed after the pending
// prefix has diverged from the original snapshot (e.g., post-append
// contamination).
if !cs.pendingPrefixMatchesRange(n, targetEpoch.snapStartSeq) {
return errPreparedSnapshotMismatch
}
// Move confirmed events to look-behind.
confirmed := make([]NormalizedEvent, n)
for i := 0; i < n; i++ {
confirmed[i] = cloneNormalizedEvent(cs.pendingEntries[i].event)
}
cs.committedLookBehind = append(cs.committedLookBehind, confirmed...)
// Remove confirmed from pending.
cs.pendingEntries = cs.pendingEntries[n:]
// Adjust rune count.
removedRunes := 0
for i := 0; i < n; i++ {
removedRunes += runeCountForEvent(confirmed[i])
}
cs.pendingRunes -= removedRunes
// Update committed cursor.
cs.committedCursor += n
// Remove confirmed fragment states (if all their entries are confirmed).
// Use the absolute confirmed range so that non-zero sequence fragments
// clean up correctly: only entries already covered by this confirm are
// removed, and their absolute sequences are preserved for subsequent
// partial/full confirms.
cs.dropFragmentEntriesBelow(targetEpoch.snapStartSeq + n)
// Trim look-behind to bounded size using the channel's effective evidence
// window (evidence_runes for rolling, maxBufferRunes otherwise).
cs.committedLookBehind = trimLookBehind(cs.committedLookBehind, cs.effectiveEvidenceRunes)
// Mark consumed when the full snapshot has been confirmed.
if n >= targetEpoch.snapshotSize {
targetEpoch.consumed = true
}
// Advance the epoch's unconfirmed start so that a subsequent re-prepare
// or partial/full confirm sees only the remaining suffix of the same epoch.
targetEpoch.snapStartSeq += n
targetEpoch.confirmedCount += n
// Clear prepared token regardless of completion.
targetEpoch.prepared = false
targetEpoch.token = ""
delete(t.preparedByChannel, channel)
return nil
}
// preparedReplacement is an opaque token representing a prepared replacement
// that has not yet been confirmed. It is immutable and cannot be tampered
// with, and it is only meaningful for the epoch it was prepared against.
type preparedReplacement struct {
token string
releaseEvents []ReleaseEvent
}
// Token returns the prepared replacement token string.
func (pr preparedReplacement) Token() string { return pr.token }
// ReleaseEvents returns a defensive copy of the replacement release events.
func (pr preparedReplacement) ReleaseEvents() []ReleaseEvent {
if pr.releaseEvents == nil {
return nil
}
out := make([]ReleaseEvent, len(pr.releaseEvents))
copy(out, pr.releaseEvents)
return out
}
// prepareReplacement validates a replacement payload against the epoch it
// substitutes and reserves the epoch's channel, without mutating pending,
// look-behind, or cursor state. Every failure mode is checked here so a caller
// can hand the payload to the sink knowing the tail is still untouched:
// the epoch must exist as a releasable record that is neither consumed nor
// invalidated, its channel must not already carry an outstanding prepared
// token, the captured sequence range must still match the pending prefix, and
// the payload must be non-empty with every event releasable and on the epoch's
// own channel.
func (t *EvidenceTail) prepareReplacement(epochID uint64, events []NormalizedEvent) (preparedReplacement, error) {
if len(events) == 0 {
return preparedReplacement{}, errors.New("streamgate: prepare replacement requires at least one event")
}
record, ok := t.epochs[epochID]
if !ok {
return preparedReplacement{}, fmt.Errorf("streamgate: prepare replacement unknown or observation-only epoch %d", epochID)
}
if record.consumed || record.invalidated {
return preparedReplacement{}, fmt.Errorf("streamgate: prepare replacement epoch %d already consumed or invalidated", epochID)
}
if t.preparedByChannel[record.channel] != "" {
return preparedReplacement{}, fmt.Errorf("streamgate: prepare replacement channel %s already has prepared token %s", record.channel, t.preparedByChannel[record.channel])
}
if record.snapEndSeq < record.snapStartSeq {
return preparedReplacement{}, fmt.Errorf("streamgate: prepare replacement epoch %d has an inverted sequence range", epochID)
}
channel := record.channel
cs := t.channelState[channel]
if cs == nil {
return preparedReplacement{}, errors.New("streamgate: prepare replacement no channel state for epoch channel")
}
expectedCount := record.snapEndSeq - record.snapStartSeq
targetCount := cs.rangeEntryCount(record.snapStartSeq, record.snapEndSeq)
if targetCount != expectedCount {
return preparedReplacement{}, errPreparedSnapshotMismatch
}
if expectedCount > 0 && !cs.pendingPrefixMatchesRange(expectedCount, record.snapStartSeq) {
return preparedReplacement{}, errPreparedSnapshotMismatch
}
releaseEvents := make([]ReleaseEvent, len(events))
payload := make([]NormalizedEvent, len(events))
for i, ev := range events {
if err := ev.Validate(); err != nil {
return preparedReplacement{}, fmt.Errorf("streamgate: prepare replacement event %d: %v", i, err)
}
if ev.Channel() != channel {
return preparedReplacement{}, fmt.Errorf("streamgate: prepare replacement event %d channel %s does not match epoch channel %s", i, ev.Channel(), channel)
}
relEv, err := normalizedToReleaseEvent(ev)
if err != nil {
return preparedReplacement{}, err
}
releaseEvents[i] = relEv
payload[i] = cloneNormalizedEvent(ev)
}
// The nonce keeps the token globally unique so a stale replacement token
// can never confirm a later replacement of the same epoch. The distinct
// prefix keeps replacement tokens out of ConfirmRelease.
t.tokenNonce++
token := fmt.Sprintf("replacement-%d-%d-%d", epochID, len(releaseEvents), t.tokenNonce)
t.replacements[token] = &replacementRecord{
token: token,
epochID: epochID,
channel: channel,
snapStartSeq: record.snapStartSeq,
snapEndSeq: record.snapEndSeq,
events: payload,
}
t.preparedByChannel[channel] = token
return preparedReplacement{
token: token,
releaseEvents: releaseEvents,
}, nil
}
// confirmReplacement settles a prepared replacement against the sink progress
// the boundary actually reported. The whole target range is consumed as
// replaced regardless of progress - the original events were superseded and can
// never reach the sink - but only the replacement prefix the sink accepted
// becomes committed look-behind and advances the cursor, so zero, partial, and
// full progress all leave the committed view equal to downstream reality.
//
// Only the target epoch's own range is touched: post-snapshot entries on the
// same channel stay pending, unrelated channels keep their pending entries,
// prepared tokens, and epoch records, and only same-channel epochs whose range
// overlaps the consumed original range are invalidated.
func (t *EvidenceTail) confirmReplacement(token string, confirmation ReleaseConfirmation) error {
if token == "" {
return errors.New("streamgate: confirm replacement token is required")
}
if !strings.HasPrefix(token, "replacement-") {
return errors.New("streamgate: confirm replacement invalid token format")
}
replacement, ok := t.replacements[token]
if !ok {
return errStalePreparedRelease
}
record, ok := t.epochs[replacement.epochID]
if !ok || record.consumed || record.invalidated {
delete(t.replacements, token)
return errStalePreparedRelease
}
if t.preparedByChannel[replacement.channel] != token {
return errStalePreparedRelease
}
cs := t.channelState[replacement.channel]
if cs == nil {
return errors.New("streamgate: confirm replacement channel state not found")
}
n := confirmation.ReleasedEvents
if n < 0 {
return errors.New("streamgate: confirm replacement released events must be non-negative")
}
if n > len(replacement.events) {
return fmt.Errorf("streamgate: confirm replacement %d exceeds prepared replacement %d", n, len(replacement.events))
}
expectedCount := replacement.snapEndSeq - replacement.snapStartSeq
targetCount := cs.rangeEntryCount(replacement.snapStartSeq, replacement.snapEndSeq)
if targetCount != expectedCount {
return errPreparedSnapshotMismatch
}
if expectedCount > 0 && !cs.pendingPrefixMatchesRange(expectedCount, replacement.snapStartSeq) {
return errPreparedSnapshotMismatch
}
// Consume the superseded original range without any downstream write.
if targetCount > 0 {
removedRunes := 0
for i := 0; i < targetCount; i++ {
removedRunes += cs.pendingEntries[i].runes
}
cs.pendingEntries = cs.pendingEntries[targetCount:]
cs.pendingRunes -= removedRunes
cs.dropFragmentEntriesBelow(replacement.snapEndSeq)
}
// Only the accepted replacement prefix becomes committed evidence.
if n > 0 {
released := make([]NormalizedEvent, n)
for i := 0; i < n; i++ {
released[i] = cloneNormalizedEvent(replacement.events[i])
}
cs.committedLookBehind = append(cs.committedLookBehind, released...)
cs.committedCursor += n
cs.committedLookBehind = trimLookBehind(cs.committedLookBehind, cs.effectiveEvidenceRunes)
}
record.consumed = true
record.prepared = false
record.token = ""
record.snapshotSize = 0
record.snapStartSeq = replacement.snapEndSeq
t.invalidateOverlappingEpochs(replacement)
delete(t.preparedByChannel, replacement.channel)
delete(t.replacements, token)
return nil
}
// invalidateOverlappingEpochs invalidates the still-open epoch records whose
// captured range on the replaced channel overlaps the consumed original range,
// so no stale token can confirm content that a replacement already superseded.
// Records on other channels, already settled records, and same-channel records
// covering only post-snapshot entries are left untouched.
func (t *EvidenceTail) invalidateOverlappingEpochs(replacement *replacementRecord) {
for _, record := range t.epochs {
if record.epochID == replacement.epochID || record.channel != replacement.channel {
continue
}
if record.consumed || record.invalidated {
continue
}
if record.snapStartSeq >= replacement.snapEndSeq || record.snapEndSeq <= replacement.snapStartSeq {
continue
}
record.invalidated = true
record.prepared = false
record.token = ""
record.snapshotSize = 0
}
}
// ResetForReplace clears all attempt-local state including pending entries,
// look-behind, fragment state, prepared tokens, and the token nonce. Used
// when a new attempt replaces the current one.
func (t *EvidenceTail) ResetForReplace() {
t.channelState = make(map[string]*channelState)
t.epochs = make(map[uint64]*epochRecord)
t.preparedByChannel = make(map[string]string)
t.replacements = make(map[string]*replacementRecord)
t.tokenNonce = 0
}
// PrepareContinuation preserves committed look-behind and cursor while
// discarding pending entries and prepared tokens. Used when resuming after
// a recovery continuation.
func (t *EvidenceTail) PrepareContinuation() {
for _, cs := range t.channelState {
cs.pendingEntries = nil
cs.pendingRunes = 0
cs.fragmentState = make(map[string]*fragmentState)
// committedLookBehind, committedCursor are preserved.
}
t.markEpochsInvalidated()
t.preparedByChannel = make(map[string]string)
t.replacements = make(map[string]*replacementRecord)
}
// DiscardPendingForTerminal discards all pending entries and prepared tokens
// without releasing anything. Used before terminal error or when the
// terminal gate overflows. No downstream write occurs.
func (t *EvidenceTail) DiscardPendingForTerminal() {
for _, cs := range t.channelState {
cs.pendingEntries = nil
cs.pendingRunes = 0
cs.fragmentState = make(map[string]*fragmentState)
// committedLookBehind and committedCursor are preserved.
}
t.markEpochsInvalidated()
t.preparedByChannel = make(map[string]string)
t.replacements = make(map[string]*replacementRecord)
}
// markEpochsInvalidated marks all epoch records as invalidated and clears
// all token/prepared/snapshot state. This prevents any stale outstanding
// token from being confirmed after recovery/continuation/terminal discard.
func (t *EvidenceTail) markEpochsInvalidated() {
for _, record := range t.epochs {
record.invalidated = true
record.prepared = false
record.token = ""
record.snapshotSize = 0
}
}
// dropChannelReplacements removes the prepared replacement records belonging to
// the given channel. Records on other channels are left untouched.
func (t *EvidenceTail) dropChannelReplacements(channel string) {
for token, replacement := range t.replacements {
if replacement.channel == channel {
delete(t.replacements, token)
}
}
}
// markChannelEpochsInvalidated marks all epoch records belonging to the given
// channel as invalidated and clears all token/prepared/snapshot state for
// those records. This prevents stale outstanding tokens on the same channel
// from being confirmed after fragment completion supersedes an older safe
// prefix snapshot. It does NOT affect epoch records on other channels.
func (t *EvidenceTail) markChannelEpochsInvalidated(channel string) {
for _, record := range t.epochs {
if record.channel == channel {
record.invalidated = true
record.prepared = false
record.token = ""
record.snapshotSize = 0
}
}
}
// validateEventUTF8 checks that text/reasoning/tool argument content in the
// event is valid UTF-8. Invalid UTF-8 is rejected with a stable error.
func validateEventUTF8(event NormalizedEvent) error {
var texts []string
switch event.Kind() {
case EventKindTextDelta:
text, err := event.AsTextDelta()
if err != nil {
return err
}
texts = append(texts, text)
case EventKindReasoningDelta:
text, err := event.AsReasoningDelta()
if err != nil {
return err
}
texts = append(texts, text)
case EventKindToolCallFragment:
tc, err := event.AsToolCallFragment()
if err != nil {
return err
}
texts = append(texts, tc.Arguments)
}
for _, text := range texts {
if text != "" && !utf8.ValidString(text) {
return fmt.Errorf("streamgate: event %s contains invalid UTF-8", event.Kind())
}
}
return nil
}
// runeCountForEvent returns the Unicode rune count for an event's text
// content. It operates on valid UTF-8 strings only.
func runeCountForEvent(event NormalizedEvent) int {
var total int
switch event.Kind() {
case EventKindTextDelta:
text, _ := event.AsTextDelta()
total += utf8.RuneCountInString(text)
case EventKindReasoningDelta:
text, _ := event.AsReasoningDelta()
total += utf8.RuneCountInString(text)
case EventKindToolCallFragment:
tc, _ := event.AsToolCallFragment()
total += utf8.RuneCountInString(tc.Arguments)
}
return total
}
// normalizedToReleaseEvent converts a NormalizedEvent to a ReleaseEvent.
// Only releasable kinds (text_delta, reasoning_delta, tool_call_fragment)
// are supported.
func normalizedToReleaseEvent(ev NormalizedEvent) (ReleaseEvent, error) {
switch ev.Kind() {
case EventKindTextDelta:
text, err := ev.AsTextDelta()
if err != nil {
return ReleaseEvent{}, err
}
return NewReleaseTextDeltaEvent(ev.Channel(), text, ev.Timestamp())
case EventKindReasoningDelta:
text, err := ev.AsReasoningDelta()
if err != nil {
return ReleaseEvent{}, err
}
return NewReleaseReasoningDeltaEvent(ev.Channel(), text, ev.Timestamp())
case EventKindToolCallFragment:
tc, err := ev.AsToolCallFragment()
if err != nil {
return ReleaseEvent{}, err
}
return NewReleaseToolCallFragmentEvent(ev.Channel(), tc.ID, tc.Name, tc.Arguments, ev.Timestamp())
default:
return ReleaseEvent{}, fmt.Errorf("streamgate: cannot release event kind %s", ev.Kind())
}
}
// CommittedCursor returns the total number of events confirmed released
// downstream for the given channel. Returns 0 if the channel has no state.
func (t *EvidenceTail) CommittedCursor(channel string) int {
cs, ok := t.channelState[channel]
if !ok {
return 0
}
return cs.committedCursor
}
// EffectiveLookBehind returns a defensive deep copy of the bounded look-behind
// events for the given channel. Returns nil if the channel has no state.
func (t *EvidenceTail) EffectiveLookBehind(channel string) []NormalizedEvent {
cs, ok := t.channelState[channel]
if !ok {
return nil
}
if len(cs.committedLookBehind) == 0 {
return nil
}
out := make([]NormalizedEvent, len(cs.committedLookBehind))
for i, ev := range cs.committedLookBehind {
out[i] = cloneNormalizedEvent(ev)
}
return out
}
// sliceEventAtRuneBoundary returns a new event with content sliced to at most
// maxRunes runes from the trailing edge. Returns an empty NormalizedEvent if
// the kind cannot be sliced (terminal, response_start, etc.).
func sliceEventAtRuneBoundary(ev NormalizedEvent, maxRunes int) (NormalizedEvent, error) {
var text string
switch ev.Kind() {
case EventKindTextDelta:
text, _ = ev.AsTextDelta()
case EventKindReasoningDelta:
text, _ = ev.AsReasoningDelta()
case EventKindToolCallFragment:
tc, _ := ev.AsToolCallFragment()
text = tc.Arguments
default:
return NormalizedEvent{}, nil
}
runes := []rune(text)
if len(runes) <= maxRunes {
return ev, nil
}
sliced := string(runes[len(runes)-maxRunes:])
switch ev.Kind() {
case EventKindTextDelta:
return NewTextDeltaEvent(ev.Channel(), sliced, ev.Timestamp())
case EventKindReasoningDelta:
return NewReasoningDeltaEvent(ev.Channel(), sliced, ev.Timestamp())
case EventKindToolCallFragment:
tc, _ := ev.AsToolCallFragment()
return NewToolCallFragmentEvent(ev.Channel(), tc.ID, tc.Name, sliced, ev.Timestamp())
}
return NormalizedEvent{}, nil
}
// trimLookBehind trims the look-behind to at most maxRunes worth of events,
// keeping the most recent events. Single oversized events are sliced at
// Unicode rune boundaries to preserve partial content as a defensive snapshot.
func trimLookBehind(events []NormalizedEvent, maxRunes int) []NormalizedEvent {
if len(events) == 0 || maxRunes <= 0 {
return events
}
// Calculate rune counts for each event.
runeCounts := make([]int, len(events))
totalRunes := 0
for i, ev := range events {
runeCounts[i] = runeCountForEvent(ev)
totalRunes += runeCounts[i]
}
if totalRunes <= maxRunes {
return events
}
// Keep from the end until we're within the limit.
var result []NormalizedEvent
remainingRunes := maxRunes
for i := len(events) - 1; i >= 0; i-- {
if remainingRunes >= runeCounts[i] {
result = append([]NormalizedEvent{events[i]}, result...)
remainingRunes -= runeCounts[i]
} else {
// Slice this event at rune boundary to preserve partial content.
sliced, err := sliceEventAtRuneBoundary(events[i], remainingRunes)
if err == nil && sliced.Kind() != "" {
result = append([]NormalizedEvent{sliced}, result...)
}
break
}
}
return result
}