- 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
952 lines
32 KiB
Go
952 lines
32 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"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// EventKind identifies the semantic role of a normalized event within the
|
|
// stream gate lifecycle.
|
|
type EventKind string
|
|
|
|
const (
|
|
// EventKindResponseStart is emitted once per channel when the provider
|
|
// announces a new response.
|
|
EventKindResponseStart EventKind = "response_start"
|
|
|
|
// EventKindTextDelta is emitted for each text chunk within a response.
|
|
EventKindTextDelta EventKind = "text_delta"
|
|
|
|
// EventKindReasoningDelta is emitted for each reasoning chunk within a
|
|
// response.
|
|
EventKindReasoningDelta EventKind = "reasoning_delta"
|
|
|
|
// EventKindToolCallFragment is emitted when a provider emits a tool call.
|
|
EventKindToolCallFragment EventKind = "tool_call_fragment"
|
|
|
|
// EventKindTerminal is emitted when the provider completes a response
|
|
// successfully.
|
|
EventKindTerminal EventKind = "terminal"
|
|
|
|
// EventKindProviderError is emitted when the provider reports a
|
|
// non-recoverable error.
|
|
EventKindProviderError EventKind = "provider_error"
|
|
)
|
|
|
|
// BaseEventDisposition describes the lifecycle stage of an event before any
|
|
// codec-specific terminal resolution.
|
|
type BaseEventDisposition string
|
|
|
|
const (
|
|
// BaseDispositionHold indicates the event is buffering and has not yet
|
|
// been committed to the downstream sink.
|
|
BaseDispositionHold BaseEventDisposition = "hold"
|
|
|
|
// BaseDispositionReleaseCandidate indicates the event has passed filter
|
|
// evaluation and is eligible for release.
|
|
BaseDispositionReleaseCandidate BaseEventDisposition = "release_candidate"
|
|
|
|
// BaseDispositionTerminalSuccessCandidate indicates the event is a
|
|
// successful terminal awaiting commit.
|
|
BaseDispositionTerminalSuccessCandidate BaseEventDisposition = "terminal_success_candidate"
|
|
|
|
// BaseDispositionTerminalErrorCandidate indicates the event is a
|
|
// provider error terminal awaiting commit.
|
|
BaseDispositionTerminalErrorCandidate BaseEventDisposition = "terminal_error_candidate"
|
|
)
|
|
|
|
var knownEventKinds = map[EventKind]struct{}{
|
|
EventKindResponseStart: {},
|
|
EventKindTextDelta: {},
|
|
EventKindReasoningDelta: {},
|
|
EventKindToolCallFragment: {},
|
|
EventKindTerminal: {},
|
|
EventKindProviderError: {},
|
|
}
|
|
|
|
var eventKindDisposition = map[EventKind]BaseEventDisposition{
|
|
EventKindResponseStart: BaseDispositionHold,
|
|
EventKindTextDelta: BaseDispositionReleaseCandidate,
|
|
EventKindReasoningDelta: BaseDispositionReleaseCandidate,
|
|
EventKindToolCallFragment: BaseDispositionReleaseCandidate,
|
|
EventKindTerminal: BaseDispositionTerminalSuccessCandidate,
|
|
EventKindProviderError: BaseDispositionTerminalErrorCandidate,
|
|
}
|
|
|
|
// BaseDispositionOf returns the base lifecycle disposition for a given event
|
|
// kind. This table is the single source of truth for codec-agnostic lifecycle
|
|
// state and must not be overridden by callers. It returns an error for unknown
|
|
// kinds.
|
|
func BaseDispositionOf(kind EventKind) (BaseEventDisposition, error) {
|
|
if _, ok := knownEventKinds[kind]; !ok {
|
|
return "", errors.New("streamgate: unknown event kind: " + string(kind))
|
|
}
|
|
return eventKindDisposition[kind], nil
|
|
}
|
|
|
|
// Validate returns nil when the event kind is a known lifecycle value.
|
|
func (k EventKind) Validate() error {
|
|
switch k {
|
|
case EventKindResponseStart, EventKindTextDelta, EventKindReasoningDelta,
|
|
EventKindToolCallFragment, EventKindTerminal, EventKindProviderError:
|
|
return nil
|
|
}
|
|
return errors.New("streamgate: unknown event kind: " + string(k))
|
|
}
|
|
|
|
// minHTTPStatus is the lowest valid HTTP status code.
|
|
const minHTTPStatus = 100
|
|
|
|
// maxHTTPStatus is the highest valid HTTP status code.
|
|
const maxHTTPStatus = 599
|
|
|
|
// forbiddenResponseStartHeaders is the case-insensitive set of hop-by-hop and
|
|
// transport metadata headers that must never be stored in the core event
|
|
// contract. These belong to the transport layer and must be stripped before
|
|
// the response-start enters the stream gate.
|
|
var forbiddenResponseStartHeaders = map[string]struct{}{
|
|
"connection": {},
|
|
"keep-alive": {},
|
|
"proxy-authenticate": {},
|
|
"proxy-authorization": {},
|
|
"proxy-connection": {},
|
|
"te": {},
|
|
"trailer": {},
|
|
"transfer-encoding": {},
|
|
"upgrade": {},
|
|
"content-length": {},
|
|
}
|
|
|
|
// validateResponseStartHeaders rejects hop-by-hop and transport metadata
|
|
// headers that must not be carried by the core event contract. Header names
|
|
// are matched case-insensitively. The header value is never included in the
|
|
// returned error.
|
|
func validateResponseStartHeaders(headers map[string]string) error {
|
|
for k := range headers {
|
|
if _, forbidden := forbiddenResponseStartHeaders[strings.ToLower(k)]; forbidden {
|
|
return errors.New("streamgate: response start header is forbidden: " + strings.ToLower(k))
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// ResponseStart describes the headers and metadata emitted once per channel
|
|
// when a provider begins a new response. All fields are private; callers use
|
|
// the accessor methods which return defensive copies.
|
|
type ResponseStart struct {
|
|
channel string
|
|
status int
|
|
headers map[string]string
|
|
timestamp time.Time
|
|
}
|
|
|
|
// NewResponseStart creates a ResponseStart with defensive copies of mutable
|
|
// inputs. The returned value is immutable to callers.
|
|
func NewResponseStart(channel string, status int, headers map[string]string, ts time.Time) (ResponseStart, error) {
|
|
if channel == "" {
|
|
return ResponseStart{}, errors.New("streamgate: response start channel is required")
|
|
}
|
|
if status < minHTTPStatus || status > maxHTTPStatus {
|
|
return ResponseStart{}, errors.New("streamgate: response start status must be between " + strconv.Itoa(minHTTPStatus) + " and " + strconv.Itoa(maxHTTPStatus))
|
|
}
|
|
if ts.IsZero() {
|
|
return ResponseStart{}, errors.New("streamgate: response start timestamp is required")
|
|
}
|
|
if err := validateResponseStartHeaders(headers); err != nil {
|
|
return ResponseStart{}, err
|
|
}
|
|
cp := make(map[string]string, len(headers))
|
|
for k, v := range headers {
|
|
cp[k] = v
|
|
}
|
|
return ResponseStart{
|
|
channel: channel,
|
|
status: status,
|
|
headers: cp,
|
|
timestamp: ts,
|
|
}, nil
|
|
}
|
|
|
|
// Validate returns nil when the ResponseStart is in a consistent state.
|
|
func (r ResponseStart) Validate() error {
|
|
if r.channel == "" {
|
|
return errors.New("streamgate: response start channel is required")
|
|
}
|
|
if r.status < minHTTPStatus || r.status > maxHTTPStatus {
|
|
return errors.New("streamgate: response start status must be between " + strconv.Itoa(minHTTPStatus) + " and " + strconv.Itoa(maxHTTPStatus))
|
|
}
|
|
if r.timestamp.IsZero() {
|
|
return errors.New("streamgate: response start timestamp is required")
|
|
}
|
|
if err := validateResponseStartHeaders(r.headers); err != nil {
|
|
return err
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Channel returns the response start channel.
|
|
func (r ResponseStart) Channel() string { return r.channel }
|
|
|
|
// Status returns the HTTP status code associated with the response start.
|
|
func (r ResponseStart) Status() int { return r.status }
|
|
|
|
// Headers returns a defensive copy of the response start headers.
|
|
func (r ResponseStart) Headers() map[string]string {
|
|
if r.headers == nil {
|
|
return nil
|
|
}
|
|
out := make(map[string]string, len(r.headers))
|
|
for k, v := range r.headers {
|
|
out[k] = v
|
|
}
|
|
return out
|
|
}
|
|
|
|
// Timestamp returns the response start timestamp.
|
|
func (r ResponseStart) Timestamp() time.Time { return r.timestamp }
|
|
|
|
// NormalizedEvent is a transport-agnostic event within the stream gate
|
|
// lifecycle. All payload fields are kind-specific and private; callers use
|
|
// the typed accessor methods. The disposition is set internally by the
|
|
// constructor based on the event kind and cannot be overridden by callers.
|
|
type NormalizedEvent struct {
|
|
kind EventKind
|
|
channel string
|
|
timestamp time.Time
|
|
disposition BaseEventDisposition
|
|
|
|
// response_start
|
|
status int
|
|
headers map[string]string
|
|
|
|
// text_delta
|
|
textDelta string
|
|
|
|
// reasoning_delta
|
|
reasoningDelta string
|
|
|
|
// tool_call_fragment
|
|
toolCallID string
|
|
toolCallName string
|
|
toolCallArgs string
|
|
|
|
// terminal / provider_error
|
|
terminalSuccess bool
|
|
externalDesc *ExternalDescriptor
|
|
failureCauses FailureCauseChain
|
|
}
|
|
|
|
// NewResponseStartEvent creates a NormalizedEvent of kind response_start.
|
|
func NewResponseStartEvent(channel string, status int, headers map[string]string, ts time.Time) (NormalizedEvent, error) {
|
|
if channel == "" {
|
|
return NormalizedEvent{}, errors.New("streamgate: response start event channel is required")
|
|
}
|
|
if status < minHTTPStatus || status > maxHTTPStatus {
|
|
return NormalizedEvent{}, errors.New("streamgate: response start event status must be between " + strconv.Itoa(minHTTPStatus) + " and " + strconv.Itoa(maxHTTPStatus))
|
|
}
|
|
if ts.IsZero() {
|
|
return NormalizedEvent{}, errors.New("streamgate: response start event timestamp is required")
|
|
}
|
|
if err := validateResponseStartHeaders(headers); err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
cp := make(map[string]string, len(headers))
|
|
for k, v := range headers {
|
|
cp[k] = v
|
|
}
|
|
disp, err := BaseDispositionOf(EventKindResponseStart)
|
|
if err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
return NormalizedEvent{
|
|
kind: EventKindResponseStart,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
disposition: disp,
|
|
status: status,
|
|
headers: cp,
|
|
}, nil
|
|
}
|
|
|
|
// NewTextDeltaEvent creates a NormalizedEvent of kind text_delta.
|
|
func NewTextDeltaEvent(channel, text string, ts time.Time) (NormalizedEvent, error) {
|
|
if channel == "" {
|
|
return NormalizedEvent{}, errors.New("streamgate: text delta event channel is required")
|
|
}
|
|
if ts.IsZero() {
|
|
return NormalizedEvent{}, errors.New("streamgate: text delta event timestamp is required")
|
|
}
|
|
disp, err := BaseDispositionOf(EventKindTextDelta)
|
|
if err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
ev := NormalizedEvent{
|
|
kind: EventKindTextDelta,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
disposition: disp,
|
|
textDelta: text,
|
|
}
|
|
if err := ev.Validate(); err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
return ev, nil
|
|
}
|
|
|
|
// NewReasoningDeltaEvent creates a NormalizedEvent of kind reasoning_delta.
|
|
func NewReasoningDeltaEvent(channel, reasoning string, ts time.Time) (NormalizedEvent, error) {
|
|
if channel == "" {
|
|
return NormalizedEvent{}, errors.New("streamgate: reasoning delta event channel is required")
|
|
}
|
|
if ts.IsZero() {
|
|
return NormalizedEvent{}, errors.New("streamgate: reasoning delta event timestamp is required")
|
|
}
|
|
disp, err := BaseDispositionOf(EventKindReasoningDelta)
|
|
if err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
ev := NormalizedEvent{
|
|
kind: EventKindReasoningDelta,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
disposition: disp,
|
|
reasoningDelta: reasoning,
|
|
}
|
|
if err := ev.Validate(); err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
return ev, nil
|
|
}
|
|
|
|
// NewToolCallFragmentEvent creates a NormalizedEvent of kind tool_call_fragment.
|
|
// The request-local raw arguments are preserved inside the typed event.
|
|
func NewToolCallFragmentEvent(channel, toolCallID, toolCallName, toolCallArgs string, ts time.Time) (NormalizedEvent, error) {
|
|
if channel == "" {
|
|
return NormalizedEvent{}, errors.New("streamgate: tool call fragment event channel is required")
|
|
}
|
|
if toolCallID == "" {
|
|
return NormalizedEvent{}, errors.New("streamgate: tool call fragment event id is required")
|
|
}
|
|
if toolCallName == "" {
|
|
return NormalizedEvent{}, errors.New("streamgate: tool call fragment event name is required")
|
|
}
|
|
if ts.IsZero() {
|
|
return NormalizedEvent{}, errors.New("streamgate: tool call fragment event timestamp is required")
|
|
}
|
|
disp, err := BaseDispositionOf(EventKindToolCallFragment)
|
|
if err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
ev := NormalizedEvent{
|
|
kind: EventKindToolCallFragment,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
disposition: disp,
|
|
toolCallID: toolCallID,
|
|
toolCallName: toolCallName,
|
|
toolCallArgs: toolCallArgs,
|
|
}
|
|
if err := ev.Validate(); err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
return ev, nil
|
|
}
|
|
|
|
// NewTerminalEvent creates a NormalizedEvent of kind terminal (success).
|
|
// Success terminals must not carry an external descriptor or failure causes.
|
|
func NewTerminalEvent(channel string, ts time.Time) (NormalizedEvent, error) {
|
|
if channel == "" {
|
|
return NormalizedEvent{}, errors.New("streamgate: terminal event channel is required")
|
|
}
|
|
if ts.IsZero() {
|
|
return NormalizedEvent{}, errors.New("streamgate: terminal event timestamp is required")
|
|
}
|
|
disp, err := BaseDispositionOf(EventKindTerminal)
|
|
if err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
return NormalizedEvent{
|
|
kind: EventKindTerminal,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
disposition: disp,
|
|
terminalSuccess: true,
|
|
}, nil
|
|
}
|
|
|
|
// NewProviderErrorEvent creates a NormalizedEvent of kind provider_error.
|
|
// It requires a validated external descriptor and an optional bounded failure
|
|
// cause chain.
|
|
func NewProviderErrorEvent(
|
|
channel string,
|
|
desc ExternalDescriptor,
|
|
causes FailureCauseChain,
|
|
ts time.Time,
|
|
) (NormalizedEvent, error) {
|
|
if channel == "" {
|
|
return NormalizedEvent{}, errors.New("streamgate: provider error event channel is required")
|
|
}
|
|
if ts.IsZero() {
|
|
return NormalizedEvent{}, errors.New("streamgate: provider error event timestamp is required")
|
|
}
|
|
if err := desc.Validate(); err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
if causes.Len() > MaxFailureCauses {
|
|
return NormalizedEvent{}, errors.New("streamgate: provider error event failure cause chain exceeds maximum")
|
|
}
|
|
for i := 0; i < causes.Len(); i++ {
|
|
c, err := causes.At(i)
|
|
if err != nil {
|
|
return NormalizedEvent{}, errors.New("streamgate: provider error event failure cause chain: " + err.Error())
|
|
}
|
|
if err := c.Validate(); err != nil {
|
|
return NormalizedEvent{}, errors.New("streamgate: provider error event failure cause chain entry " + strconv.Itoa(i) + ": " + err.Error())
|
|
}
|
|
}
|
|
disp, err := BaseDispositionOf(EventKindProviderError)
|
|
if err != nil {
|
|
return NormalizedEvent{}, err
|
|
}
|
|
descCopy := desc
|
|
return NormalizedEvent{
|
|
kind: EventKindProviderError,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
disposition: disp,
|
|
externalDesc: &descCopy,
|
|
failureCauses: causes.Copy(),
|
|
}, nil
|
|
}
|
|
|
|
// Validate returns nil when the NormalizedEvent is in a consistent state.
|
|
// Each kind's required payload must be present and every forbidden payload
|
|
// field must be empty or nil. Nested StableToken, FailureCause, and
|
|
// ExternalDescriptor values are re-validated against the same grammar.
|
|
func (e NormalizedEvent) Validate() error {
|
|
if e.kind == "" {
|
|
return errors.New("streamgate: normalized event kind is required")
|
|
}
|
|
if err := e.kind.Validate(); err != nil {
|
|
return err
|
|
}
|
|
if e.channel == "" {
|
|
return errors.New("streamgate: normalized event channel is required")
|
|
}
|
|
if e.timestamp.IsZero() {
|
|
return errors.New("streamgate: normalized event timestamp is required")
|
|
}
|
|
expectedDisp, _ := BaseDispositionOf(e.kind)
|
|
if e.disposition != expectedDisp {
|
|
return errors.New("streamgate: normalized event disposition does not match kind")
|
|
}
|
|
switch e.kind {
|
|
case EventKindResponseStart:
|
|
if e.status < minHTTPStatus || e.status > maxHTTPStatus {
|
|
return errors.New("streamgate: response start event status must be between " + strconv.Itoa(minHTTPStatus) + " and " + strconv.Itoa(maxHTTPStatus))
|
|
}
|
|
if err := validateResponseStartHeaders(e.headers); err != nil {
|
|
return err
|
|
}
|
|
if e.textDelta != "" {
|
|
return errors.New("streamgate: response start event must not have text delta")
|
|
}
|
|
if e.reasoningDelta != "" {
|
|
return errors.New("streamgate: response start event must not have reasoning delta")
|
|
}
|
|
if e.toolCallID != "" || e.toolCallName != "" || e.toolCallArgs != "" {
|
|
return errors.New("streamgate: response start event must not have tool call fragment")
|
|
}
|
|
if e.terminalSuccess {
|
|
return errors.New("streamgate: response start event must not be terminal")
|
|
}
|
|
if e.externalDesc != nil {
|
|
return errors.New("streamgate: response start event must not have external descriptor")
|
|
}
|
|
if e.failureCauses.Len() > 0 {
|
|
return errors.New("streamgate: response start event must not have failure causes")
|
|
}
|
|
case EventKindTextDelta:
|
|
if e.textDelta == "" {
|
|
return errors.New("streamgate: text delta event requires content")
|
|
}
|
|
if e.status != 0 {
|
|
return errors.New("streamgate: text delta event must not have response start status")
|
|
}
|
|
if e.headers != nil {
|
|
return errors.New("streamgate: text delta event must not have response start headers")
|
|
}
|
|
if e.reasoningDelta != "" {
|
|
return errors.New("streamgate: text delta event must not have reasoning delta")
|
|
}
|
|
if e.toolCallID != "" || e.toolCallName != "" || e.toolCallArgs != "" {
|
|
return errors.New("streamgate: text delta event must not have tool call fragment")
|
|
}
|
|
if e.terminalSuccess {
|
|
return errors.New("streamgate: text delta event must not be terminal")
|
|
}
|
|
if e.externalDesc != nil {
|
|
return errors.New("streamgate: text delta event must not have external descriptor")
|
|
}
|
|
if e.failureCauses.Len() > 0 {
|
|
return errors.New("streamgate: text delta event must not have failure causes")
|
|
}
|
|
case EventKindReasoningDelta:
|
|
if e.reasoningDelta == "" {
|
|
return errors.New("streamgate: reasoning delta event requires content")
|
|
}
|
|
if e.status != 0 {
|
|
return errors.New("streamgate: reasoning delta event must not have response start status")
|
|
}
|
|
if e.headers != nil {
|
|
return errors.New("streamgate: reasoning delta event must not have response start headers")
|
|
}
|
|
if e.textDelta != "" {
|
|
return errors.New("streamgate: reasoning delta event must not have text delta")
|
|
}
|
|
if e.toolCallID != "" || e.toolCallName != "" || e.toolCallArgs != "" {
|
|
return errors.New("streamgate: reasoning delta event must not have tool call fragment")
|
|
}
|
|
if e.terminalSuccess {
|
|
return errors.New("streamgate: reasoning delta event must not be terminal")
|
|
}
|
|
if e.externalDesc != nil {
|
|
return errors.New("streamgate: reasoning delta event must not have external descriptor")
|
|
}
|
|
if e.failureCauses.Len() > 0 {
|
|
return errors.New("streamgate: reasoning delta event must not have failure causes")
|
|
}
|
|
case EventKindToolCallFragment:
|
|
if e.toolCallID == "" {
|
|
return errors.New("streamgate: tool call fragment event requires id")
|
|
}
|
|
if e.toolCallName == "" {
|
|
return errors.New("streamgate: tool call fragment event requires name")
|
|
}
|
|
if e.toolCallArgs == "" {
|
|
return errors.New("streamgate: tool call fragment event requires arguments")
|
|
}
|
|
if e.status != 0 {
|
|
return errors.New("streamgate: tool call fragment event must not have response start status")
|
|
}
|
|
if e.headers != nil {
|
|
return errors.New("streamgate: tool call fragment event must not have response start headers")
|
|
}
|
|
if e.textDelta != "" {
|
|
return errors.New("streamgate: tool call fragment event must not have text delta")
|
|
}
|
|
if e.reasoningDelta != "" {
|
|
return errors.New("streamgate: tool call fragment event must not have reasoning delta")
|
|
}
|
|
if e.terminalSuccess {
|
|
return errors.New("streamgate: tool call fragment event must not be terminal")
|
|
}
|
|
if e.externalDesc != nil {
|
|
return errors.New("streamgate: tool call fragment event must not have external descriptor")
|
|
}
|
|
if e.failureCauses.Len() > 0 {
|
|
return errors.New("streamgate: tool call fragment event must not have failure causes")
|
|
}
|
|
case EventKindTerminal:
|
|
if !e.terminalSuccess {
|
|
return errors.New("streamgate: terminal event must be success")
|
|
}
|
|
if e.status != 0 {
|
|
return errors.New("streamgate: terminal event must not have response start status")
|
|
}
|
|
if e.headers != nil {
|
|
return errors.New("streamgate: terminal event must not have response start headers")
|
|
}
|
|
if e.textDelta != "" {
|
|
return errors.New("streamgate: terminal event must not have text delta")
|
|
}
|
|
if e.reasoningDelta != "" {
|
|
return errors.New("streamgate: terminal event must not have reasoning delta")
|
|
}
|
|
if e.toolCallID != "" || e.toolCallName != "" || e.toolCallArgs != "" {
|
|
return errors.New("streamgate: terminal event must not have tool call fragment")
|
|
}
|
|
if e.externalDesc != nil {
|
|
return errors.New("streamgate: terminal event must not have external descriptor")
|
|
}
|
|
if e.failureCauses.Len() > 0 {
|
|
return errors.New("streamgate: terminal event must not have failure causes")
|
|
}
|
|
case EventKindProviderError:
|
|
if e.externalDesc == nil {
|
|
return errors.New("streamgate: provider error event requires external descriptor")
|
|
}
|
|
if e.terminalSuccess {
|
|
return errors.New("streamgate: provider error event must not be success")
|
|
}
|
|
if e.status != 0 {
|
|
return errors.New("streamgate: provider error event must not have response start status")
|
|
}
|
|
if e.headers != nil {
|
|
return errors.New("streamgate: provider error event must not have response start headers")
|
|
}
|
|
if e.textDelta != "" {
|
|
return errors.New("streamgate: provider error event must not have text delta")
|
|
}
|
|
if e.reasoningDelta != "" {
|
|
return errors.New("streamgate: provider error event must not have reasoning delta")
|
|
}
|
|
if e.toolCallID != "" || e.toolCallName != "" || e.toolCallArgs != "" {
|
|
return errors.New("streamgate: provider error event must not have tool call fragment")
|
|
}
|
|
if err := e.externalDesc.Validate(); err != nil {
|
|
return errors.New("streamgate: provider error event external descriptor: " + err.Error())
|
|
}
|
|
if e.failureCauses.Len() > MaxFailureCauses {
|
|
return errors.New("streamgate: provider error event failure cause chain exceeds maximum")
|
|
}
|
|
for i := 0; i < e.failureCauses.Len(); i++ {
|
|
c, err := e.failureCauses.At(i)
|
|
if err != nil {
|
|
return errors.New("streamgate: provider error event failure cause chain: " + err.Error())
|
|
}
|
|
if err := c.Validate(); err != nil {
|
|
return errors.New("streamgate: provider error event failure cause chain entry " + strconv.Itoa(i) + ": " + err.Error())
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Kind returns the event kind.
|
|
func (e NormalizedEvent) Kind() EventKind { return e.kind }
|
|
|
|
// Channel returns the event channel.
|
|
func (e NormalizedEvent) Channel() string { return e.channel }
|
|
|
|
// Timestamp returns the event timestamp.
|
|
func (e NormalizedEvent) Timestamp() time.Time { return e.timestamp }
|
|
|
|
// Disposition returns the base event disposition.
|
|
func (e NormalizedEvent) Disposition() BaseEventDisposition { return e.disposition }
|
|
|
|
// AsResponseStart returns the response start data, or an error if the event
|
|
// is not of kind response_start.
|
|
func (e NormalizedEvent) AsResponseStart() (ResponseStart, error) {
|
|
if e.kind != EventKindResponseStart {
|
|
return ResponseStart{}, errors.New("streamgate: event is not response_start")
|
|
}
|
|
return ResponseStart{
|
|
channel: e.channel,
|
|
status: e.status,
|
|
headers: e.copyHeaders(),
|
|
timestamp: e.timestamp,
|
|
}, nil
|
|
}
|
|
|
|
// AsTextDelta returns the text delta content, or an error if the event is
|
|
// not of kind text_delta.
|
|
func (e NormalizedEvent) AsTextDelta() (string, error) {
|
|
if e.kind != EventKindTextDelta {
|
|
return "", errors.New("streamgate: event is not text_delta")
|
|
}
|
|
return e.textDelta, nil
|
|
}
|
|
|
|
// AsReasoningDelta returns the reasoning delta content, or an error if the
|
|
// event is not of kind reasoning_delta.
|
|
func (e NormalizedEvent) AsReasoningDelta() (string, error) {
|
|
if e.kind != EventKindReasoningDelta {
|
|
return "", errors.New("streamgate: event is not reasoning_delta")
|
|
}
|
|
return e.reasoningDelta, nil
|
|
}
|
|
|
|
// AsToolCallFragment returns the tool call fragment data, or an error if the
|
|
// event is not of kind tool_call_fragment.
|
|
func (e NormalizedEvent) AsToolCallFragment() (ToolCall, error) {
|
|
if e.kind != EventKindToolCallFragment {
|
|
return ToolCall{}, errors.New("streamgate: event is not tool_call_fragment")
|
|
}
|
|
return ToolCall{
|
|
ID: e.toolCallID,
|
|
Name: e.toolCallName,
|
|
Arguments: e.toolCallArgs,
|
|
}, nil
|
|
}
|
|
|
|
// AsTerminal returns a success TerminalResult, or an error if the event is
|
|
// not of kind terminal.
|
|
func (e NormalizedEvent) AsTerminal() (TerminalResult, error) {
|
|
if e.kind != EventKindTerminal {
|
|
return TerminalResult{}, errors.New("streamgate: event is not terminal")
|
|
}
|
|
tr, err := NewSuccessTerminalResult(e.channel, e.timestamp)
|
|
if err != nil {
|
|
return TerminalResult{}, err
|
|
}
|
|
return tr, nil
|
|
}
|
|
|
|
// AsProviderError returns an error TerminalResult, or an error if the event
|
|
// is not of kind provider_error.
|
|
func (e NormalizedEvent) AsProviderError() (TerminalResult, error) {
|
|
if e.kind != EventKindProviderError {
|
|
return TerminalResult{}, errors.New("streamgate: event is not provider_error")
|
|
}
|
|
if e.externalDesc == nil {
|
|
return TerminalResult{}, errors.New("streamgate: provider error event requires external descriptor")
|
|
}
|
|
descCopy := *e.externalDesc
|
|
tr, err := NewErrorTerminalResult(e.channel, descCopy, e.failureCauses.Copy(), e.timestamp)
|
|
if err != nil {
|
|
return TerminalResult{}, err
|
|
}
|
|
return tr, nil
|
|
}
|
|
|
|
// copyHeaders returns a defensive copy of the event headers.
|
|
func (e NormalizedEvent) copyHeaders() map[string]string {
|
|
if e.headers == nil {
|
|
return nil
|
|
}
|
|
out := make(map[string]string, len(e.headers))
|
|
for k, v := range e.headers {
|
|
out[k] = v
|
|
}
|
|
return out
|
|
}
|
|
|
|
// ReleaseEvent is a transport-agnostic release payload emitted by the stream
|
|
// gate host to the downstream sink. ReleaseEvent carries only safe,
|
|
// downstream-transmittable content: text_delta, reasoning_delta, and
|
|
// tool_call_fragment. Response-start and terminal payloads are handled
|
|
// exclusively through ReleaseSink.CommitResponseStart and
|
|
// ReleaseSink.CommitTerminal and must never be stored in or exposed via
|
|
// ReleaseEvent.
|
|
type ReleaseEvent struct {
|
|
kind EventKind
|
|
channel string
|
|
timestamp time.Time
|
|
|
|
// text_delta
|
|
textDelta string
|
|
|
|
// reasoning_delta
|
|
reasoningDelta string
|
|
|
|
// tool_call_fragment
|
|
toolCallID string
|
|
toolCallName string
|
|
toolCallArgs string
|
|
}
|
|
|
|
// NewReleaseTextDeltaEvent creates a ReleaseEvent of kind text_delta.
|
|
func NewReleaseTextDeltaEvent(channel, text string, ts time.Time) (ReleaseEvent, error) {
|
|
if channel == "" {
|
|
return ReleaseEvent{}, errors.New("streamgate: release event channel is required")
|
|
}
|
|
if ts.IsZero() {
|
|
return ReleaseEvent{}, errors.New("streamgate: release event timestamp is required")
|
|
}
|
|
ev := ReleaseEvent{
|
|
kind: EventKindTextDelta,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
textDelta: text,
|
|
}
|
|
if err := ev.Validate(); err != nil {
|
|
return ReleaseEvent{}, err
|
|
}
|
|
return ev, nil
|
|
}
|
|
|
|
// NewReleaseReasoningDeltaEvent creates a ReleaseEvent of kind reasoning_delta.
|
|
func NewReleaseReasoningDeltaEvent(channel, reasoning string, ts time.Time) (ReleaseEvent, error) {
|
|
if channel == "" {
|
|
return ReleaseEvent{}, errors.New("streamgate: release event channel is required")
|
|
}
|
|
if ts.IsZero() {
|
|
return ReleaseEvent{}, errors.New("streamgate: release event timestamp is required")
|
|
}
|
|
ev := ReleaseEvent{
|
|
kind: EventKindReasoningDelta,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
reasoningDelta: reasoning,
|
|
}
|
|
if err := ev.Validate(); err != nil {
|
|
return ReleaseEvent{}, err
|
|
}
|
|
return ev, nil
|
|
}
|
|
|
|
// NewReleaseToolCallFragmentEvent creates a ReleaseEvent of kind
|
|
// tool_call_fragment. The tool call arguments are preserved through the
|
|
// release path.
|
|
func NewReleaseToolCallFragmentEvent(channel, toolCallID, toolCallName, toolCallArgs string, ts time.Time) (ReleaseEvent, error) {
|
|
if channel == "" {
|
|
return ReleaseEvent{}, errors.New("streamgate: release event channel is required")
|
|
}
|
|
if toolCallID == "" {
|
|
return ReleaseEvent{}, errors.New("streamgate: release event tool call id is required")
|
|
}
|
|
if toolCallName == "" {
|
|
return ReleaseEvent{}, errors.New("streamgate: release event tool call name is required")
|
|
}
|
|
if ts.IsZero() {
|
|
return ReleaseEvent{}, errors.New("streamgate: release event timestamp is required")
|
|
}
|
|
ev := ReleaseEvent{
|
|
kind: EventKindToolCallFragment,
|
|
channel: channel,
|
|
timestamp: ts,
|
|
toolCallID: toolCallID,
|
|
toolCallName: toolCallName,
|
|
toolCallArgs: toolCallArgs,
|
|
}
|
|
if err := ev.Validate(); err != nil {
|
|
return ReleaseEvent{}, err
|
|
}
|
|
return ev, nil
|
|
}
|
|
|
|
// Validate returns nil when the ReleaseEvent is in a consistent state.
|
|
// ReleaseEvent may only carry text_delta, reasoning_delta, or
|
|
// tool_call_fragment. Response-start and terminal payloads are forbidden.
|
|
func (r ReleaseEvent) Validate() error {
|
|
if r.kind == "" {
|
|
return errors.New("streamgate: release event kind is required")
|
|
}
|
|
if err := r.kind.Validate(); err != nil {
|
|
return err
|
|
}
|
|
if r.channel == "" {
|
|
return errors.New("streamgate: release event channel is required")
|
|
}
|
|
if r.timestamp.IsZero() {
|
|
return errors.New("streamgate: release event timestamp is required")
|
|
}
|
|
switch r.kind {
|
|
case EventKindTextDelta:
|
|
if r.textDelta == "" {
|
|
return errors.New("streamgate: release event text_delta requires content")
|
|
}
|
|
if r.reasoningDelta != "" {
|
|
return errors.New("streamgate: release event text_delta must not have reasoning content")
|
|
}
|
|
if r.toolCallID != "" || r.toolCallName != "" || r.toolCallArgs != "" {
|
|
return errors.New("streamgate: release event text_delta must not have tool call data")
|
|
}
|
|
case EventKindReasoningDelta:
|
|
if r.reasoningDelta == "" {
|
|
return errors.New("streamgate: release event reasoning_delta requires content")
|
|
}
|
|
if r.textDelta != "" {
|
|
return errors.New("streamgate: release event reasoning_delta must not have text content")
|
|
}
|
|
if r.toolCallID != "" || r.toolCallName != "" || r.toolCallArgs != "" {
|
|
return errors.New("streamgate: release event reasoning_delta must not have tool call data")
|
|
}
|
|
case EventKindToolCallFragment:
|
|
if r.toolCallID == "" {
|
|
return errors.New("streamgate: release event tool_call_fragment requires id")
|
|
}
|
|
if r.toolCallName == "" {
|
|
return errors.New("streamgate: release event tool_call_fragment requires name")
|
|
}
|
|
if r.toolCallArgs == "" {
|
|
return errors.New("streamgate: release event tool_call_fragment requires arguments")
|
|
}
|
|
if r.textDelta != "" {
|
|
return errors.New("streamgate: release event tool_call_fragment must not have text content")
|
|
}
|
|
if r.reasoningDelta != "" {
|
|
return errors.New("streamgate: release event tool_call_fragment must not have reasoning content")
|
|
}
|
|
default:
|
|
return errors.New("streamgate: release event kind is not releasable")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// Kind returns the release event kind.
|
|
func (r ReleaseEvent) Kind() EventKind { return r.kind }
|
|
|
|
// Channel returns the release event channel.
|
|
func (r ReleaseEvent) Channel() string { return r.channel }
|
|
|
|
// Timestamp returns the release event timestamp.
|
|
func (r ReleaseEvent) Timestamp() time.Time { return r.timestamp }
|
|
|
|
// AsTextDelta returns the release text delta content.
|
|
func (r ReleaseEvent) AsTextDelta() (string, error) {
|
|
if r.kind != EventKindTextDelta {
|
|
return "", errors.New("streamgate: release event is not text_delta")
|
|
}
|
|
return r.textDelta, nil
|
|
}
|
|
|
|
// AsReasoningDelta returns the release reasoning delta content.
|
|
func (r ReleaseEvent) AsReasoningDelta() (string, error) {
|
|
if r.kind != EventKindReasoningDelta {
|
|
return "", errors.New("streamgate: release event is not reasoning_delta")
|
|
}
|
|
return r.reasoningDelta, nil
|
|
}
|
|
|
|
// AsToolCallFragment returns the release tool call fragment data.
|
|
func (r ReleaseEvent) AsToolCallFragment() (ToolCall, error) {
|
|
if r.kind != EventKindToolCallFragment {
|
|
return ToolCall{}, errors.New("streamgate: release event is not tool_call_fragment")
|
|
}
|
|
return ToolCall{
|
|
ID: r.toolCallID,
|
|
Name: r.toolCallName,
|
|
Arguments: r.toolCallArgs,
|
|
}, nil
|
|
}
|
|
|
|
// ToolCall describes a single tool invocation emitted by a provider.
|
|
type ToolCall struct {
|
|
ID string
|
|
Name string
|
|
Arguments string
|
|
}
|
|
|
|
// NewToolCall creates a ToolCall with validation.
|
|
func NewToolCall(id, name string) (ToolCall, error) {
|
|
if id == "" {
|
|
return ToolCall{}, errors.New("streamgate: tool call id is required")
|
|
}
|
|
if name == "" {
|
|
return ToolCall{}, errors.New("streamgate: tool call name is required")
|
|
}
|
|
return ToolCall{
|
|
ID: id,
|
|
Name: name,
|
|
}, nil
|
|
}
|
|
|
|
// Validate returns nil when the ToolCall is in a consistent state.
|
|
func (tc ToolCall) Validate() error {
|
|
if tc.ID == "" {
|
|
return errors.New("streamgate: tool call id is required")
|
|
}
|
|
if tc.Name == "" {
|
|
return errors.New("streamgate: tool call name is required")
|
|
}
|
|
return nil
|
|
}
|
|
|
|
// copyExternalDescriptor returns a defensive copy of the external descriptor.
|
|
func copyExternalDescriptor(d *ExternalDescriptor) *ExternalDescriptor {
|
|
if d == nil {
|
|
return nil
|
|
}
|
|
cp := *d
|
|
return &cp
|
|
}
|