iop/apps/edge/internal/openai/hot_path_light.go

1262 lines
48 KiB
Go

package openai
import (
"context"
"encoding/json"
"fmt"
"net/http"
"strconv"
"strings"
"sync"
"time"
edgeservice "iop/apps/edge/internal/service"
"iop/packages/go/config"
)
const defaultHotPathLightCapacity = 1024
const hotPathOutputCapMetadata = "iop_hot_path_output_token_cap"
func hotPathOutputTokenCap(metadata map[string]string) int {
if metadata == nil {
return 0
}
cap, err := strconv.Atoi(strings.TrimSpace(metadata[hotPathOutputCapMetadata]))
if err != nil || cap < 1 {
return 0
}
return cap
}
// applyHotPathOutputTokenCap replaces any caller metadata value with the
// validated endpoint field. A missing field removes the internal key so
// metadata cannot manufacture a trusted output budget.
func applyHotPathOutputTokenCap(metadata map[string]string, candidates ...*int) {
if metadata == nil {
return
}
delete(metadata, hotPathOutputCapMetadata)
for _, candidate := range candidates {
if candidate != nil && *candidate > 0 {
metadata[hotPathOutputCapMetadata] = strconv.Itoa(*candidate)
return
}
}
}
type hotPathLightPhase string
const (
hotPathPhaseAwaitArtifacts hotPathLightPhase = "await_artifacts"
hotPathPhaseLocalActive hotPathLightPhase = "local_active"
hotPathPhaseReviewActive hotPathLightPhase = "review_active"
hotPathPhaseReviewRepair hotPathLightPhase = "review_repair_active"
hotPathPhaseCleanupPending hotPathLightPhase = "cleanup_pending"
)
type hotPathPendingKind string
const (
hotPathPendingLocalTools hotPathPendingKind = "local_tools"
hotPathPendingLocalHandoff hotPathPendingKind = "local_handoff"
hotPathPendingReviewInspection hotPathPendingKind = "review_inspection"
hotPathPendingReviewRepair hotPathPendingKind = "review_repair"
hotPathPendingCleanup hotPathPendingKind = "cleanup"
)
type hotPathStageToolResult struct {
ProviderCallID string
Body string
IsError bool
}
type hotPathStageExchange struct {
Output normalizedStageOutput
Results []hotPathStageToolResult
}
type hotPathPendingCall struct {
publicCallID string
providerCallID string
payload *workspaceEncodedPayload
}
type hotPathLightRecord struct {
requestID string
ownerEdgeID string
principalRef string
protocol string
lineage logicalRequestLineage
immutableTask string
tools []any
binding *workspaceBinding
preset config.ExecutionPreset
dispatch routeDispatch
selectorStageID string
selectorCommit hotPathStageCorrelation
localStageID string
localCommit hotPathStageCorrelation
reviewStageID string
cleanupStageID string
phase hotPathLightPhase
artifactReady bool
localPlanRead bool
workerReviewWritten bool
reviewerPlanRead bool
reviewerReviewRead bool
reviewerInspected bool
pendingLocalCommit *hotPathStageCorrelation
running bool
pendingKind hotPathPendingKind
pending map[string]hotPathPendingCall
pendingHash string
pendingOutput normalizedStageOutput
consumedHashes map[string]struct{}
consumedIDs map[string]struct{}
localTranscript []hotPathStageExchange
reviewTranscript []hotPathStageExchange
cleanupTransitions int
terminalIntent *hotPathTerminalIntent
terminalDisposition *hotPathTerminalDisposition
}
type hotPathLightStore struct {
mu sync.Mutex
capacity int
records map[string]*hotPathLightRecord
}
type hotPathDispatchSnapshot struct {
RequestID string
OwnerEdgeID string
PrincipalRef string
Protocol string
Phase hotPathLightPhase
StageID string
Stage config.ExecutionRouteStage
Route routeDispatch
PresetRoute routeDispatch
Input hotPathStageInput
Tools []any
Transcript []hotPathStageExchange
Stream bool
// RequiresCollectedProjection prevents canonical provider tool deltas from
// escaping before a command-mode caller binding can encode them.
RequiresCollectedProjection bool
// OutputBudget is recalculated from the request-local outer accumulator
// before every stage. Limited, remaining, and exhausted are distinct so an
// exhausted turn cannot be encoded as a one-token provider request.
OutputBudget hotPathOutputBudget
}
type hotPathLightDisposition struct {
RequestID string
StageID string
Phase hotPathLightPhase
TransitionFrom hotPathLightPhase
Terminal *hotPathTerminalIntent
}
func newHotPathLightStore(capacity int) *hotPathLightStore {
if capacity <= 0 {
capacity = defaultHotPathLightCapacity
}
return &hotPathLightStore{capacity: capacity, records: make(map[string]*hotPathLightRecord)}
}
func (s *hotPathLightStore) pin(
requestID, ownerEdgeID, principalRef, protocol, selectorStageID string,
lineage logicalRequestLineage,
task string,
tools any,
binding *workspaceBinding,
preset config.ExecutionPreset,
dispatch routeDispatch,
) error {
if s == nil || binding == nil {
return fmt.Errorf("light flow binding is unavailable")
}
if !validLogicalRequestID(requestID) || !validLogicalRequestID(selectorStageID) {
return fmt.Errorf("light flow identity is invalid")
}
immutableTools, err := cloneHotPathTools(tools)
if err != nil {
return err
}
if strings.TrimSpace(task) == "" {
return fmt.Errorf("light flow immutable task is empty")
}
s.mu.Lock()
defer s.mu.Unlock()
if _, exists := s.records[requestID]; exists {
return fmt.Errorf("light flow already exists")
}
if len(s.records) >= s.capacity {
return fmt.Errorf("light flow capacity reached")
}
s.records[requestID] = &hotPathLightRecord{
requestID: requestID, ownerEdgeID: ownerEdgeID, principalRef: principalRef,
protocol: protocol, lineage: lineage, immutableTask: strings.TrimSpace(task),
tools: immutableTools, binding: binding, preset: preset.Clone(), dispatch: cloneHotPathDispatch(dispatch),
selectorStageID: selectorStageID, phase: hotPathPhaseAwaitArtifacts,
consumedHashes: make(map[string]struct{}), consumedIDs: make(map[string]struct{}),
}
return nil
}
func cloneHotPathTools(tools any) ([]any, error) {
raw, err := json.Marshal(tools)
if err != nil {
return nil, fmt.Errorf("clone light flow tools: %w", err)
}
var out []any
decoder := json.NewDecoder(strings.NewReader(string(raw)))
decoder.UseNumber()
if err := decoder.Decode(&out); err != nil {
return nil, fmt.Errorf("clone light flow tools: %w", err)
}
return out, nil
}
func cloneHotPathDispatch(dispatch routeDispatch) routeDispatch {
out := dispatch
out.Preset = dispatch.Preset.Clone()
if dispatch.PresetResolvedBindings != nil {
out.PresetResolvedBindings = make(map[string]routeDispatch, len(dispatch.PresetResolvedBindings))
for key, binding := range dispatch.PresetResolvedBindings {
binding.Preset = binding.Preset.Clone()
binding.PresetResolvedBindings = nil
out.PresetResolvedBindings[key] = binding
}
}
return out
}
func (s *hotPathLightStore) remove(requestID, ownerEdgeID string) {
if s == nil || requestID == "" {
return
}
s.mu.Lock()
defer s.mu.Unlock()
if record := s.records[requestID]; record != nil && record.ownerEdgeID == ownerEdgeID {
delete(s.records, requestID)
}
}
func (s *hotPathLightStore) has(requestID, ownerEdgeID string) bool {
if s == nil || requestID == "" {
return false
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
return record != nil && record.ownerEdgeID == ownerEdgeID
}
func (s *hotPathLightStore) updateArtifactLineage(requestID, ownerEdgeID string, lineage logicalRequestLineage, localEligible bool) error {
if s == nil {
return fmt.Errorf("light flow is unavailable")
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID {
return fmt.Errorf("light flow state is unavailable")
}
record.lineage = lineage
if localEligible {
record.artifactReady = true
}
return nil
}
func (s *hotPathLightStore) commitSelector(requestID, ownerEdgeID string, output normalizedStageOutput, gate hotPathSelectorGate) error {
if s == nil {
return fmt.Errorf("light flow is unavailable")
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID || record.phase != hotPathPhaseAwaitArtifacts {
return fmt.Errorf("light flow selector commit is unavailable")
}
if strings.TrimSpace(output.ResponseID) == "" || strings.TrimSpace(gate.RunID) == "" {
return fmt.Errorf("light flow selector correlation is incomplete")
}
record.selectorCommit = hotPathStageCorrelation{
StageID: record.selectorStageID, ResponseID: output.ResponseID, RunID: gate.RunID,
ProviderID: gate.ProviderID, Terminal: output.TerminalReason,
}
return nil
}
func (s *hotPathLightStore) startLocal(requestID, ownerEdgeID string, coordinator *logicalRequestCoordinator) (hotPathLightDisposition, error) {
if s == nil || coordinator == nil {
return hotPathLightDisposition{}, fmt.Errorf("light flow is unavailable")
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID {
return hotPathLightDisposition{}, fmt.Errorf("light flow state is unavailable")
}
if record.phase != hotPathPhaseAwaitArtifacts || !record.artifactReady || strings.TrimSpace(record.selectorCommit.ResponseID) == "" {
return hotPathLightDisposition{}, fmt.Errorf("light flow is not eligible for local execution")
}
stageID, err := coordinator.newStageID()
if err != nil {
return hotPathLightDisposition{}, err
}
if _, err := coordinator.activateStage(requestID, ownerEdgeID, stageID); err != nil {
return hotPathLightDisposition{}, err
}
record.localStageID = stageID
record.phase = hotPathPhaseLocalActive
return hotPathLightDisposition{RequestID: requestID, StageID: stageID, Phase: record.phase}, nil
}
func (s *hotPathLightStore) beginDispatch(requestID, ownerEdgeID string, stream bool) (hotPathDispatchSnapshot, error) {
if s == nil {
return hotPathDispatchSnapshot{}, fmt.Errorf("light flow is unavailable")
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID {
return hotPathDispatchSnapshot{}, fmt.Errorf("light flow state is unavailable")
}
if record.running || record.pending != nil || record.phase == hotPathPhaseCleanupPending || record.phase == hotPathPhaseAwaitArtifacts {
return hotPathDispatchSnapshot{}, fmt.Errorf("light flow stage is not dispatchable")
}
stage, route, stageID, input, transcript, err := record.dispatchValues()
if err != nil {
return hotPathDispatchSnapshot{}, err
}
record.running = true
return hotPathDispatchSnapshot{
RequestID: requestID, OwnerEdgeID: ownerEdgeID, PrincipalRef: record.principalRef,
Protocol: record.protocol, Phase: record.phase, StageID: stageID, Stage: stage,
Route: route, PresetRoute: cloneHotPathDispatch(record.dispatch), Input: input,
Tools: cloneAnySlice(record.tools), Transcript: cloneStageTranscript(transcript), Stream: stream,
RequiresCollectedProjection: record.binding.requiresCollectedProjection(),
}, nil
}
func (r *hotPathLightRecord) dispatchValues() (config.ExecutionRouteStage, routeDispatch, string, hotPathStageInput, []hotPathStageExchange, error) {
route, ok := r.preset.Routes[config.ModeLight]
if !ok || len(route.Stages) != 2 {
return config.ExecutionRouteStage{}, routeDispatch{}, "", hotPathStageInput{}, nil, fmt.Errorf("light route requires local and review stages")
}
paths := newReservedPaths(r.requestID)
switch r.phase {
case hotPathPhaseLocalActive:
stage := route.Stages[0].Clone()
binding, ok := r.dispatch.PresetResolvedBindings[stage.Model]
if !ok {
return config.ExecutionRouteStage{}, routeDispatch{}, "", hotPathStageInput{}, nil, fmt.Errorf("local stage binding is unavailable")
}
return stage, binding, r.localStageID, buildLocalStageInput(r.immutableTask, paths, r.selectorCommit), r.localTranscript, nil
case hotPathPhaseReviewActive, hotPathPhaseReviewRepair:
stage := route.Stages[1].Clone()
binding, ok := r.dispatch.PresetResolvedBindings[stage.Model]
if !ok {
return config.ExecutionRouteStage{}, routeDispatch{}, "", hotPathStageInput{}, nil, fmt.Errorf("review stage binding is unavailable")
}
return stage, binding, r.reviewStageID, buildReviewStageInput(r.immutableTask, paths, r.selectorCommit, r.localCommit), r.reviewTranscript, nil
default:
return config.ExecutionRouteStage{}, routeDispatch{}, "", hotPathStageInput{}, nil, fmt.Errorf("phase %q is not dispatchable", r.phase)
}
}
func cloneAnySlice(values []any) []any {
if values == nil {
return nil
}
out := make([]any, len(values))
for i, value := range values {
out[i] = cloneAnyValue(value)
}
return out
}
func cloneStageTranscript(values []hotPathStageExchange) []hotPathStageExchange {
out := make([]hotPathStageExchange, len(values))
for i, value := range values {
out[i].Output = cloneNormalizedStageOutput(value.Output)
out[i].Results = append([]hotPathStageToolResult(nil), value.Results...)
}
return out
}
func cloneNormalizedStageOutput(value normalizedStageOutput) normalizedStageOutput {
out := value
out.Deltas = append([]normalizedStageDelta(nil), value.Deltas...)
out.ToolCalls = make([]normalizedToolCall, len(value.ToolCalls))
for i, call := range value.ToolCalls {
out.ToolCalls[i] = call
out.ToolCalls[i].Arguments = cloneAnyMap(call.Arguments)
}
out.Usage = cloneRawJSON(value.Usage)
if value.OpenAIUsage != nil {
usage := *value.OpenAIUsage
out.OpenAIUsage = &usage
}
return out
}
func (s *hotPathLightStore) abortDispatch(requestID, ownerEdgeID string) {
if s == nil {
return
}
s.mu.Lock()
defer s.mu.Unlock()
if record := s.records[requestID]; record != nil && record.ownerEdgeID == ownerEdgeID {
record.running = false
}
}
func (s *hotPathLightStore) abortWithDisposition(requestID, ownerEdgeID string, disposition hotPathTerminalDisposition) {
if s == nil || !disposition.valid() {
return
}
s.mu.Lock()
defer s.mu.Unlock()
if record := s.records[requestID]; record != nil && record.ownerEdgeID == ownerEdgeID {
record.running = false
selected := disposition
record.terminalDisposition = &selected
}
}
func (s *hotPathLightStore) issueTools(
ctx context.Context,
requestID, ownerEdgeID string,
output normalizedStageOutput,
visible normalizedStageOutput,
kind hotPathPendingKind,
outer *hotPathOuterTurn,
coordinator *logicalRequestCoordinator,
localCommit *hotPathStageCorrelation,
) (normalizedStageOutput, error) {
if s == nil || coordinator == nil {
return normalizedStageOutput{}, fmt.Errorf("light flow is unavailable")
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID || !record.running || record.pending != nil {
return normalizedStageOutput{}, fmt.Errorf("light flow tool frontier is unavailable")
}
preallocated := make(map[string]string)
if outer != nil && output.ProgressivelyReleased {
for _, call := range outer.accumulator().ToolCalls {
preallocated[call.ProviderCallID] = call.ID
}
}
mapped, pending, err := mapHotPathStageCalls(record, output, kind, coordinator, preallocated)
if err != nil {
return normalizedStageOutput{}, err
}
stageID := record.localStageID
if kind != hotPathPendingLocalTools && kind != hotPathPendingLocalHandoff {
stageID = record.reviewStageID
}
if outer != nil {
if !output.ProgressivelyReleased && !output.CallerStageOnly {
if err := runHotPathCollectedStage(ctx, outer, stageID, mapped); err != nil {
return normalizedStageOutput{}, fmt.Errorf("collect light tool outer turn: %w", err)
}
}
current := hotPathCompatibilityOutput(outer, mapped, record.protocol)
if len(current.ToolCalls) == 0 && outer.outputBudget().Exhausted {
outer.commitLengthTerminal()
return hotPathCompatibilityOutput(outer, mapped.StageResponseOverlay(visible), record.protocol), nil
}
if !output.CallerStageOnly {
if err := outer.projectToolIdentities(mapped.ToolCalls); err != nil {
return normalizedStageOutput{}, err
}
}
}
mapped = mapped.StageResponseOverlay(visible)
if outer != nil {
mapped = hotPathCompatibilityOutput(outer, mapped, record.protocol)
}
issuedHash, err := directIssuedCallHash(record.protocol, mapped)
if err != nil {
return normalizedStageOutput{}, err
}
expected := make([]logicalRequestExpectedTool, 0, len(mapped.ToolCalls))
for _, call := range mapped.ToolCalls {
expected = append(expected, logicalRequestExpectedTool{PublicCallID: call.ID, ProviderCallID: call.ProviderCallID})
}
if _, err := coordinator.awaitToolResults(requestID, ownerEdgeID, stageID, expected, issuedHash); err != nil {
return normalizedStageOutput{}, err
}
if kind == hotPathPendingLocalHandoff {
if localCommit == nil {
return normalizedStageOutput{}, fmt.Errorf("worker review handoff commit correlation is unavailable")
}
commit := *localCommit
commit.StageID = record.localStageID
record.pendingLocalCommit = &commit
}
record.pendingKind = kind
record.pending = pending
record.pendingHash = issuedHash
record.pendingOutput = cloneNormalizedStageOutput(output)
record.running = false
return mapped, nil
}
func mapHotPathStageCalls(record *hotPathLightRecord, output normalizedStageOutput, kind hotPathPendingKind, coordinator *logicalRequestCoordinator, preallocated map[string]string) (normalizedStageOutput, map[string]hotPathPendingCall, error) {
if len(output.ToolCalls) == 0 {
return normalizedStageOutput{}, nil, fmt.Errorf("light flow tool output is empty")
}
mappedCalls := make([]normalizedToolCall, 0, len(output.ToolCalls))
pending := make(map[string]hotPathPendingCall, len(output.ToolCalls))
paths := newReservedPaths(record.requestID)
for _, call := range output.ToolCalls {
providerID := strings.TrimSpace(call.ProviderCallID)
if providerID == "" {
providerID = strings.TrimSpace(call.ID)
}
if !validLogicalRequestID(providerID) {
return normalizedStageOutput{}, nil, fmt.Errorf("stage provider tool id is invalid")
}
publicID := strings.TrimSpace(preallocated[providerID])
if publicID != "" && !validLogicalRequestID(publicID) {
return normalizedStageOutput{}, nil, fmt.Errorf("stage public tool id is invalid")
}
operation, requiredPath, reserved, err := hotPathWorkspaceCall(kind, paths, call)
if err != nil {
return normalizedStageOutput{}, nil, err
}
var mapped normalizedToolCall
var payload *workspaceEncodedPayload
if reserved {
bound := record.binding.operation(operation)
if bound == nil || strings.TrimSpace(bound.toolName) == "" {
return normalizedStageOutput{}, nil, fmt.Errorf("artifact %s binding is unavailable", operation)
}
callerCall := call
callerCall.Name = bound.toolName
callerCall.Arguments = cloneAnyMap(call.Arguments)
setMappedArgument(callerCall.Arguments, bound.pathField, requiredPath)
mapped, payload, err = mapArtifactCall(record.binding, callerCall, operation, requiredPath, coordinator)
if err != nil {
return normalizedStageOutput{}, nil, err
}
if publicID != "" {
mapped.ID = publicID
payload.publicCallID = publicID
payload.correlationDigest = computePayloadCorrelationDigest(payload)
}
} else {
if !hotPathToolAllowed(record.tools, call.Name) {
return normalizedStageOutput{}, nil, fmt.Errorf("stage tool %q is not in the immutable caller tool set", call.Name)
}
if publicID == "" {
var allocErr error
publicID, allocErr = coordinator.newCallID()
if allocErr != nil {
return normalizedStageOutput{}, nil, allocErr
}
}
mapped = call
mapped.ID = publicID
mapped.ProviderCallID = providerID
mapped.Arguments = cloneAnyMap(call.Arguments)
}
mappedCalls = append(mappedCalls, mapped)
pending[mapped.ID] = hotPathPendingCall{publicCallID: mapped.ID, providerCallID: providerID, payload: payload}
}
mapped := cloneNormalizedStageOutput(output)
mapped.ToolCalls = mappedCalls
if record.protocol == "anthropic" {
mapped.TerminalReason = "tool_use"
} else {
mapped.TerminalReason = "tool_calls"
}
return mapped, pending, nil
}
func hotPathToolAllowed(tools []any, name string) bool {
schemas, err := normalizeToolSchemas(tools)
if err != nil {
return false
}
_, ok := schemas[strings.TrimSpace(name)]
return ok
}
func hotPathWorkspaceCall(kind hotPathPendingKind, paths reservedPaths, call normalizedToolCall) (workspaceOperationKind, string, bool, error) {
reserved := reservedPathsFromToolCall(call)
if len(reserved) != 1 {
return "", "", false, nil
}
observed := cleanRelativePath(reserved[0])
switch kind {
case hotPathPendingLocalTools:
if observed == cleanRelativePath(paths.PlanPath) {
return opKindRead, observed, true, nil
}
case hotPathPendingLocalHandoff:
if observed == cleanRelativePath(paths.ReviewPath) {
return opKindWrite, paths.ReviewPath, true, nil
}
case hotPathPendingReviewInspection:
if observed == cleanRelativePath(paths.PlanPath) || observed == cleanRelativePath(paths.ReviewPath) {
return opKindRead, observed, true, nil
}
}
return "", "", false, nil
}
func (s *hotPathLightStore) consumeChat(ownerEdgeID, principalRef string, rawBody []byte, lineage logicalRequestContinuationLineage, coordinator *logicalRequestCoordinator) (logicalRequestSnapshot, hotPathLightDisposition, bool, error) {
results, err := decodeChatWorkspaceResults(rawBody)
if err != nil {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, err
}
return s.consume(ownerEdgeID, principalRef, "openai", lineage, results, coordinator)
}
func (s *hotPathLightStore) consumeAnthropic(ownerEdgeID, principalRef string, rawBody []byte, lineage logicalRequestContinuationLineage, coordinator *logicalRequestCoordinator) (logicalRequestSnapshot, hotPathLightDisposition, bool, error) {
results, err := decodeAnthropicWorkspaceResults(rawBody)
if err != nil {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, err
}
return s.consume(ownerEdgeID, principalRef, "anthropic", lineage, results, coordinator)
}
func (s *hotPathLightStore) consume(ownerEdgeID, principalRef, protocol string, lineage logicalRequestContinuationLineage, results []workspaceResult, coordinator *logicalRequestCoordinator) (logicalRequestSnapshot, hotPathLightDisposition, bool, error) {
if s == nil || coordinator == nil {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, false, nil
}
s.mu.Lock()
defer s.mu.Unlock()
record, matched, err := s.matchRecordLocked(ownerEdgeID, principalRef, protocol, lineage)
if !matched || err != nil {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, matched, err
}
if record.phase == hotPathPhaseCleanupPending && record.pendingKind == hotPathPendingCleanup {
return s.consumeCleanupLocked(record, lineage, results, coordinator)
}
if record.pending == nil || record.pendingHash == "" || len(results) != len(record.pending) {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, fmt.Errorf("light tool result set mismatch")
}
byPublic := make(map[string]workspaceResult, len(results))
for _, result := range results {
pending, ok := record.pending[result.callID]
if !ok {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, fmt.Errorf("light tool result id is not pending")
}
if _, duplicate := byPublic[result.callID]; duplicate {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, fmt.Errorf("light tool result id is duplicated")
}
if pending.payload != nil {
receipt := matchResultReceipt(record.binding, pending.payload, result)
if !receipt.matched {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, fmt.Errorf("light workspace receipt rejected: %s", receipt.mismatchReason)
}
}
byPublic[result.callID] = result
}
// Command-mode SDKs may reserialize earlier message history between tool
// turns. The pending receipt already binds this continuation to the exact
// request, owner, principal, endpoint, toolset, issued calls, and result IDs.
snap, err := coordinator.consumeArtifactContinuationByLineage(record.requestID, ownerEdgeID, principalRef, lineage)
if err != nil {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, err
}
stageResults := make([]hotPathStageToolResult, 0, len(record.pendingOutput.ToolCalls))
for _, providerCall := range record.pendingOutput.ToolCalls {
providerID := strings.TrimSpace(providerCall.ProviderCallID)
if providerID == "" {
providerID = providerCall.ID
}
var pending hotPathPendingCall
var result workspaceResult
for publicID, item := range record.pending {
if item.providerCallID == providerID {
pending = item
result = byPublic[publicID]
break
}
}
if pending.providerCallID == "" {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, fmt.Errorf("light provider result correlation is unavailable")
}
stageResults = append(stageResults, hotPathStageToolResult{ProviderCallID: providerID, Body: string(result.body), IsError: result.status == "error"})
}
exchange := hotPathStageExchange{Output: cloneNormalizedStageOutput(record.pendingOutput), Results: stageResults}
pendingKind := record.pendingKind
if pendingKind == hotPathPendingLocalTools || pendingKind == hotPathPendingLocalHandoff {
record.localTranscript = append(record.localTranscript, exchange)
} else {
record.reviewTranscript = append(record.reviewTranscript, exchange)
}
for publicID, pending := range record.pending {
result := byPublic[publicID]
if result.status == "error" {
continue
}
if pending.payload == nil {
if pendingKind == hotPathPendingReviewInspection || pendingKind == hotPathPendingReviewRepair {
record.reviewerInspected = true
}
continue
}
path := cleanRelativePath(pending.payload.safePath)
switch pendingKind {
case hotPathPendingLocalTools:
if path == cleanRelativePath(newReservedPaths(record.requestID).PlanPath) {
record.localPlanRead = true
}
case hotPathPendingLocalHandoff:
record.workerReviewWritten = true
case hotPathPendingReviewInspection:
paths := newReservedPaths(record.requestID)
if path == cleanRelativePath(paths.PlanPath) {
record.reviewerPlanRead = true
}
if path == cleanRelativePath(paths.ReviewPath) {
record.reviewerReviewRead = true
}
}
}
for id := range record.pending {
record.consumedIDs[id] = struct{}{}
}
record.consumedHashes[record.pendingHash] = struct{}{}
record.lineage = lineage.Committed
record.pending = nil
record.pendingHash = ""
record.pendingOutput = normalizedStageOutput{}
previousPhase := record.phase
if pendingKind == hotPathPendingLocalHandoff {
if !record.workerReviewWritten || record.pendingLocalCommit == nil {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, fmt.Errorf("worker completion review projection failed")
}
reviewStageID, err := coordinator.newStageID()
if err != nil {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, err
}
record.localCommit = *record.pendingLocalCommit
record.pendingLocalCommit = nil
record.reviewStageID = reviewStageID
record.phase = hotPathPhaseReviewActive
} else {
record.phase = phaseAfterHotPathResult(pendingKind)
}
record.pendingKind = ""
stageID := record.localStageID
if record.phase != hotPathPhaseLocalActive {
stageID = record.reviewStageID
}
if _, err := coordinator.activateStage(record.requestID, record.ownerEdgeID, stageID); err != nil {
return logicalRequestSnapshot{}, hotPathLightDisposition{}, true, err
}
return snap, hotPathLightDisposition{
RequestID: record.requestID, StageID: stageID, Phase: record.phase, TransitionFrom: previousPhase,
}, true, nil
}
func (s *hotPathLightStore) cleanupStage(requestID, ownerEdgeID string) string {
if s == nil {
return ""
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID || record.phase != hotPathPhaseCleanupPending {
return ""
}
return record.cleanupStageID
}
func phaseAfterHotPathResult(kind hotPathPendingKind) hotPathLightPhase {
switch kind {
case hotPathPendingLocalTools, hotPathPendingLocalHandoff:
return hotPathPhaseLocalActive
case hotPathPendingReviewInspection:
return hotPathPhaseReviewActive
case hotPathPendingReviewRepair:
return hotPathPhaseReviewRepair
default:
return ""
}
}
func (s *hotPathLightStore) matchRecordLocked(ownerEdgeID, principalRef, protocol string, lineage logicalRequestContinuationLineage) (*hotPathLightRecord, bool, error) {
var candidates []*hotPathLightRecord
for _, record := range s.records {
pendingRelated := record.pending != nil && (record.pendingHash == lineage.IssuedCallHash || hotPathPendingIDsIntersect(record, lineage.ResultIDs) || record.lineage == lineage.Prefix)
_, consumedHash := record.consumedHashes[lineage.IssuedCallHash]
if pendingRelated || consumedHash || hotPathConsumedIDsIntersect(record, lineage.ResultIDs) {
candidates = append(candidates, record)
}
}
if len(candidates) == 0 {
return nil, false, nil
}
for _, record := range candidates {
if _, replay := record.consumedHashes[lineage.IssuedCallHash]; replay {
return nil, true, fmt.Errorf("light tool frontier replay rejected")
}
}
for _, record := range candidates {
if record.phase != hotPathPhaseCleanupPending || record.pendingKind != hotPathPendingCleanup || !hotPathPendingIDsIntersect(record, lineage.ResultIDs) {
continue
}
if record.ownerEdgeID != ownerEdgeID {
return nil, true, errLogicalRequestOwnerMismatch
}
if record.principalRef != principalRef {
return nil, true, errLogicalRequestPrincipal
}
if record.protocol != protocol || record.lineage.Endpoint != lineage.Prefix.Endpoint || record.lineage.ToolsetDigest != lineage.Prefix.ToolsetDigest {
return nil, true, errLogicalRequestLineage
}
return record, true, nil
}
for _, record := range candidates {
if record.pendingHash != lineage.IssuedCallHash {
continue
}
if record.ownerEdgeID != ownerEdgeID {
return nil, true, errLogicalRequestOwnerMismatch
}
if record.principalRef != principalRef {
return nil, true, errLogicalRequestPrincipal
}
if record.protocol != protocol {
return nil, true, fmt.Errorf("%w: protocol changed", errLogicalRequestLineage)
}
if record.lineage.Endpoint != lineage.Prefix.Endpoint || record.lineage.ToolsetDigest != lineage.Prefix.ToolsetDigest {
return nil, true, describeArtifactPrefixMismatch(record.lineage, lineage.Prefix)
}
return record, true, nil
}
for _, record := range candidates {
if record.ownerEdgeID != ownerEdgeID || record.principalRef != principalRef || record.protocol != protocol {
continue
}
if record.lineage != lineage.Prefix {
return nil, true, describeArtifactPrefixMismatch(record.lineage, lineage.Prefix)
}
if record.pendingHash != lineage.IssuedCallHash {
return nil, true, fmt.Errorf("%w: issued tool calls changed", errLogicalRequestLineage)
}
}
return nil, true, errLogicalRequestLineage
}
func hotPathPendingIDsIntersect(record *hotPathLightRecord, ids []string) bool {
for _, id := range ids {
if _, ok := record.pending[id]; ok {
return true
}
}
return false
}
func hotPathConsumedIDsIntersect(record *hotPathLightRecord, ids []string) bool {
for _, id := range ids {
if _, ok := record.consumedIDs[id]; ok {
return true
}
}
return false
}
func (s *hotPathLightStore) commitLocal(requestID, ownerEdgeID string, output normalizedStageOutput, correlation hotPathStageCorrelation, coordinator *logicalRequestCoordinator) (hotPathLightDisposition, error) {
if s == nil || coordinator == nil {
return hotPathLightDisposition{}, fmt.Errorf("light flow is unavailable")
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID || record.phase != hotPathPhaseLocalActive || !record.running || len(output.ToolCalls) != 0 {
return hotPathLightDisposition{}, fmt.Errorf("local completion cannot transition to review")
}
reviewStageID, err := coordinator.newStageID()
if err != nil {
return hotPathLightDisposition{}, err
}
if _, err := coordinator.transitionStage(requestID, ownerEdgeID, record.localStageID, reviewStageID); err != nil {
return hotPathLightDisposition{}, err
}
correlation.StageID = record.localStageID
correlation.ResponseID = output.ResponseID
correlation.Terminal = output.TerminalReason
record.localCommit = correlation
record.reviewStageID = reviewStageID
record.phase = hotPathPhaseReviewActive
record.running = false
return hotPathLightDisposition{RequestID: requestID, StageID: reviewStageID, Phase: record.phase}, nil
}
func (s *hotPathLightStore) classifyLocalOutput(requestID, ownerEdgeID string, output normalizedStageOutput) (hotPathPendingKind, error) {
if s == nil {
return "", fmt.Errorf("light flow is unavailable")
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID || record.phase != hotPathPhaseLocalActive || !record.running {
return "", fmt.Errorf("local flow state is unavailable")
}
if len(output.ToolCalls) == 0 {
return "", nil
}
paths := newReservedPaths(requestID)
reviewWrites := make([]normalizedToolCall, 0, 1)
for _, call := range output.ToolCalls {
if isWorkerReviewHandoffCall(record.binding, paths, call) {
reviewWrites = append(reviewWrites, call)
}
}
if len(reviewWrites) == 0 {
return hotPathPendingLocalTools, nil
}
for _, call := range reviewWrites {
if _, err := artifactWriteContent(record.binding, call); err != nil {
return "", err
}
}
return hotPathPendingLocalHandoff, nil
}
func isWorkerReviewHandoffCall(binding *workspaceBinding, paths reservedPaths, call normalizedToolCall) bool {
write := binding.operation(opKindWrite)
if write == nil || strings.TrimSpace(write.contentField) == "" {
return false
}
if _, ok := lookupMappedArgument(call.Arguments, write.contentField); !ok {
return false
}
for _, observed := range reservedPathsFromToolCall(call) {
if cleanRelativePath(observed) == cleanRelativePath(paths.ReviewPath) {
return true
}
}
return false
}
type hotPathReviewEvidence struct {
planRead bool
reviewRead bool
inspected bool
}
func (s *hotPathLightStore) reviewEvidence(requestID, ownerEdgeID string) (hotPathReviewEvidence, error) {
if s == nil {
return hotPathReviewEvidence{}, fmt.Errorf("light flow is unavailable")
}
s.mu.Lock()
defer s.mu.Unlock()
record := s.records[requestID]
if record == nil || record.ownerEdgeID != ownerEdgeID {
return hotPathReviewEvidence{}, fmt.Errorf("review flow state is unavailable")
}
return hotPathReviewEvidence{
planRead: record.reviewerPlanRead, reviewRead: record.reviewerReviewRead,
inspected: record.reviewerInspected,
}, nil
}
func (s *Server) runHotPathLocalEligible(w http.ResponseWriter, r *http.Request, dispatch routeDispatch, protocol string, stream bool, metadata map[string]string) error {
requestID := strings.TrimSpace(metadata["iop_logical_request_id"])
if requestID == "" {
return s.writeHotPathLightError(w, protocol, http.StatusBadRequest, "light flow request identity is unavailable")
}
if _, err := s.lightFlows.startLocal(requestID, s.edgeIDValue(), s.requestCoordinator); err != nil {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadRequest, err.Error()))
}
return s.runHotPathLightStage(w, r, dispatch, protocol, stream, requestID, hotPathOutputTokenCap(metadata))
}
func (s *Server) runHotPathLightContinuation(w http.ResponseWriter, r *http.Request, dispatch routeDispatch, protocol string, stream bool, metadata map[string]string) error {
requestID := strings.TrimSpace(metadata["iop_logical_request_id"])
if requestID == "" {
return s.writeHotPathLightError(w, protocol, http.StatusBadRequest, "light flow request identity is unavailable")
}
return s.runHotPathLightStage(w, r, dispatch, protocol, stream, requestID, hotPathOutputTokenCap(metadata))
}
func (s *Server) runHotPathLightStage(w http.ResponseWriter, r *http.Request, dispatch routeDispatch, protocol string, stream bool, requestID string, outputTokenCap int) error {
// This object is deliberately request-local. It is never stored in the
// logical-request record: a caller tool result starts a new HTTP turn and
// therefore must not retain the previous response writer or terminal.
outer := hotPathCallerOuterTurn(r, protocol, "", outputTokenCap)
if protocol == "openai" && stream {
if err := outer.setToolIDAllocator(s.requestCoordinator.newCallID); err != nil {
return err
}
if codec := hotPathChatOuterCodecFromRequest(r); codec != nil {
if err := codec.prepareProgressiveWriter(w, outer); err != nil {
return err
}
}
}
if protocol == "anthropic" && stream {
if codec := hotPathAnthropicCodecFromRequest(r); codec != nil {
if err := codec.prepareProgressiveWriter(w, outer, true); err != nil {
return err
}
}
}
var visible normalizedStageOutput
for transitions := 0; transitions < 2; transitions++ {
budget := outer.outputBudget()
if budget.Exhausted {
return s.writeHotPathLightLengthTerminal(w, r, dispatch, protocol, stream, requestID, visible)
}
if budget.MissingUsage {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadGateway,
"provider output usage is required before a later Hot Path stage"))
}
snapshot, err := s.lightFlows.beginDispatch(requestID, s.edgeIDValue(), stream)
if err != nil {
// A failed dispatch acquisition does not own the record's running
// stage, so it must not abort or transfer another caller's work.
return s.writeHotPathLightError(w, protocol, http.StatusBadRequest, err.Error())
}
snapshot.OutputBudget = budget
stageStart := time.Now()
output, correlation, err := s.dispatchHotPathStage(r.Context(), r, snapshot, outer)
stageDuration := time.Since(stageStart).Seconds()
attemptDisposition := hotPathDispositionForSuccess(output.TerminalReason, len(output.ToolCalls) > 0)
if err != nil {
attemptDisposition = hotPathDispositionForError(err)
if disposition, ok := hotPathDispositionFromError(err); ok {
attemptDisposition = disposition.Kind
}
}
// Every acquired provider attempt owns exactly one stage projection,
// including provider errors, timeouts, and caller cancellation.
s.observeHotPathStage(r.Context(), hotPathModeLight, hotPathStageKindForPhase(snapshot.Phase),
hotPathAttemptBucketForTranscript(snapshot.Transcript),
hotPathTerminalDispositionFromKind(attemptDisposition), snapshot.RequestID, snapshot.StageID,
dispatch.Preset.ID, stageDuration)
if err != nil {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointErrorForCause(protocol, http.StatusBadGateway, snapshot.StageID, err))
}
collected := false
if snapshot.Phase == hotPathPhaseLocalActive && len(output.ToolCalls) == 0 && !output.ProgressivelyReleased && !output.CallerStageOnly {
if err := runHotPathCollectedStage(r.Context(), outer, snapshot.StageID, output); err != nil {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadGateway, err.Error()))
}
collected = true
}
if snapshot.Phase == hotPathPhaseLocalActive && len(output.ToolCalls) == 0 && outer.outputBudget().Exhausted {
return s.writeHotPathLightLengthTerminal(w, r, dispatch, protocol, stream, requestID, output)
}
visible = mergeVisibleStageOutput(visible, output)
// The collector compatibility path remains the endpoint renderer until
// endpoint codecs consume released deltas directly. Feed the same
// output into the sequencer now so its usage and terminal boundary span
// local→review transitions in this HTTP turn.
if len(output.ToolCalls) == 0 && !output.ProgressivelyReleased && !output.CallerStageOnly && !collected {
if err := runHotPathCollectedStage(r.Context(), outer, snapshot.StageID, output); err != nil {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadGateway, err.Error()))
}
}
if len(output.ToolCalls) == 0 && hotPathIsProviderLengthTerminal(output.TerminalReason) {
return s.writeHotPathLightLengthTerminal(w, r, dispatch, protocol, stream, requestID, output)
}
switch snapshot.Phase {
case hotPathPhaseLocalActive:
if len(output.ToolCalls) > 0 {
kind, err := s.lightFlows.classifyLocalOutput(requestID, s.edgeIDValue(), output)
if err != nil {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadRequest, err.Error()))
}
mapped, err := s.lightFlows.issueTools(r.Context(), requestID, s.edgeIDValue(), output, visible, kind, outer, s.requestCoordinator, &correlation)
if err != nil {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadRequest, err.Error()))
}
if len(mapped.ToolCalls) == 0 && outer.outputBudget().Exhausted {
return s.writeHotPathLightLengthTerminal(w, r, dispatch, protocol, stream, requestID, mapped)
}
outer.commitTerminalSuccess(mapped.TerminalReason)
return s.writeHotPathStageResponse(w, r, dispatch, protocol, stream, requestID,
hotPathCompatibilityOutput(outer, mapped, protocol))
}
if outer.outputBudget().Exhausted {
return s.writeHotPathLightLengthTerminal(w, r, dispatch, protocol, stream, requestID, output)
}
if outer.outputBudget().MissingUsage {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadGateway,
"provider output usage is required before a later Hot Path stage"))
}
if disposition, err := s.lightFlows.commitLocal(requestID, s.edgeIDValue(), output, correlation, s.requestCoordinator); err != nil {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadRequest, err.Error()))
} else {
// Emit the local→review transition observation exactly once. The
// review stage id and bounded mode/stage-kind join the lifecycle.
s.observeHotPathLightTransition(r.Context(), hotPathStageKindReview, hotPathAttemptFirst,
disposition.RequestID, disposition.StageID, dispatch.Preset.ID)
}
// Only the reviewer owns caller-visible terminal output. Local-stage
// prose is not carried into the review completion projection.
visible = normalizedStageOutput{}
continue
default:
final, done, err := s.advanceHotPathReview(r.Context(), requestID, snapshot.Phase, output, visible, outer, protocol)
if err != nil {
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusBadRequest, err.Error()))
}
if s.lightFlows.cleanupStage(requestID, s.edgeIDValue()) != "" {
s.observeHotPathCleanupTransition(r.Context(), requestID, dispatch.Preset.ID)
}
if done {
if len(final.ToolCalls) == 0 && outer.outputBudget().Exhausted {
return s.writeHotPathLightLengthTerminal(w, r, dispatch, protocol, stream, requestID, final)
}
outer.commitTerminalSuccess(final.TerminalReason)
return s.writeHotPathStageResponse(w, r, dispatch, protocol, stream, requestID,
hotPathCompatibilityOutput(outer, final, protocol))
}
}
}
message := "light flow exceeded the fixed internal transition bound"
return s.writeHotPathPrimaryError(w, r, dispatch, protocol, stream, requestID,
hotPathLightEndpointError(protocol, http.StatusInternalServerError, message))
}
// writeHotPathLightLengthTerminal writes the endpoint response for a light-mode
// request that terminates by provider length or output-budget exhaustion without
// entering the cleanup phase, then emits its exactly-once outer terminal
// observation with the winning disposition. It is the non-cleanup peer of
// writeHotPathTerminal's cleanup-ending terminal owner: the two light sub-paths
// are disjoint (cleanup-ending vs length/budget), so a light request still emits
// exactly one terminal. Following the cleanup post-write ownership rule, the
// intended length terminal is resolved against the endpoint write result through
// resolveHotPathObservedDisposition, so a caller-canceled or timed-out response
// write wins over length instead of publishing length before the caller
// disposition can be selected. Preset state is closed before the write and the
// response write error is preserved as the return value. The resolved
// disposition is a closed enum, so raw error text never reaches logs or metric
// labels (SDD S15).
func (s *Server) writeHotPathLightLengthTerminal(w http.ResponseWriter, r *http.Request, dispatch routeDispatch, protocol string, stream bool, requestID string, output normalizedStageOutput) error {
outer := hotPathCurrentCallerOuterTurn(r, protocol)
outer.commitLengthTerminal()
s.terminalPresetRequest(requestID, s.edgeIDValue())
endpointWriteErr := s.writeHotPathStageResponse(w, r, dispatch, protocol, stream, requestID,
hotPathCompatibilityOutput(outer, output, protocol))
winning := resolveHotPathObservedDisposition(outer, hotPathTerminalDisposition{
Kind: hotPathDispositionLength, Source: "light_length",
}, endpointWriteErr)
s.observeHotPathTerminal(r.Context(), hotPathModeLight,
hotPathTerminalDispositionFromKind(winning.Kind), requestID, winning.StageID, dispatch.Preset.ID)
return endpointWriteErr
}
func hotPathIsProviderLengthTerminal(reason string) bool {
switch strings.TrimSpace(reason) {
case "length", "max_tokens":
return true
default:
return false
}
}
func (output normalizedStageOutput) StageResponseOverlay(visible normalizedStageOutput) normalizedStageOutput {
visible.ResponseID = output.ResponseID
visible.Created = output.Created
visible.ToolCalls = cloneNormalizedStageOutput(output).ToolCalls
visible.TerminalReason = output.TerminalReason
visible.Usage = cloneRawJSON(output.Usage)
visible.OpenAIUsage = output.OpenAIUsage
return visible
}
func mergeVisibleStageOutput(left, right normalizedStageOutput) normalizedStageOutput {
if strings.TrimSpace(left.ResponseID) == "" {
return cloneNormalizedStageOutput(right)
}
out := cloneNormalizedStageOutput(right)
out.Content = joinVisibleText(left.Content, right.Content)
out.Reasoning = joinVisibleText(left.Reasoning, right.Reasoning)
return out
}
func joinVisibleText(left, right string) string {
if left == "" {
return right
}
if right == "" {
return left
}
return left + "\n" + right
}
func (s *Server) writeHotPathStageResponse(w http.ResponseWriter, r *http.Request, dispatch routeDispatch, protocol string, stream bool, requestID string, output normalizedStageOutput) error {
turn := &hotPathTurn{
RequestID: requestID, OwnerEdgeID: s.edgeIDValue(), Dispatch: dispatch,
Protocol: protocol, Stream: stream, PublicModelID: dispatch.ExternalModelID,
Writer: w, Request: r,
}
turn.OuterTurn = hotPathCurrentCallerOuterTurn(r, protocol)
return s.writeDirectResponse(turn, output)
}
func hotPathCurrentCallerOuterTurn(r *http.Request, protocol string) *hotPathOuterTurn {
switch protocol {
case "openai":
if codec := hotPathChatOuterCodecFromRequest(r); codec != nil {
return codec.currentOuterTurn()
}
case "anthropic":
if codec := hotPathAnthropicCodecFromRequest(r); codec != nil {
return codec.currentOuterTurn()
}
}
return nil
}
func (s *Server) writeHotPathLightError(w http.ResponseWriter, protocol string, status int, message string) error {
disposition := hotPathTerminalDisposition{
Kind: hotPathDispositionProviderError, Cause: message, Source: "light_flow",
}
if status >= http.StatusBadRequest && status < http.StatusInternalServerError {
disposition.Kind = hotPathDispositionValidationError
}
if protocol == "anthropic" {
policy := anthropicHotPathPolicy(disposition)
writeAnthropicError(w, policy.status, policy.errorType, message)
} else {
policy := chatHotPathPolicy(disposition)
writeError(w, policy.status, policy.errorType, message)
}
return fmt.Errorf("%s", message)
}
func (s *Server) dispatchHotPathStage(ctx context.Context, r *http.Request, snapshot hotPathDispatchSnapshot, outer *hotPathOuterTurn) (normalizedStageOutput, hotPathStageCorrelation, error) {
return s.submitHotPathStage(ctx, r, snapshot, outer)
}
// Compile-time assertion that the stage dispatcher still uses the same
// surface-neutral service request type as selector dispatch.
var _ = edgeservice.ProviderPoolDispatchRequest{}