- Refactor plan, code-review, finalize-task-routing, refine-local-plans, router skills - Add agent-workflow-loop-orchestration skill and plan agent configs - Update roadmap: knowledge-tool-optimization milestones, stream-evidence-gate-core SDD - Add stream-evidence-gate-core task, archive, and Go streamgate package - Update dev-test inventory (edge/node smoke), agent-contract, edge-local-dev-guide - Deprecate USER_REVIEW for output-validation-filters SDD
494 lines
18 KiB
Go
494 lines
18 KiB
Go
package streamgate
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"sort"
|
|
"time"
|
|
)
|
|
|
|
// FilterEvaluator is the minimal execution seam that the parallel evaluator
|
|
// consumes. Implementations are produced by the filter-registry Task and
|
|
// adapted to this seam. The seam is intentionally small: only an immutable
|
|
// identity and a synchronous, context-aware evaluation call.
|
|
type FilterEvaluator interface {
|
|
// ID returns a stable, ASCII-only identifier for this evaluator.
|
|
ID() string
|
|
|
|
// Evaluate runs the filter against the given batch and returns a
|
|
// sanitized FilterDecision or a non-nil error. Implementations must not
|
|
// retain the batch beyond the call.
|
|
Evaluate(context.Context, EvidenceBatch) (FilterDecision, error)
|
|
}
|
|
|
|
// FilterEnforcement identifies how a filter's decision is enforced when the
|
|
// current epoch requires evaluation. The parallel evaluator reads this value
|
|
// to decide whether a blocking violation must terminate the stream.
|
|
type FilterEnforcement string
|
|
|
|
const (
|
|
// FilterEnforcementBlocking indicates the filter may halt release of the
|
|
// current event downstream.
|
|
FilterEnforcementBlocking FilterEnforcement = "blocking"
|
|
|
|
// FilterEnforcementObserveOnly indicates the filter observes the event
|
|
// but must not halt release even on a violation.
|
|
FilterEnforcementObserveOnly FilterEnforcement = "observe_only"
|
|
)
|
|
|
|
// knownFilterEnforcements is the closed set of valid FilterEnforcement values.
|
|
var knownFilterEnforcements = map[FilterEnforcement]struct{}{
|
|
FilterEnforcementBlocking: {},
|
|
FilterEnforcementObserveOnly: {},
|
|
}
|
|
|
|
// Validate returns nil when the FilterEnforcement is a known lifecycle value.
|
|
func (e FilterEnforcement) Validate() error {
|
|
switch e {
|
|
case FilterEnforcementBlocking, FilterEnforcementObserveOnly:
|
|
return nil
|
|
}
|
|
return errors.New("streamgate: unknown filter enforcement: " + string(e))
|
|
}
|
|
|
|
// EpochFilter binds a FilterEvaluator to its execution contract for a single
|
|
// epoch. It encodes whether the filter applies, whether its trigger is ready,
|
|
// and how its decision is enforced. Construction validates every field and
|
|
// normalises readiness into a stable predicate that the parallel evaluator
|
|
// reads before dispatching to the guardrail.
|
|
type EpochFilter struct {
|
|
id StableToken
|
|
evaluator FilterEvaluator
|
|
subscribed bool
|
|
triggerReady bool
|
|
blocksRelease bool
|
|
enforcement FilterEnforcement
|
|
timeout time.Duration
|
|
}
|
|
|
|
// NewEpochFilter creates a validated EpochFilter. All fields are enforced:
|
|
//
|
|
// - evaluator must be non-nil.
|
|
// - evaluator.ID() must be a valid stable token and non-empty.
|
|
// - enforcement must be a known value.
|
|
// - timeout must be strictly positive.
|
|
//
|
|
// The returned value is immutable to callers.
|
|
func NewEpochFilter(
|
|
evaluator FilterEvaluator,
|
|
subscribedEventPresent, triggerReady, blocksRelease bool,
|
|
enforcement FilterEnforcement,
|
|
timeout time.Duration,
|
|
) (EpochFilter, error) {
|
|
if evaluator == nil {
|
|
return EpochFilter{}, errors.New("streamgate: epoch filter evaluator is required")
|
|
}
|
|
id, err := NewStableTokenRequired("filterID", evaluator.ID())
|
|
if err != nil {
|
|
return EpochFilter{}, errors.New("streamgate: epoch filter evaluator id: " + err.Error())
|
|
}
|
|
if err := enforcement.Validate(); err != nil {
|
|
return EpochFilter{}, errors.New("streamgate: epoch filter enforcement: " + err.Error())
|
|
}
|
|
if timeout <= 0 {
|
|
return EpochFilter{}, errors.New("streamgate: epoch filter evaluation timeout must be positive")
|
|
}
|
|
return EpochFilter{
|
|
id: id,
|
|
evaluator: evaluator,
|
|
subscribed: subscribedEventPresent,
|
|
triggerReady: triggerReady,
|
|
blocksRelease: blocksRelease,
|
|
enforcement: enforcement,
|
|
timeout: timeout,
|
|
}, nil
|
|
}
|
|
|
|
// ID returns the stable identifier snapshot at construction time.
|
|
// The underlying evaluator may be mutated without affecting this value.
|
|
func (f EpochFilter) ID() string { return f.id.String() }
|
|
|
|
// EnforcedAs returns the enforcement mode recorded at construction.
|
|
func (f EpochFilter) EnforcedAs() FilterEnforcement { return f.enforcement }
|
|
|
|
// SubscribedEventPresent returns whether the current epoch subscribed to this
|
|
// filter's event kind.
|
|
func (f EpochFilter) SubscribedEventPresent() bool { return f.subscribed }
|
|
|
|
// TriggerReady returns whether the filter's trigger condition is satisfied
|
|
// for the current epoch.
|
|
func (f EpochFilter) TriggerReady() bool { return f.triggerReady }
|
|
|
|
// BlocksRelease returns whether this filter blocks release of the current
|
|
// event on a violation.
|
|
func (f EpochFilter) BlocksRelease() bool { return f.blocksRelease }
|
|
|
|
// EvaluationTimeout returns the positive duration allocated to a single
|
|
// Evaluate call.
|
|
func (f EpochFilter) EvaluationTimeout() time.Duration { return f.timeout }
|
|
|
|
// EvaluatedForEpoch reports whether this filter is a ready evaluation target
|
|
// for the current epoch. The normalisation rules are:
|
|
//
|
|
// - unsubscribed epochs are not_applicable_for_epoch.
|
|
// - subscribed epochs where the trigger is not ready and the filter blocks
|
|
// release are deferred_by_requirement.
|
|
// - subscribed epochs where the trigger is not ready and the filter does
|
|
// not block release are not_applicable_for_epoch.
|
|
// - subscribed epochs where the trigger is ready are ready for evaluation.
|
|
//
|
|
// The parallel evaluator dispatches to the guardrail only when this returns
|
|
// true.
|
|
func (f EpochFilter) EvaluatedForEpoch() bool {
|
|
if !f.subscribed {
|
|
return false
|
|
}
|
|
if !f.triggerReady {
|
|
if f.blocksRelease {
|
|
// deferred_by_requirement — not evaluated this epoch.
|
|
return false
|
|
}
|
|
// nonblocking trigger unready — treated as not_applicable.
|
|
return false
|
|
}
|
|
return true
|
|
}
|
|
|
|
// NormalizeOutcome derives the canonical FilterOutcome for the current epoch
|
|
// from the binding state. The returned outcome captures exactly one of
|
|
// evaluated, evaluation_error, not_applicable_for_epoch, or
|
|
// deferred_by_requirement with no decision and no raw error payload.
|
|
func (f EpochFilter) NormalizeOutcome() FilterOutcome {
|
|
if !f.subscribed {
|
|
return NewFilterOutcomeNotApplicableForEpoch()
|
|
}
|
|
if !f.triggerReady {
|
|
if f.blocksRelease {
|
|
return NewFilterOutcomeDeferredByRequirement()
|
|
}
|
|
return NewFilterOutcomeNotApplicableForEpoch()
|
|
}
|
|
// Ready paths do not produce a pre-evaluation outcome; the evaluator
|
|
// result will be wrapped separately by the parallel coordinator.
|
|
return FilterOutcome{}
|
|
}
|
|
|
|
// Evaluator returns the bound FilterEvaluator for direct invocation. The
|
|
// parallel coordinator uses this seam to fan out to the guardrail.
|
|
func (f EpochFilter) Evaluator() FilterEvaluator { return f.evaluator }
|
|
|
|
// EvaluationFailureDisposition identifies how an evaluation error is
|
|
// propagated by the parallel coordinator. It is derived from the filter's
|
|
// enforcement mode and whether the outcome was an actual evaluation error.
|
|
type EvaluationFailureDisposition string
|
|
|
|
const (
|
|
// EvaluationFailureDispositionNone indicates no evaluation error
|
|
// occurred or the outcome was a successful evaluated decision.
|
|
EvaluationFailureDispositionNone EvaluationFailureDisposition = "none"
|
|
|
|
// EvaluationFailureDispositionBlockingFatal indicates a blocking
|
|
// filter produced a violation or fatal decision. The coordinator
|
|
// must terminate the stream.
|
|
EvaluationFailureDispositionBlockingFatal EvaluationFailureDisposition = "blocking_fatal"
|
|
|
|
// EvaluationFailureDispositionObserveError indicates an observe-only
|
|
// filter produced a violation or fatal decision. The coordinator
|
|
// logs the error code but does not terminate the stream.
|
|
EvaluationFailureDispositionObserveError EvaluationFailureDisposition = "observe_error"
|
|
)
|
|
|
|
// knownEvaluationFailureDispositions is the closed set of valid
|
|
// EvaluationFailureDisposition values.
|
|
var knownEvaluationFailureDispositions = map[EvaluationFailureDisposition]struct{}{
|
|
EvaluationFailureDispositionNone: {},
|
|
EvaluationFailureDispositionBlockingFatal: {},
|
|
EvaluationFailureDispositionObserveError: {},
|
|
}
|
|
|
|
// Validate returns nil when the EvaluationFailureDisposition is a known
|
|
// lifecycle value.
|
|
func (d EvaluationFailureDisposition) Validate() error {
|
|
switch d {
|
|
case EvaluationFailureDispositionNone,
|
|
EvaluationFailureDispositionBlockingFatal,
|
|
EvaluationFailureDispositionObserveError:
|
|
return nil
|
|
}
|
|
return errors.New("streamgate: unknown evaluation failure disposition: " + string(d))
|
|
}
|
|
|
|
// EpochFilterOutcome records a filter evaluation result for the immutable
|
|
// EvaluationSet. Raw errors are not preserved; the stable filter id, canonical
|
|
// FilterOutcome, enforcement mode, and failure disposition are retained. This
|
|
// guarantees that downstream Arbiter inputs cannot depend on completion order
|
|
// or leak raw error details.
|
|
type EpochFilterOutcome struct {
|
|
filterID StableToken
|
|
outcome FilterOutcome
|
|
enforcement FilterEnforcement
|
|
failureDisposition EvaluationFailureDisposition
|
|
}
|
|
|
|
// EvaluationSet is an immutable, stable-ordered collection of
|
|
// EpochFilterOutcome entries. Every accessor returns a defensive copy so
|
|
// that caller mutation cannot alter the set. The set is sorted by filter
|
|
// id in ascending stable-token order at construction time and remains
|
|
// stable regardless of completion order. The raw error payload is never
|
|
// preserved; only the stable error code survives in the outcome kind.
|
|
type EvaluationSet struct {
|
|
outcomes []EpochFilterOutcome
|
|
}
|
|
|
|
// NewEpochFilterOutcome constructs a validated EpochFilterOutcome.
|
|
//
|
|
// The caller provides the filter id, the raw FilterOutcome, and the
|
|
// enforcement mode. The constructor derives the failure disposition from
|
|
// these inputs:
|
|
//
|
|
// - evaluated + blocking violation/fatal → blocking_fatal
|
|
// - evaluated + observe_only violation/fatal → observe_error
|
|
// - evaluation_error + blocking → blocking_fatal
|
|
// - evaluation_error + observe_only → observe_error
|
|
// - not_applicable / deferred / evaluated with pass/observe/replacement → none
|
|
//
|
|
// The constructor rejects:
|
|
// - empty filter id
|
|
// - unknown filter id grammar
|
|
// - unknown enforcement
|
|
// - FilterOutcome that fails its own Validate()
|
|
//
|
|
// The returned value is immutable to callers.
|
|
func NewEpochFilterOutcome(
|
|
filterID string,
|
|
outcome FilterOutcome,
|
|
enforcement FilterEnforcement,
|
|
) (EpochFilterOutcome, error) {
|
|
if filterID == "" {
|
|
return EpochFilterOutcome{}, errors.New("streamgate: epoch filter outcome filter id is required")
|
|
}
|
|
fid, err := NewStableTokenRequired("filterID", filterID)
|
|
if err != nil {
|
|
return EpochFilterOutcome{}, errors.New("streamgate: epoch filter outcome filter id: " + err.Error())
|
|
}
|
|
if err := enforcement.Validate(); err != nil {
|
|
return EpochFilterOutcome{}, errors.New("streamgate: epoch filter outcome enforcement: " + err.Error())
|
|
}
|
|
if err := outcome.Validate(); err != nil {
|
|
return EpochFilterOutcome{}, errors.New("streamgate: epoch filter outcome: " + err.Error())
|
|
}
|
|
|
|
// Enforce that an evaluated outcome's decision filter id matches the
|
|
// wrapper's filter id. This prevents package-local forging from decoupling
|
|
// the outcome from its wrapper identity.
|
|
if outcome.Kind() == FilterOutcomeKindEvaluated && outcome.Decision() != nil {
|
|
if outcome.Decision().FilterID() != fid.String() {
|
|
return EpochFilterOutcome{}, errors.New("streamgate: evaluated outcome filter id mismatch")
|
|
}
|
|
}
|
|
|
|
disposition, err := deriveEvaluationFailureDisposition(outcome, enforcement)
|
|
if err != nil {
|
|
return EpochFilterOutcome{}, err
|
|
}
|
|
|
|
return EpochFilterOutcome{
|
|
filterID: fid,
|
|
outcome: outcome,
|
|
enforcement: enforcement,
|
|
failureDisposition: disposition,
|
|
}, nil
|
|
}
|
|
|
|
// FilterID returns the stable filter identifier.
|
|
func (o EpochFilterOutcome) FilterID() string { return o.filterID.value }
|
|
|
|
// Outcome returns a defensive copy of the canonical FilterOutcome.
|
|
func (o EpochFilterOutcome) Outcome() FilterOutcome {
|
|
cp := o.outcome
|
|
return cp
|
|
}
|
|
|
|
// Enforcement returns the enforcement mode at construction time.
|
|
func (o EpochFilterOutcome) Enforcement() FilterEnforcement { return o.enforcement }
|
|
|
|
// FailureDisposition returns the derived failure disposition.
|
|
func (o EpochFilterOutcome) FailureDisposition() EvaluationFailureDisposition {
|
|
return o.failureDisposition
|
|
}
|
|
|
|
// NewEvaluationSet constructs an immutable, stable-ordered EvaluationSet from
|
|
// the provided EpochFilterOutcome entries. The constructor:
|
|
//
|
|
// - Validates every entry via EpochFilterOutcome.Validate().
|
|
// - Rejects duplicate filter ids.
|
|
// - Rejects entries whose outcome kind is evaluated but whose decision is
|
|
// nil (the evaluator must produce a valid FilterDecision).
|
|
// - Rejects entries whose outcome kind is evaluation_error but whose error
|
|
// code is empty.
|
|
// - Sorts by filter id in ascending stable-token order.
|
|
// - Returns defensive copies so that caller mutation cannot alter the set.
|
|
//
|
|
// The returned value's accessor methods also return defensive copies.
|
|
func NewEvaluationSet(outcomes []EpochFilterOutcome) (EvaluationSet, error) {
|
|
if len(outcomes) == 0 {
|
|
return EvaluationSet{}, nil
|
|
}
|
|
|
|
seen := make(map[string]struct{}, len(outcomes))
|
|
cp := make([]EpochFilterOutcome, len(outcomes))
|
|
for i, o := range outcomes {
|
|
if err := o.Validate(); err != nil {
|
|
return EvaluationSet{}, errors.New("streamgate: evaluation set entry " + string(rune('0'+i)) + ": " + err.Error())
|
|
}
|
|
if _, dup := seen[o.filterID.value]; dup {
|
|
return EvaluationSet{}, errors.New("streamgate: evaluation set duplicate filter id: " + o.filterID.value)
|
|
}
|
|
seen[o.filterID.value] = struct{}{}
|
|
cp[i] = o
|
|
}
|
|
|
|
// Stable ascending order by filter id.
|
|
sort.SliceStable(cp, func(i, j int) bool {
|
|
return cp[i].filterID.value < cp[j].filterID.value
|
|
})
|
|
|
|
return EvaluationSet{outcomes: cp}, nil
|
|
}
|
|
|
|
// Validate returns nil when the EvaluationSet is in a consistent state.
|
|
func (s EvaluationSet) Validate() error {
|
|
seen := make(map[string]struct{}, len(s.outcomes))
|
|
for _, o := range s.outcomes {
|
|
if err := o.Validate(); err != nil {
|
|
return err
|
|
}
|
|
if _, dup := seen[o.filterID.value]; dup {
|
|
return errors.New("streamgate: evaluation set duplicate filter id: " + o.filterID.value)
|
|
}
|
|
seen[o.filterID.value] = struct{}{}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Len returns the number of outcomes in the set.
|
|
func (s EvaluationSet) Len() int { return len(s.outcomes) }
|
|
|
|
// At returns a defensive copy of the outcome at index i. Callers must not
|
|
// mutate the returned value.
|
|
func (s EvaluationSet) At(i int) (EpochFilterOutcome, error) {
|
|
if i < 0 || i >= len(s.outcomes) {
|
|
return EpochFilterOutcome{}, errors.New("streamgate: evaluation set index out of range")
|
|
}
|
|
return s.outcomes[i], nil
|
|
}
|
|
|
|
// ByID returns a defensive copy of the outcome whose filter id matches the
|
|
// given id, or an error if no such entry exists.
|
|
func (s EvaluationSet) ByID(id string) (EpochFilterOutcome, error) {
|
|
for _, o := range s.outcomes {
|
|
if o.filterID.value == id {
|
|
return o, nil
|
|
}
|
|
}
|
|
return EpochFilterOutcome{}, errors.New("streamgate: evaluation set filter id not found: " + id)
|
|
}
|
|
|
|
// All returns a defensive copy of all outcomes in stable id order.
|
|
func (s EvaluationSet) All() []EpochFilterOutcome {
|
|
if s.outcomes == nil {
|
|
return nil
|
|
}
|
|
cp := make([]EpochFilterOutcome, len(s.outcomes))
|
|
copy(cp, s.outcomes)
|
|
return cp
|
|
}
|
|
|
|
// Validate returns nil when the EpochFilterOutcome is in a consistent state.
|
|
func (o EpochFilterOutcome) Validate() error {
|
|
if o.filterID.value == "" {
|
|
return errors.New("streamgate: epoch filter outcome filter id is required")
|
|
}
|
|
if err := validateStableTokenRequired("filterID", o.filterID.value); err != nil {
|
|
return err
|
|
}
|
|
if err := o.outcome.Validate(); err != nil {
|
|
return err
|
|
}
|
|
if err := o.enforcement.Validate(); err != nil {
|
|
return err
|
|
}
|
|
if err := o.failureDisposition.Validate(); err != nil {
|
|
return err
|
|
}
|
|
switch o.outcome.Kind() {
|
|
case FilterOutcomeKindEvaluated:
|
|
if o.outcome.Decision() == nil {
|
|
return errors.New("streamgate: evaluated outcome must carry a decision")
|
|
}
|
|
if o.outcome.Decision().FilterID() != o.filterID.value {
|
|
return errors.New("streamgate: evaluated outcome filter id mismatch")
|
|
}
|
|
case FilterOutcomeKindEvaluationError:
|
|
if o.outcome.ErrorCode() == "" {
|
|
return errors.New("streamgate: evaluation_error outcome must carry an error code")
|
|
}
|
|
case FilterOutcomeKindNotApplicableForEpoch, FilterOutcomeKindDeferredByRequirement:
|
|
// These are always valid; nothing to check beyond Validate.
|
|
default:
|
|
return errors.New("streamgate: unknown outcome kind")
|
|
}
|
|
expected, err := deriveEvaluationFailureDisposition(o.outcome, o.enforcement)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
if o.failureDisposition != expected {
|
|
return errors.New("streamgate: epoch filter outcome failure disposition mismatch")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// deriveEvaluationFailureDisposition derives the canonical failure disposition
|
|
// from the given outcome and enforcement mode. The rules are the single
|
|
// source of truth shared between NewEpochFilterOutcome and
|
|
// EpochFilterOutcome.Validate():
|
|
//
|
|
// - evaluated + blocking violation/fatal → blocking_fatal
|
|
// - evaluated + observe_only violation/fatal → observe_error
|
|
// - evaluation_error + blocking → blocking_fatal
|
|
// - evaluation_error + observe_only → observe_error
|
|
// - not_applicable / deferred / evaluated with pass/observe/replacement → none
|
|
//
|
|
// Returns an error only for outcomes whose kind is not in the known set.
|
|
func deriveEvaluationFailureDisposition(
|
|
outcome FilterOutcome,
|
|
enforcement FilterEnforcement,
|
|
) (EvaluationFailureDisposition, error) {
|
|
switch outcome.Kind() {
|
|
case FilterOutcomeKindEvaluated:
|
|
decision := outcome.Decision()
|
|
if decision == nil {
|
|
return "", errors.New("streamgate: evaluated outcome must carry a decision")
|
|
}
|
|
switch decision.Kind() {
|
|
case FilterDecisionKindPass, FilterDecisionKindObserve, FilterDecisionKindReplacement:
|
|
return EvaluationFailureDispositionNone, nil
|
|
case FilterDecisionKindViolation, FilterDecisionKindFatal:
|
|
if enforcement == FilterEnforcementBlocking {
|
|
return EvaluationFailureDispositionBlockingFatal, nil
|
|
}
|
|
return EvaluationFailureDispositionObserveError, nil
|
|
default:
|
|
return EvaluationFailureDispositionNone, nil
|
|
}
|
|
case FilterOutcomeKindEvaluationError:
|
|
if enforcement == FilterEnforcementBlocking {
|
|
return EvaluationFailureDispositionBlockingFatal, nil
|
|
}
|
|
return EvaluationFailureDispositionObserveError, nil
|
|
case FilterOutcomeKindNotApplicableForEpoch, FilterOutcomeKindDeferredByRequirement:
|
|
return EvaluationFailureDispositionNone, nil
|
|
default:
|
|
return "", errors.New("streamgate: epoch filter outcome unknown kind")
|
|
}
|
|
}
|