From b3f2558093b3c2788513b8ff4e27024288c79f90 Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Wed, 7 Oct 2026 17:53:54 +0000 Subject: [PATCH] Make Turn one interface agent.Turn is now the single interface the Router calls on a running turn: Cancel, CancellationOutcome, SteerWithReceipt, SubmitFunctionResult and AwaitSettlement. Session, DurableSteerer, Steerer and FunctionResultSubmitter are deleted, so the Router no longer discovers support through type assertions. An adapter that cannot perform an operation returns agent.ErrUnsupportedOperation, and the admitted capability declaration decides before the Router calls it. Every steer is durable. Core always sent durable_receipt=true, so the field, its producer in execution/delivery.go and the Router's non-durable branch (Steerer.Steer under a call timeout) are deleted, along with each adapter's Steer method and the onboarding fixture's. SteerWithReceipt always reports the write callback; the dispatch session.(agent.Steerer) and session.(agent.FunctionResultSubmitter) assertions and their contract_violation fallback are gone. Each adapter's contracts.go asserts only agent.Executor and agent.Turn, and the stray per-file assertions are removed. contracttest drives SteerWithReceipt directly. Test fakes implement the full Turn; durable_only_test.go tested the deleted field and is removed. harness-onboarding.md and runtime-protocol.md, with their zh translations, describe the one interface and durable-only steering. proto.Version is unchanged; the final merge bumps it. --- .../internal/agent/claudesdk/cancellation.go | 3 - .../agent/claudesdk/cancellation_test.go | 2 +- .../internal/agent/claudesdk/contracts.go | 10 +- .../internal/agent/claudesdk/executor.go | 3 - .../claudesdk/executor_confirmation_test.go | 2 +- .../agent/claudesdk/functions_test.go | 8 +- .../internal/agent/claudesdk/steering.go | 15 +-- .../internal/agent/claudesdk/steering_test.go | 8 +- apps/daemon/internal/agent/codex/contracts.go | 10 +- .../internal/agent/codex/executor_turn.go | 2 - apps/daemon/internal/agent/codex/session.go | 2 - .../internal/agent/codex/session_steering.go | 24 +---- .../codex/session_steering_lifecycle_test.go | 23 ++-- .../codex/session_steering_receipt_test.go | 100 +++++++----------- .../agent/codex/session_steering_test.go | 10 +- .../agent/contract_declarations_test.go | 8 +- .../internal/agent/contracttest/text.go | 6 +- apps/daemon/internal/agent/harness.go | 65 ++++-------- apps/daemon/internal/agent/mcode/contracts.go | 10 +- apps/daemon/internal/agent/mcode/executor.go | 3 - .../agent/mcode/executor_backpressure_test.go | 3 +- apps/daemon/internal/agent/mcode/session.go | 2 - apps/daemon/internal/agent/mcode/steering.go | 10 +- .../internal/agenthost/view_linux_test.go | 4 + .../internal/dispatch/cancellation_test.go | 2 +- .../dispatch/capability_admission_test.go | 4 +- .../internal/dispatch/durable_only_test.go | 69 ------------ .../dispatch/executor_cancel_receipt_test.go | 6 +- .../dispatch/executor_handoff_test.go | 5 +- .../daemon/internal/dispatch/executor_test.go | 3 + apps/daemon/internal/dispatch/functions.go | 28 ++--- .../internal/dispatch/functions_test.go | 2 +- .../dispatch/preparation_cancel_test.go | 9 +- .../dispatch/preparation_cleanup_test.go | 7 +- .../preparation_executor_fixture_test.go | 39 +++---- .../internal/dispatch/preparation_test.go | 10 +- .../prepared_handoff_mutation_test.go | 12 +-- .../dispatch/prepared_handoff_test.go | 15 ++- .../internal/dispatch/receipt_order_test.go | 8 +- .../dispatch/receipt_shutdown_test.go | 4 +- apps/daemon/internal/dispatch/router_test.go | 15 ++- apps/daemon/internal/dispatch/steering.go | 24 +---- .../internal/dispatch/steering_lifetime.go | 4 +- .../dispatch/steering_lifetime_test.go | 22 ++-- .../daemon/internal/dispatch/steering_test.go | 25 ++--- .../dispatch/workspace_directory_test.go | 3 +- .../internal/wireconformance/wire_test.go | 5 +- apps/daemon/testdata/onboarding/main.go | 7 +- contracts/agents-api/harness-onboarding.md | 21 ++-- contracts/agents-api/zh/harness-onboarding.md | 23 ++-- docs/runtime-protocol.md | 3 +- docs/zh/runtime-protocol.md | 5 +- internal/agentdaemon/proto/steering.go | 5 +- services/core/internal/execution/delivery.go | 2 +- .../integration/prepared_dispatch_test.go | 2 +- .../integration/steering_receipts_test.go | 4 +- 56 files changed, 245 insertions(+), 481 deletions(-) delete mode 100644 apps/daemon/internal/dispatch/durable_only_test.go diff --git a/apps/daemon/internal/agent/claudesdk/cancellation.go b/apps/daemon/internal/agent/claudesdk/cancellation.go index 01d18b381..1eda40d88 100644 --- a/apps/daemon/internal/agent/claudesdk/cancellation.go +++ b/apps/daemon/internal/agent/claudesdk/cancellation.go @@ -3,7 +3,6 @@ package claudesdk import ( "context" "errors" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -63,5 +62,3 @@ func (s *session) CancellationOutcome() proto.DonePayload { return proto.DonePayload{} } } - -var _ agent.Turn = (*session)(nil) diff --git a/apps/daemon/internal/agent/claudesdk/cancellation_test.go b/apps/daemon/internal/agent/claudesdk/cancellation_test.go index 5b2c991e7..401f34b67 100644 --- a/apps/daemon/internal/agent/claudesdk/cancellation_test.go +++ b/apps/daemon/internal/agent/claudesdk/cancellation_test.go @@ -44,7 +44,7 @@ func TestCancellationWaitsForDrainAndPublishesOutcome(t *testing.T) { if got := running.CancellationOutcome(); !reflect.DeepEqual(got, proto.DonePayload{}) { t.Fatal("unsettled outcome was exposed", got) } - if err := running.(*session).Steer(ctx, proto.PromptSteerPayload{InputID: "later", Input: proto.TextInput("later")}); !errors.Is(err, agent.ErrSteeringInactive) { + if err := running.SteerWithReceipt(ctx, proto.PromptSteerPayload{InputID: "later", Input: proto.TextInput("later")}, func() {}); !errors.Is(err, agent.ErrSteeringInactive) { t.Fatal("cancelled execution accepted steering", err) } if err := os.WriteFile(filepath.Join(config.StateDir, "release"), nil, 0o600); err != nil { diff --git a/apps/daemon/internal/agent/claudesdk/contracts.go b/apps/daemon/internal/agent/claudesdk/contracts.go index e425fc638..9d20175b6 100644 --- a/apps/daemon/internal/agent/claudesdk/contracts.go +++ b/apps/daemon/internal/agent/claudesdk/contracts.go @@ -2,13 +2,7 @@ package claudesdk import "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" -// Every public Harness implements each small contract explicitly. Unsupported -// extensions return agent.ErrUnsupportedOperation without native effects. var ( - _ agent.Executor = (*executor)(nil) - _ agent.Turn = (*session)(nil) - _ agent.Session = (*session)(nil) - _ agent.DurableSteerer = (*session)(nil) - _ agent.Steerer = (*session)(nil) - _ agent.FunctionResultSubmitter = (*session)(nil) + _ agent.Executor = (*executor)(nil) + _ agent.Turn = (*session)(nil) ) diff --git a/apps/daemon/internal/agent/claudesdk/executor.go b/apps/daemon/internal/agent/claudesdk/executor.go index 6c68f9387..0224611ff 100644 --- a/apps/daemon/internal/agent/claudesdk/executor.go +++ b/apps/daemon/internal/agent/claudesdk/executor.go @@ -300,9 +300,6 @@ func (s *session) AwaitSettlement(ctx context.Context) (agent.TurnSettlement, er } } -var _ agent.Executor = (*executor)(nil) -var _ agent.Turn = (*session)(nil) - // A timed-out operation retains its original Turn identity. Its callback cannot // retire a healthy successor after that operation has otherwise settled. func (s *session) invalidate() { diff --git a/apps/daemon/internal/agent/claudesdk/executor_confirmation_test.go b/apps/daemon/internal/agent/claudesdk/executor_confirmation_test.go index bba0b7099..44d703232 100644 --- a/apps/daemon/internal/agent/claudesdk/executor_confirmation_test.go +++ b/apps/daemon/internal/agent/claudesdk/executor_confirmation_test.go @@ -38,7 +38,7 @@ func TestExecutorNativeConfirmationSurvivesCleanup(t *testing.T) { if mode == "pending_input" { inputReceipt = make(chan error, 1) go func() { - inputReceipt <- turn.(*session).SteerWithReceipt(t.Context(), proto.PromptSteerPayload{InputID: "input", Input: proto.TextInput("extra")}, nil) + inputReceipt <- turn.SteerWithReceipt(t.Context(), proto.PromptSteerPayload{InputID: "input", Input: proto.TextInput("extra")}, func() {}) }() if event := <-out; event.Type != proto.TypeDelta { t.Fatal("input write barrier missing") diff --git a/apps/daemon/internal/agent/claudesdk/functions_test.go b/apps/daemon/internal/agent/claudesdk/functions_test.go index c813e715f..05eedef18 100644 --- a/apps/daemon/internal/agent/claudesdk/functions_test.go +++ b/apps/daemon/internal/agent/claudesdk/functions_test.go @@ -12,7 +12,6 @@ import ( "testing" "time" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -31,7 +30,6 @@ func TestFunctionTurnNativeReceipts(t *testing.T) { t.Fatal(err) } defer running.Cancel(context.Background()) - submitter := running.(agent.FunctionResultSubmitter) submissions := make(chan error, 2) calls := 0 failed := false @@ -48,17 +46,17 @@ func TestFunctionTurnNativeReceipts(t *testing.T) { } calls++ invalid := proto.FunctionResultPayload{DeliveryID: "delivery-" + call.CallID, CallID: call.CallID, Success: true} - if err := submitter.SubmitFunctionResult(ctx, invalid); err == nil { + if err := running.SubmitFunctionResult(ctx, invalid); err == nil { t.Fatal("missing content consumed call") } image := "https://example.invalid/image" invalid.Content = []proto.InputContent{{Type: "input_image", ImageURL: &image}} - if err := submitter.SubmitFunctionResult(ctx, invalid); err == nil { + if err := running.SubmitFunctionResult(ctx, invalid); err == nil { t.Fatal("image should fail before delivery") } first, second := "first-"+call.CallID, "second-"+call.CallID value := proto.FunctionResultPayload{DeliveryID: "delivery-" + call.CallID, CallID: call.CallID, Success: call.CallID == "b", Content: []proto.InputContent{{Type: "input_text", Text: &first}, {Type: "input_text", Text: &second}}} - go func() { submissions <- submitter.SubmitFunctionResult(ctx, value) }() + go func() { submissions <- running.SubmitFunctionResult(ctx, value) }() case proto.TypeToolCall: var tool proto.ToolCallPayload if err := event.DecodePayload(&tool); err != nil { diff --git a/apps/daemon/internal/agent/claudesdk/steering.go b/apps/daemon/internal/agent/claudesdk/steering.go index c64b55d72..4db08b1b3 100644 --- a/apps/daemon/internal/agent/claudesdk/steering.go +++ b/apps/daemon/internal/agent/claudesdk/steering.go @@ -24,15 +24,8 @@ type steeringState struct { seen map[string]bool } -var _ agent.Steerer = (*session)(nil) - -// Steer waits for native consumption, which may occur in a later native turn -// within this one SDK query. A completed stdin write is not a receipt. -func (s *session) Steer(ctx context.Context, input proto.PromptSteerPayload) error { - return s.SteerWithReceipt(ctx, input, nil) -} - -// SteerWithReceipt separates a complete bridge write from native consumption. +// SteerWithReceipt separates a complete bridge write from native consumption, +// which may occur in a later native turn within this one SDK query. func (s *session) SteerWithReceipt(ctx context.Context, input proto.PromptSteerPayload, written func()) error { select { case <-s.cancelOutput: @@ -91,9 +84,7 @@ func (s *session) SteerWithReceipt(ctx context.Context, input proto.PromptSteerP s.invalidate() return fmt.Errorf("claudesdk: input transport failed") } - if written != nil { - written() - } + written() select { case err := <-pending.receipt: return err diff --git a/apps/daemon/internal/agent/claudesdk/steering_test.go b/apps/daemon/internal/agent/claudesdk/steering_test.go index 25343a2e2..523422115 100644 --- a/apps/daemon/internal/agent/claudesdk/steering_test.go +++ b/apps/daemon/internal/agent/claudesdk/steering_test.go @@ -67,7 +67,7 @@ func TestSteeringReceiptsAndLifecycle(t *testing.T) { } }) } else { - reply <- s.Steer(receiptCtx, input) + reply <- s.SteerWithReceipt(receiptCtx, input, func() {}) } }() if mode == "blocked-write" { @@ -144,7 +144,7 @@ func TestSteeringReceiptsAndLifecycle(t *testing.T) { if mode == "duplicate-usage" && measurements != 1 { t.Fatal("duplicate measurement was published") } - if err := s.Steer(ctx, input); !errors.Is(err, agent.ErrSteeringInactive) { + if err := s.SteerWithReceipt(ctx, input, func() {}); !errors.Is(err, agent.ErrSteeringInactive) { t.Fatal("completed execution accepted input", err) } if _, err := s.AwaitSettlement(ctx); err != nil && (mode == "success" || mode == "phased" || mode == "timeout") { @@ -159,11 +159,11 @@ func TestSteeringReceiptsAndLifecycle(t *testing.T) { func TestSteeringDoesNotSendBeforeReadiness(t *testing.T) { s := &session{process: &clirunner.Process{}} - if err := s.Steer(context.Background(), proto.PromptSteerPayload{InputID: "one", Input: proto.TextInput("hello")}); !errors.Is(err, agent.ErrSteeringNotReady) { + if err := s.SteerWithReceipt(context.Background(), proto.PromptSteerPayload{InputID: "one", Input: proto.TextInput("hello")}, func() {}); !errors.Is(err, agent.ErrSteeringNotReady) { t.Fatal(err) } s.stopSteering() - if err := s.Steer(context.Background(), proto.PromptSteerPayload{InputID: "one", Input: proto.TextInput("hello")}); !errors.Is(err, agent.ErrSteeringInactive) { + if err := s.SteerWithReceipt(context.Background(), proto.PromptSteerPayload{InputID: "one", Input: proto.TextInput("hello")}, func() {}); !errors.Is(err, agent.ErrSteeringInactive) { t.Fatal(err) } } diff --git a/apps/daemon/internal/agent/codex/contracts.go b/apps/daemon/internal/agent/codex/contracts.go index 2514d88f1..a9bbbfdbb 100644 --- a/apps/daemon/internal/agent/codex/contracts.go +++ b/apps/daemon/internal/agent/codex/contracts.go @@ -2,13 +2,7 @@ package codex import "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" -// Every public Harness implements each small contract explicitly. Unsupported -// extensions return agent.ErrUnsupportedOperation without native effects. var ( - _ agent.Executor = (*Executor)(nil) - _ agent.Turn = (*Session)(nil) - _ agent.Session = (*Session)(nil) - _ agent.DurableSteerer = (*Session)(nil) - _ agent.Steerer = (*Session)(nil) - _ agent.FunctionResultSubmitter = (*Session)(nil) + _ agent.Executor = (*Executor)(nil) + _ agent.Turn = (*Session)(nil) ) diff --git a/apps/daemon/internal/agent/codex/executor_turn.go b/apps/daemon/internal/agent/codex/executor_turn.go index 60e5f6319..7a848578a 100644 --- a/apps/daemon/internal/agent/codex/executor_turn.go +++ b/apps/daemon/internal/agent/codex/executor_turn.go @@ -139,8 +139,6 @@ func (s *Session) Cancel(ctx context.Context) error { return err } -var _ agent.Turn = (*Session)(nil) - // Server callbacks retain their originating Turn, whose ID must match exactly. func (s *Session) onServerRequest(method string, handler ServerRequestHandler) { s.rpc.OnServerRequest(method, func(raw json.RawMessage, id any) (any, error) { diff --git a/apps/daemon/internal/agent/codex/session.go b/apps/daemon/internal/agent/codex/session.go index 4c05f9cc9..a4488676b 100644 --- a/apps/daemon/internal/agent/codex/session.go +++ b/apps/daemon/internal/agent/codex/session.go @@ -107,8 +107,6 @@ type Session struct { outcome cancellationOutcomeState } -var _ agent.Session = (*Session)(nil) - // --------------------------------------------------------------------------- // notification handlers // --------------------------------------------------------------------------- diff --git a/apps/daemon/internal/agent/codex/session_steering.go b/apps/daemon/internal/agent/codex/session_steering.go index 53fabc084..446b449b9 100644 --- a/apps/daemon/internal/agent/codex/session_steering.go +++ b/apps/daemon/internal/agent/codex/session_steering.go @@ -25,19 +25,9 @@ type TurnSteerParams struct { Input []UserInput `json:"input"` } -var _ agent.Steerer = (*Session)(nil) - -// Steer returns success only after Codex accepts input for this native turn. -func (s *Session) Steer(ctx context.Context, input proto.PromptSteerPayload) error { - return s.steer(ctx, input, nil) -} - -// SteerWithReceipt waits under the Run context after reporting the complete write. +// SteerWithReceipt returns success only after Codex accepts input for this +// native turn. It waits under the Run context after reporting the complete write. func (s *Session) SteerWithReceipt(ctx context.Context, input proto.PromptSteerPayload, written func()) error { - return s.steer(ctx, input, written) -} - -func (s *Session) steer(ctx context.Context, input proto.PromptSteerPayload, written func()) error { if !s.beginOperation() { return agent.ErrSteeringInactive } @@ -67,19 +57,13 @@ func (s *Session) steer(ctx context.Context, input proto.PromptSteerPayload, wri return agent.ErrSteeringRejected } params := TurnSteerParams{ThreadID: threadID, ExpectedTurnID: turnID, Input: content} - timeout := s.rpc.cfg.RequestTimeout - if written != nil { - timeout = 0 - } raw, err := s.rpc.requestWithTimeout(ctx, "turn/steer", params, func(frame any) error { if err := s.rpc.writeFrameContext(ctx, frame); err != nil { return err } - if written != nil { - written() - } + written() return nil - }, timeout, nil) + }, 0, nil) if err != nil { var rejected *JsonRpcError if errors.As(err, &rejected) { diff --git a/apps/daemon/internal/agent/codex/session_steering_lifecycle_test.go b/apps/daemon/internal/agent/codex/session_steering_lifecycle_test.go index de8ea9aaa..00dd4fa2b 100644 --- a/apps/daemon/internal/agent/codex/session_steering_lifecycle_test.go +++ b/apps/daemon/internal/agent/codex/session_steering_lifecycle_test.go @@ -25,9 +25,9 @@ func TestSteeringReceiptTimeoutAndCompletionKeepProcessAlive(t *testing.T) { } ctx, cancel := context.WithTimeout(context.Background(), timeout) defer cancel() - done := make(chan error, 1) + done, written := make(chan error, 1), make(chan struct{}) go func() { - done <- s.Steer(ctx, proto.PromptSteerPayload{InputID: "input", Input: proto.TextInput("extra")}) + done <- s.SteerWithReceipt(ctx, proto.PromptSteerPayload{InputID: "input", Input: proto.TextInput("extra")}, func() { close(written) }) }() var request JsonRpcRequest if err := json.NewDecoder(server.FromClient).Decode(&request); err != nil { @@ -36,19 +36,10 @@ func TestSteeringReceiptTimeoutAndCompletionKeepProcessAlive(t *testing.T) { // Withhold the response after reading the entire request frame. if complete { // Reading the pipe does not mean its writer has returned yet. - for { - client.pendingMu.Lock() - pending := client.pending[request.ID] - waiting := pending != nil && pending.timer != nil - client.pendingMu.Unlock() - if waiting { - break - } - select { - case <-ctx.Done(): - t.Fatal("steering request did not reach response wait") - case <-time.After(time.Millisecond): - } + select { + case <-written: + case <-ctx.Done(): + t.Fatal("steering request did not reach response wait") } s.onTurnCompleted(json.RawMessage(`{"threadId":"thread","turn":{"id":"turn","status":"completed"}}`)) } @@ -125,7 +116,7 @@ func TestBlockedSteeringWriteEndsRunWithTerminalFrames(t *testing.T) { s := turn.(*Session) callCtx, callCancel := context.WithTimeout(ctx, 50*time.Millisecond) defer callCancel() - if err := s.Steer(callCtx, proto.PromptSteerPayload{InputID: "blocked", Input: proto.TextInput("extra")}); err == nil { + if err := s.SteerWithReceipt(callCtx, proto.PromptSteerPayload{InputID: "blocked", Input: proto.TextInput("extra")}, func() {}); err == nil { t.Fatal("blocked write accepted") } if client.Alive() { diff --git a/apps/daemon/internal/agent/codex/session_steering_receipt_test.go b/apps/daemon/internal/agent/codex/session_steering_receipt_test.go index bb14579d9..7432577ce 100644 --- a/apps/daemon/internal/agent/codex/session_steering_receipt_test.go +++ b/apps/daemon/internal/agent/codex/session_steering_receipt_test.go @@ -10,66 +10,46 @@ import ( ) func TestDurableSteeringBypassesOnlyNativeResponseDeadline(t *testing.T) { - for _, durable := range []bool{false, true} { - t.Run(map[bool]string{false: "legacy", true: "durable"}[durable], func(t *testing.T) { - client, server, cleanup := NewTestClient() - defer cleanup() - client.cfg.RequestTimeout = 20 * time.Millisecond - s := &Session{rpc: client.JSONRPCClient, cancelCtx: context.Background()} - s.setThreadID("thread") - s.onTurnStarted(json.RawMessage(`{"threadId":"thread","turn":{"id":"turn"}}`)) - ctx, cancel := context.WithTimeout(context.Background(), time.Second) - defer cancel() - written := make(chan struct{}) - reply := make(chan error, 1) - go func() { - input := proto.PromptSteerPayload{InputID: "extra", Input: proto.TextInput("text")} - if durable { - reply <- s.SteerWithReceipt(ctx, input, func() { close(written) }) - } else { - reply <- s.Steer(ctx, input) - } - }() - if durable { - select { - case <-written: - t.Fatal("written before frame was read") - default: - } - } - var request JsonRpcRequest - if err := json.NewDecoder(server.FromClient).Decode(&request); err != nil { - t.Fatal(err) - } - if durable { - select { - case <-written: - case <-ctx.Done(): - t.Fatal("missing write phase") - } - } - select { - case err := <-reply: - if durable || err == nil { - t.Fatal("unexpected response deadline", err) - } - case <-time.After(60 * time.Millisecond): - if !durable { - t.Fatal("legacy timeout disappeared") - } - } - if durable { - if err := json.NewEncoder(server.ToClient).Encode(map[string]any{"id": request.ID, "result": map[string]string{"turnId": "turn"}}); err != nil { - t.Fatal(err) - } - if err := <-reply; err != nil { - t.Fatal(err) - } - } - if !client.Alive() { - t.Fatal("receipt wait killed process") - } - }) + client, server, cleanup := NewTestClient() + defer cleanup() + client.cfg.RequestTimeout = 20 * time.Millisecond + s := &Session{rpc: client.JSONRPCClient, cancelCtx: context.Background()} + s.setThreadID("thread") + s.onTurnStarted(json.RawMessage(`{"threadId":"thread","turn":{"id":"turn"}}`)) + ctx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + written := make(chan struct{}) + reply := make(chan error, 1) + go func() { + reply <- s.SteerWithReceipt(ctx, proto.PromptSteerPayload{InputID: "extra", Input: proto.TextInput("text")}, func() { close(written) }) + }() + select { + case <-written: + t.Fatal("written before frame was read") + default: + } + var request JsonRpcRequest + if err := json.NewDecoder(server.FromClient).Decode(&request); err != nil { + t.Fatal(err) + } + select { + case <-written: + case <-ctx.Done(): + t.Fatal("missing write phase") + } + select { + case err := <-reply: + t.Fatal("unexpected response deadline", err) + case <-time.After(60 * time.Millisecond): + } + if err := json.NewEncoder(server.ToClient).Encode(map[string]any{"id": request.ID, "result": map[string]string{"turnId": "turn"}}); err != nil { + t.Fatal(err) + } + if err := <-reply; err != nil { + t.Fatal(err) + } + if !client.Alive() { + t.Fatal("receipt wait killed process") } } diff --git a/apps/daemon/internal/agent/codex/session_steering_test.go b/apps/daemon/internal/agent/codex/session_steering_test.go index 1fe0fcc47..524398938 100644 --- a/apps/daemon/internal/agent/codex/session_steering_test.go +++ b/apps/daemon/internal/agent/codex/session_steering_test.go @@ -33,7 +33,7 @@ func TestSteeringUsesNativeActiveTurnAndReceipt(t *testing.T) { s.onTurnStarted(json.RawMessage(`{"threadId":"native-thread","turn":{"id":"native-turn"}}`)) done := make(chan error, 1) go func() { - done <- s.Steer(ctx, proto.PromptSteerPayload{InputID: "input-1", Input: proto.TextInput("追加输入")}) + done <- s.SteerWithReceipt(ctx, proto.PromptSteerPayload{InputID: "input-1", Input: proto.TextInput("追加输入")}, func() {}) }() var request struct { ID string `json:"id"` @@ -78,7 +78,7 @@ func TestSteeringDeadlineReleasesBlockedNativeWrite(t *testing.T) { defer cancel() done := make(chan error, 1) go func() { - done <- s.Steer(ctx, proto.PromptSteerPayload{InputID: "blocked", Input: proto.TextInput("extra")}) + done <- s.SteerWithReceipt(ctx, proto.PromptSteerPayload{InputID: "blocked", Input: proto.TextInput("extra")}, func() {}) }() // No reader drains the pipe, so the request never reaches its response wait. select { @@ -103,18 +103,18 @@ func TestSteeringDoesNotStartOrReviveTurns(t *testing.T) { // A nil RPC client proves these states do not issue a request. for _, notification := range []json.RawMessage{nil, json.RawMessage(`{"threadId":"other-thread","turn":{"id":"other-turn"}}`)} { s.onTurnStarted(notification) - if err := s.Steer(ctx, input); !errors.Is(err, agent.ErrSteeringNotReady) { + if err := s.SteerWithReceipt(ctx, input, func() {}); !errors.Is(err, agent.ErrSteeringNotReady) { t.Fatalf("starting turn: %v", err) } } s.stopSteering() s.onTurnStarted(json.RawMessage(`{"threadId":"native-thread","turn":{"id":"late-turn"}}`)) - if err := s.Steer(ctx, input); !errors.Is(err, agent.ErrSteeringInactive) { + if err := s.SteerWithReceipt(ctx, input, func() {}); !errors.Is(err, agent.ErrSteeringInactive) { t.Fatalf("terminal turn revived: %v", err) } s = &Session{cancelCtx: ctx} cancel() - if err := s.Steer(context.Background(), input); !errors.Is(err, agent.ErrSteeringInactive) { + if err := s.SteerWithReceipt(context.Background(), input, func() {}); !errors.Is(err, agent.ErrSteeringInactive) { t.Fatalf("cancelled run: %v", err) } } diff --git a/apps/daemon/internal/agent/contract_declarations_test.go b/apps/daemon/internal/agent/contract_declarations_test.go index fe8ec74e7..f007af24a 100644 --- a/apps/daemon/internal/agent/contract_declarations_test.go +++ b/apps/daemon/internal/agent/contract_declarations_test.go @@ -16,12 +16,8 @@ import ( // decision. Each adapter's compile assertions then enforce its actual methods. func TestPublicHarnessContractDeclarations(t *testing.T) { roles := map[string][]string{ - "Executor": {"executor"}, - "Turn": {"session"}, - "Session": {"session"}, - "DurableSteerer": {"session"}, - "Steerer": {"session"}, - "FunctionResultSubmitter": {"session"}, + "Executor": {"executor"}, + "Turn": {"session"}, } files, err := filepath.Glob("*.go") if err != nil { diff --git a/apps/daemon/internal/agent/contracttest/text.go b/apps/daemon/internal/agent/contracttest/text.go index c389163a7..ee38f754b 100644 --- a/apps/daemon/internal/agent/contracttest/text.go +++ b/apps/daemon/internal/agent/contracttest/text.go @@ -66,12 +66,8 @@ func TextLifecycle(t *testing.T, fixture TextFixture) { if err := first.Cancel(ctx); err != nil { t.Fatal(err) } - steerer, ok := second.(agent.DurableSteerer) - if !ok { - t.Fatal("public text Turn must implement DurableSteerer") - } var written atomic.Int32 - if err := steerer.SteerWithReceipt(ctx, proto.PromptSteerPayload{InputID: "contract-input", Input: fixture.SteeringInput, DurableReceipt: true}, func() { written.Add(1) }); err != nil { + if err := second.SteerWithReceipt(ctx, proto.PromptSteerPayload{InputID: "contract-input", Input: fixture.SteeringInput}, func() { written.Add(1) }); err != nil { t.Fatalf("active input after stale cancellation: %v", err) } if written.Load() != 1 { diff --git a/apps/daemon/internal/agent/harness.go b/apps/daemon/internal/agent/harness.go index 719485404..565a95816 100644 --- a/apps/daemon/internal/agent/harness.go +++ b/apps/daemon/internal/agent/harness.go @@ -3,11 +3,10 @@ // contracts/agents-api/harness-onboarding.md. // // Required lifecycle: ExecutorFactory prepares a fixed Session-owned Executor; -// each StartTurn returns a fresh Turn with its own output and settlement. -// Turn includes the public text DurableSteerer requirement. Keep extension -// interfaces separate, but implement each explicitly: unsupported operations -// return ErrUnsupportedOperation before any native effects. Interface presence -// does not advertise support; the capability declaration controls admission. MCP, +// each StartTurn returns a fresh Turn with its own output and settlement. Turn +// is one interface: an operation the adapter does not support returns +// ErrUnsupportedOperation before any native effects, and the capability +// declaration, not the method, decides whether the Runtime calls it. MCP, // images, structured output and Subagent observations use protocol messages // rather than additional Go interfaces; qualify and advertise them separately. // @@ -552,27 +551,10 @@ type Executor interface { } // Turn owns one output stream and never retargets cancellation to a successor. -type Turn interface { - Session - DurableSteerer - // Success confirms closed output and settled native input, function and - // child-work obligations. Errors cannot prove cancellation. - AwaitSettlement(context.Context) (TurnSettlement, error) -} - -// TurnSettlement describes native readiness after this Turn has fully drained. -// A false result must carry a reason; it requires confirmed Executor.Close. -type TurnSettlement struct { - Reusable bool - Reason string -} - -// Session is the cancellation and outcome surface of a Turn. Every owner -// exposes observed state. // For Executor-owned Turns, AwaitSettlement and Executor.Close define settlement // and resource retirement; Cancel alone does not transfer resource ownership. -type Session interface { - // Cancel signals the session to abort. Idempotent. Actual teardown +type Turn interface { + // Cancel signals the Turn to abort. Idempotent. Actual teardown // happens asynchronously and is signalled via the out channel close. Cancel(ctx context.Context) error // CancellationOutcome snapshots observed native identity, usage and output. @@ -580,28 +562,23 @@ type Session interface { // result means no observed evidence, not unsupported cancellation or success. // Reading it does not wait for or establish native settlement. CancellationOutcome() proto.DonePayload + // SteerWithReceipt delivers active input: it calls written once the complete + // input is written, then waits for the native application receipt. Every + // public Harness implements it; returning Unsupported is not a receipt. + SteerWithReceipt(ctx context.Context, input proto.PromptSteerPayload, written func()) error + // SubmitFunctionResult delivers a result for an outstanding native call. + // A successful return requires its native application receipt. + SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error + // Success confirms closed output and settled native input, function and + // child-work obligations. Errors cannot prove cancellation. + AwaitSettlement(context.Context) (TurnSettlement, error) } -// Turn extension contracts. Every public Harness implements each interface; -// unsupported operations return ErrUnsupportedOperation with a fixed safe reason. - -// DurableSteerer reports one complete write synchronously, then waits for the native receipt. -// It is independent of Steerer and is mandatory on every public Turn. -// Required input receipts cannot be implemented by returning Unsupported. -type DurableSteerer interface { - SteerWithReceipt(context.Context, proto.PromptSteerPayload, func()) error -} - -// Steerer delivers non-durable input when qualified; otherwise it explicitly -// returns ErrUnsupportedOperation without submitting input. -type Steerer interface { - Steer(context.Context, proto.PromptSteerPayload) error -} - -// FunctionResultSubmitter delivers a result for an outstanding native call. -// A successful return requires its native application receipt. -type FunctionResultSubmitter interface { - SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error +// TurnSettlement describes native readiness after this Turn has fully drained. +// A false result must carry a reason; it requires confirmed Executor.Close. +type TurnSettlement struct { + Reusable bool + Reason string } // Kind registration. diff --git a/apps/daemon/internal/agent/mcode/contracts.go b/apps/daemon/internal/agent/mcode/contracts.go index 500870957..a4fd04e0b 100644 --- a/apps/daemon/internal/agent/mcode/contracts.go +++ b/apps/daemon/internal/agent/mcode/contracts.go @@ -2,13 +2,7 @@ package mcode import "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" -// Every public Harness implements each small contract explicitly. Unsupported -// extensions return agent.ErrUnsupportedOperation without native effects. var ( - _ agent.Executor = (*executor)(nil) - _ agent.Turn = (*Session)(nil) - _ agent.Session = (*Session)(nil) - _ agent.DurableSteerer = (*Session)(nil) - _ agent.Steerer = (*Session)(nil) - _ agent.FunctionResultSubmitter = (*Session)(nil) + _ agent.Executor = (*executor)(nil) + _ agent.Turn = (*Session)(nil) ) diff --git a/apps/daemon/internal/agent/mcode/executor.go b/apps/daemon/internal/agent/mcode/executor.go index d40b8964c..50294a84a 100644 --- a/apps/daemon/internal/agent/mcode/executor.go +++ b/apps/daemon/internal/agent/mcode/executor.go @@ -24,9 +24,6 @@ type executor struct { starting, closed, invalid bool } -var _ agent.Executor = (*executor)(nil) -var _ agent.Turn = (*Session)(nil) - // NewExecutorFactory fixes the deployment workspace once; nil selects none. func NewExecutorFactory(config *WorkspaceConfig) agent.ExecutorFactory { var frozen *WorkspaceConfig diff --git a/apps/daemon/internal/agent/mcode/executor_backpressure_test.go b/apps/daemon/internal/agent/mcode/executor_backpressure_test.go index 4651d46db..ba0806d2a 100644 --- a/apps/daemon/internal/agent/mcode/executor_backpressure_test.go +++ b/apps/daemon/internal/agent/mcode/executor_backpressure_test.go @@ -77,8 +77,7 @@ func TestExecutorUnknownSteeringOutcomeCannotBecomeAppliedCancellation(t *testin case <-ctx.Done(): t.Fatal("native prompt did not start") } - session := turn.(*Session) - if err := session.Steer(ctx, proto.PromptSteerPayload{InputID: "input-unknown", Input: proto.TextInput("continue")}); err == nil { + if err := turn.SteerWithReceipt(ctx, proto.PromptSteerPayload{InputID: "input-unknown", Input: proto.TextInput("continue")}, func() {}); err == nil { t.Fatal("unknown native receipt was accepted") } if err := turn.Cancel(ctx); err == nil { diff --git a/apps/daemon/internal/agent/mcode/session.go b/apps/daemon/internal/agent/mcode/session.go index ec43b1482..69bfebca1 100644 --- a/apps/daemon/internal/agent/mcode/session.go +++ b/apps/daemon/internal/agent/mcode/session.go @@ -50,8 +50,6 @@ type Session struct { subagentHistoryReady bool } -var _ agent.Session = (*Session)(nil) - func launch(ctx context.Context, req proto.PromptRequestPayload, opts launchOptions, binary string) (*Session, error) { start, args := opts.start, []string{"acp"} if start == nil { diff --git a/apps/daemon/internal/agent/mcode/steering.go b/apps/daemon/internal/agent/mcode/steering.go index 005f5c859..12870daa4 100644 --- a/apps/daemon/internal/agent/mcode/steering.go +++ b/apps/daemon/internal/agent/mcode/steering.go @@ -10,12 +10,6 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -var _ agent.DurableSteerer = (*Session)(nil) - -func (s *Session) Steer(ctx context.Context, input proto.PromptSteerPayload) error { - return s.SteerWithReceipt(ctx, input, nil) -} - // Native acceptance belongs to the active ACP Turn; it does not promise model consumption. func (s *Session) SteerWithReceipt(ctx context.Context, input proto.PromptSteerPayload, written func()) error { s.mu.Lock() @@ -53,9 +47,7 @@ func (s *Session) SteerWithReceipt(ctx context.Context, input proto.PromptSteerP s.markInputUncertain() return fmt.Errorf("mcode: input transport failed") } - if written != nil { - written() - } + written() var frame rpcFrame var stopped error select { diff --git a/apps/daemon/internal/agenthost/view_linux_test.go b/apps/daemon/internal/agenthost/view_linux_test.go index f99b145d3..0848b38b8 100644 --- a/apps/daemon/internal/agenthost/view_linux_test.go +++ b/apps/daemon/internal/agenthost/view_linux_test.go @@ -730,6 +730,10 @@ func (t *testTurn) SteerWithReceipt(context.Context, proto.PromptSteerPayload, f return agent.ErrUnsupportedOperation } +func (t *testTurn) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { + return agent.ErrUnsupportedOperation +} + func (t *testTurn) AwaitSettlement(ctx context.Context) (agent.TurnSettlement, error) { select { case <-t.settled: diff --git a/apps/daemon/internal/dispatch/cancellation_test.go b/apps/daemon/internal/dispatch/cancellation_test.go index 7161b9ef5..23346ced7 100644 --- a/apps/daemon/internal/dispatch/cancellation_test.go +++ b/apps/daemon/internal/dispatch/cancellation_test.go @@ -34,7 +34,7 @@ func (s *cancelReceiptSession) Cancel(ctx context.Context) error { } func registerCancelReceiptKind(h *harness, sess *cancelReceiptSession) { - registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { sess.fakeSession = &fakeSession{out: out, closeOutOnCancel: true} return sess, nil }) diff --git a/apps/daemon/internal/dispatch/capability_admission_test.go b/apps/daemon/internal/dispatch/capability_admission_test.go index 92c0a596a..0a10601f0 100644 --- a/apps/daemon/internal/dispatch/capability_admission_test.go +++ b/apps/daemon/internal/dispatch/capability_admission_test.go @@ -15,8 +15,8 @@ func TestSteeringDoesNotReplayUnsupportedImplementation(t *testing.T) { h := newHarness(t) defer h.router.Shutdown(context.Background()) var calls atomic.Int32 - factory := func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { - return &steeringSession{fakeSession: &fakeSession{out: out, closeOutOnCancel: true}, steer: func(context.Context, proto.PromptSteerPayload) error { + factory := func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { + return &steeringSession{fakeSession: &fakeSession{out: out, closeOutOnCancel: true}, steer: func(context.Context, proto.PromptSteerPayload, func()) error { calls.Add(1) return fmt.Errorf("%w: fixture has no active input", agent.ErrUnsupportedOperation) }}, nil diff --git a/apps/daemon/internal/dispatch/durable_only_test.go b/apps/daemon/internal/dispatch/durable_only_test.go deleted file mode 100644 index d8707a13d..000000000 --- a/apps/daemon/internal/dispatch/durable_only_test.go +++ /dev/null @@ -1,69 +0,0 @@ -package dispatch_test - -import ( - "context" - "sync/atomic" - "testing" - - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" - "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" -) - -type durableOnlyExecutor struct { - *reusableExecutor - applied atomic.Int32 -} - -func (e *durableOnlyExecutor) StartTurn(ctx context.Context, id string, input proto.MessageInput, out chan<- proto.Envelope) (agent.Turn, error) { - turn, err := e.reusableExecutor.StartTurn(ctx, id, input, out) - return &durableOnlyTurn{Turn: turn, applied: &e.applied}, err -} - -type durableOnlyTurn struct { - agent.Turn - applied *atomic.Int32 -} - -func (t *durableOnlyTurn) SteerWithReceipt(_ context.Context, _ proto.PromptSteerPayload, written func()) error { - written() - t.applied.Add(1) - return nil -} - -func TestPublicTextSteeringDoesNotRequireOptionalSteerer(t *testing.T) { - owner := &durableOnlyExecutor{reusableExecutor: &reusableExecutor{starts: make(chan *reusableTurn, 1)}} - router, sender := poolRouter(t, func(context.Context, proto.PromptRequestPayload) (agent.Executor, error) { return owner, nil }) - assign(t, router, preparationSessionID, "") - admission := executorAdmission(t, router, sender, "prepare", executorRequest()) - startExecutorTurn(t, router, sender, "prepare", "run", admission) - turn := <-owner.starts - defer turn.finish() - input := proto.PromptSteerPayload{InputID: "extra", Input: proto.TextInput("follow-up"), DurableReceipt: true} - if err := router.Handle(t.Context(), mustEnv(t, proto.TypePromptSteer, "run", input)); err != nil { - t.Fatal(err) - } - waitFor(t, func() bool { - for _, event := range sender.snapshot() { - var ack proto.PromptSteerAckPayload - if event.Type == proto.TypePromptSteerAck && event.DecodePayload(&ack) == nil && ack.InputID == input.InputID && ack.Accepted { - return true - } - } - return false - }, "durable input applied without optional Steerer") - if owner.applied.Load() != 1 { - t.Fatal("input was not applied exactly once") - } - input.DurableReceipt = false - input.InputID = "ordinary" - if err := router.Handle(t.Context(), mustEnv(t, proto.TypePromptSteer, "run", input)); err != nil { - t.Fatal(err) - } - for _, event := range sender.snapshot() { - var ack proto.PromptSteerAckPayload - if event.Type == proto.TypePromptSteerAck && event.DecodePayload(&ack) == nil && ack.InputID == "ordinary" && ack.ErrorCode == "unsupported" { - return - } - } - t.Fatal("unimplemented optional Steerer was not rejected") -} diff --git a/apps/daemon/internal/dispatch/executor_cancel_receipt_test.go b/apps/daemon/internal/dispatch/executor_cancel_receipt_test.go index 131c9a3fb..253c548e4 100644 --- a/apps/daemon/internal/dispatch/executor_cancel_receipt_test.go +++ b/apps/daemon/internal/dispatch/executor_cancel_receipt_test.go @@ -101,8 +101,8 @@ func (t *receiptCancelTurn) AwaitSettlement(ctx context.Context) (agent.TurnSett func (t *receiptCancelTurn) CancellationOutcome() proto.DonePayload { return proto.DonePayload{Content: "observed", Metadata: map[string]any{proto.DoneMetaAgentSessionID: "native-session"}} } -func (t *receiptCancelTurn) Steer(ctx context.Context, input proto.PromptSteerPayload) error { - return t.SteerWithReceipt(ctx, input, nil) +func (t *receiptCancelTurn) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { + return agent.ErrUnknownFunctionCall } func (t *receiptCancelTurn) SteerWithReceipt(ctx context.Context, _ proto.PromptSteerPayload, written func()) error { written() @@ -155,7 +155,7 @@ func TestExecutorCancellationReachesNativeBeforeDurableReceiptJoin(t *testing.T) p := executorAdmission(t, r, sender.recSender, "first", executorRequest()) startExecutorTurn(t, r, sender.recSender, "first", "run", p) turn := <-owner.turn - input := proto.PromptSteerPayload{InputID: "pending", Input: proto.TextInput("second input"), DurableReceipt: true} + input := proto.PromptSteerPayload{InputID: "pending", Input: proto.TextInput("second input")} if err := r.Handle(t.Context(), mustEnv(t, proto.TypePromptSteer, "run", input)); err != nil { t.Fatal(err) } diff --git a/apps/daemon/internal/dispatch/executor_handoff_test.go b/apps/daemon/internal/dispatch/executor_handoff_test.go index b7f89a573..918c27bc0 100644 --- a/apps/daemon/internal/dispatch/executor_handoff_test.go +++ b/apps/daemon/internal/dispatch/executor_handoff_test.go @@ -227,7 +227,10 @@ func TestPreparedDonePublishesAfterExecutorHandoff(t *testing.T) { } } -// These fixtures exercise settlement only; active input is deliberately rejected. +// These fixtures exercise settlement only; active input and results are deliberately rejected. func (*terminalHandoffTurn) SteerWithReceipt(context.Context, proto.PromptSteerPayload, func()) error { return agent.ErrSteeringRejected } +func (*terminalHandoffTurn) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { + return agent.ErrUnknownFunctionCall +} diff --git a/apps/daemon/internal/dispatch/executor_test.go b/apps/daemon/internal/dispatch/executor_test.go index 680fd02c5..3a84bac60 100644 --- a/apps/daemon/internal/dispatch/executor_test.go +++ b/apps/daemon/internal/dispatch/executor_test.go @@ -338,3 +338,6 @@ func TestExecutorRejectsOutputFromAnotherTurn(t *testing.T) { func (*reusableTurn) SteerWithReceipt(context.Context, proto.PromptSteerPayload, func()) error { return agent.ErrSteeringRejected } +func (*reusableTurn) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { + return agent.ErrUnknownFunctionCall +} diff --git a/apps/daemon/internal/dispatch/functions.go b/apps/daemon/internal/dispatch/functions.go index 569a26aaf..026aa6b9a 100644 --- a/apps/daemon/internal/dispatch/functions.go +++ b/apps/daemon/internal/dispatch/functions.go @@ -54,30 +54,20 @@ func (r *Router) handleFunctionResult(ctx context.Context, env proto.Envelope) e } } session, finishOperation, ready := r.preparedOperationLocked(state) - var submitter agent.FunctionResultSubmitter - if ready { - submitter, _ = session.(agent.FunctionResultSubmitter) - } r.mu.Unlock() - if state != nil && !ready { - return r.sendInteractionDecisionAck(ctx, env, result.DeliveryID, false, "not_ready", "function call is waiting for the native session") + if state == nil { + return r.sendInteractionDecisionAck(ctx, env, result.DeliveryID, false, "not_pending", "function call is no longer pending") } - if finishOperation != nil { - defer finishOperation() - var stop context.CancelFunc - ctx, stop = r.shutdownContext(ctx) - defer stop() + if !ready { + return r.sendInteractionDecisionAck(ctx, env, result.DeliveryID, false, "not_ready", "function call is waiting for the native session") } - if state != nil && !state.capabilities.FunctionTools.IsSupported() { + defer finishOperation() + ctx, stop := r.shutdownContext(ctx) + defer stop() + if !state.capabilities.FunctionTools.IsSupported() { return r.sendInteractionDecisionAck(ctx, env, result.DeliveryID, false, "unsupported", "The runtime declaration does not support function results.") } - if ready && submitter == nil { - return r.sendInteractionDecisionAck(ctx, env, result.DeliveryID, false, "contract_violation", "Declared function capability has no implementation.") - } - if submitter == nil { - return r.sendInteractionDecisionAck(ctx, env, result.DeliveryID, false, "not_pending", "function call is no longer pending") - } - if err := submitter.SubmitFunctionResult(ctx, result); err != nil { + if err := session.SubmitFunctionResult(ctx, result); err != nil { code := "runtime_error" if errors.Is(err, agent.ErrUnsupportedOperation) { code = "contract_violation" diff --git a/apps/daemon/internal/dispatch/functions_test.go b/apps/daemon/internal/dispatch/functions_test.go index 55fa0b339..2c6474ddc 100644 --- a/apps/daemon/internal/dispatch/functions_test.go +++ b/apps/daemon/internal/dispatch/functions_test.go @@ -38,7 +38,7 @@ func TestFunctionReceiptsScopeRetriesAndConflicts(t *testing.T) { reg := agent.NewRegistry() sender := &recSender{} sessions := map[string]*functionSession{} - registerSession(reg, proto.SupportedAgentKind{Kind: "function-test", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{FunctionTools: proto.CapabilitySupported})}, func(ctx context.Context, p proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(reg, proto.SupportedAgentKind{Kind: "function-test", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{FunctionTools: proto.CapabilitySupported})}, func(ctx context.Context, p proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { s := &functionSession{fakeSession: &fakeSession{out: out, ctx: ctx, closeOutOnCancel: true}} sessions[p.RunID] = s return s, nil diff --git a/apps/daemon/internal/dispatch/preparation_cancel_test.go b/apps/daemon/internal/dispatch/preparation_cancel_test.go index 5804650cb..0502727e6 100644 --- a/apps/daemon/internal/dispatch/preparation_cancel_test.go +++ b/apps/daemon/internal/dispatch/preparation_cancel_test.go @@ -9,7 +9,6 @@ import ( "testing" "time" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/dispatch" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -82,7 +81,7 @@ func TestPreparedCancellationWaitsForOutputAndCleanup(t *testing.T) { Content: "observed", Usage: proto.Usage{Tokens: &proto.TokenUsage{InputTokens: 7, OutputTokens: 3, TotalTokens: 10}}, Metadata: map[string]any{proto.DoneMetaAgentSessionID: "observed-native"}, }} var session *fakeSession - p.start = func(_ context.Context, id string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, id string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session = &fakeSession{out: out, closeOutOnCancel: true, postCancelEnvelopes: []proto.Envelope{mustEnv(t, proto.TypeDone, id, p.outcome)}} out <- mustEnv(t, proto.TypeDelta, id, proto.DeltaPayload{Delta: "observed"}) @@ -150,7 +149,7 @@ func TestPreparedCancellationBeforeTransferPreservesUnknownOutcome(t *testing.T) cancelOnce.Do(func() { close(cancelled) }) return p.Close() } - p.start = func(_ context.Context, _ string, _ proto.MessageInput, _ chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, _ chan<- proto.Envelope) (fixtureSession, error) { close(entered) <-cancelled <-allowReturn @@ -204,7 +203,7 @@ func TestPreparedCancellationFailuresRemainConservative(t *testing.T) { entered, allowReturn := make(chan struct{}), make(chan struct{}) var session *fakeSession p := &cancellationPreparation{controlledPreparation: &controlledPreparation{closed: make(chan struct{})}} - p.start = func(_ context.Context, id string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, id string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session = &fakeSession{out: out, closeOutOnCancel: true} out <- mustEnv(t, proto.TypeUsage, id, proto.UsagePayload{Usage: proto.Usage{InputTokens: 3}}) close(entered) @@ -255,7 +254,7 @@ func TestPreparedCancellationTimeoutKeepsCapacityUntilStartReturns(t *testing.T) entered, allowReturn := make(chan struct{}), make(chan struct{}) p := &cancellationPreparation{controlledPreparation: &controlledPreparation{closed: make(chan struct{})}} p.cancel = func(context.Context) error { return p.Close() } - p.start = func(ctx context.Context, _ string, _ proto.MessageInput, _ chan<- proto.Envelope) (agent.Session, error) { + p.start = func(ctx context.Context, _ string, _ proto.MessageInput, _ chan<- proto.Envelope) (fixtureSession, error) { close(entered) <-allowReturn return nil, ctx.Err() diff --git a/apps/daemon/internal/dispatch/preparation_cleanup_test.go b/apps/daemon/internal/dispatch/preparation_cleanup_test.go index b3aa40fcf..94227d7f7 100644 --- a/apps/daemon/internal/dispatch/preparation_cleanup_test.go +++ b/apps/daemon/internal/dispatch/preparation_cleanup_test.go @@ -9,7 +9,6 @@ import ( "testing" "time" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -37,7 +36,7 @@ func TestPreparedCancellationDoesNotAcknowledgeFailedCleanup(t *testing.T) { close(cancelled) return nil }}} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, _ chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, _ chan<- proto.Envelope) (fixtureSession, error) { close(entered) <-cancelled return nil, context.Canceled @@ -64,7 +63,7 @@ func TestShutdownRetriesFailedPreparedCancellationOnSameTarget(t *testing.T) { sender := &recSender{} var session *fakeSession p := &cancellationPreparation{controlledPreparation: &controlledPreparation{closed: make(chan struct{})}} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session = &fakeSession{out: out, closeOutOnCancel: true} return session, nil } @@ -165,7 +164,7 @@ func TestPublishedPreparedRunRetainsRetryAfterHandleRetirement(t *testing.T) { sender := &recSender{} session := &fakeSession{closeOutOnCancel: true} p := &cancellationPreparation{controlledPreparation: &controlledPreparation{closed: make(chan struct{})}} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out return session, nil } diff --git a/apps/daemon/internal/dispatch/preparation_executor_fixture_test.go b/apps/daemon/internal/dispatch/preparation_executor_fixture_test.go index 5e9e778d2..a8c952438 100644 --- a/apps/daemon/internal/dispatch/preparation_executor_fixture_test.go +++ b/apps/daemon/internal/dispatch/preparation_executor_fixture_test.go @@ -11,13 +11,23 @@ import ( // preparedFixture is a disposable fault-injection resource; Start transfers it // to one Session. A cancellablePreparation follows that Session across Start. type preparedFixture interface { - Start(context.Context, string, proto.MessageInput, chan<- proto.Envelope) (agent.Session, error) + Start(context.Context, string, proto.MessageInput, chan<- proto.Envelope) (fixtureSession, error) Close() error } type cancellablePreparation interface { preparedFixture - agent.Session + Cancel(context.Context) error + CancellationOutcome() proto.DonePayload +} + +// fixtureSession is the agent.Turn a fixture starts, without the settlement +// that preparationTurn supplies. +type fixtureSession interface { + Cancel(context.Context) error + CancellationOutcome() proto.DonePayload + SteerWithReceipt(context.Context, proto.PromptSteerPayload, func()) error + SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error } type preparationFactory func(context.Context, proto.PromptRequestPayload) (preparedFixture, error) @@ -77,11 +87,11 @@ func (e *preparationExecutor) StartTurn(ctx context.Context, id string, input pr // This fixture's old Start may have produced output before rejecting; retain its Turn. return &preparationTurn{owner: e, settled: terminal}, err } - return &preparationTurn{Session: session, owner: e, settled: terminal}, err + return &preparationTurn{fixtureSession: session, owner: e, settled: terminal}, err } type preparationTurn struct { - agent.Session + fixtureSession owner *preparationExecutor settled <-chan struct{} } @@ -101,24 +111,3 @@ func (t *preparationTurn) AwaitSettlement(ctx context.Context) (agent.TurnSettle return agent.TurnSettlement{}, ctx.Err() } } -func (t *preparationTurn) SubmitFunctionResult(ctx context.Context, p proto.FunctionResultPayload) error { - if target, ok := t.Session.(agent.FunctionResultSubmitter); ok { - return target.SubmitFunctionResult(ctx, p) - } - return agent.ErrUnknownFunctionCall -} -func (t *preparationTurn) Steer(ctx context.Context, p proto.PromptSteerPayload) error { - if target, ok := t.Session.(agent.Steerer); ok { - return target.Steer(ctx, p) - } - return agent.ErrSteeringInactive -} -func (t *preparationTurn) SteerWithReceipt(ctx context.Context, p proto.PromptSteerPayload, write func()) error { - if target, ok := t.Session.(agent.DurableSteerer); ok { - return target.SteerWithReceipt(ctx, p, write) - } - if target, ok := t.Session.(agent.Steerer); ok { - return target.Steer(ctx, p) - } - return agent.ErrSteeringInactive -} diff --git a/apps/daemon/internal/dispatch/preparation_test.go b/apps/daemon/internal/dispatch/preparation_test.go index 8c8b9d072..081ebd227 100644 --- a/apps/daemon/internal/dispatch/preparation_test.go +++ b/apps/daemon/internal/dispatch/preparation_test.go @@ -21,9 +21,9 @@ type controlledPreparation struct { closed chan struct{} once sync.Once mu sync.Mutex - session agent.Session + session fixtureSession starts atomic.Int32 - start func(context.Context, string, proto.MessageInput, chan<- proto.Envelope) (agent.Session, error) + start func(context.Context, string, proto.MessageInput, chan<- proto.Envelope) (fixtureSession, error) closeHook func() } @@ -38,7 +38,7 @@ func (p *controlledPreparation) Close() error { }) return nil } -func (p *controlledPreparation) Start(ctx context.Context, id string, prompt proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { +func (p *controlledPreparation) Start(ctx context.Context, id string, prompt proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { p.starts.Add(1) session, err := p.start(ctx, id, prompt, out) if session != nil { @@ -211,7 +211,7 @@ func TestPreparationSingleTransferAndReleaseDoesNotCancelRun(t *testing.T) { sender := &recSender{} gotSession := make(chan *fakeSession, 1) p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(ctx context.Context, id string, prompt proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(ctx context.Context, id string, prompt proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { if id != "real-run" || *prompt[0].Content[0].Text != "actual input" { t.Error("start identity or prompt changed") } @@ -265,7 +265,7 @@ func TestPreparationCancelDuringStartClosesLateSession(t *testing.T) { entered, cancelEntered, allowReturn := make(chan struct{}), make(chan struct{}), make(chan struct{}) lateSession := make(chan *fakeSession, 1) p := &cancellationPreparation{controlledPreparation: &controlledPreparation{closed: make(chan struct{})}} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { close(entered) <-allowReturn s := &fakeSession{out: out, closeOutOnCancel: true} diff --git a/apps/daemon/internal/dispatch/prepared_handoff_mutation_test.go b/apps/daemon/internal/dispatch/prepared_handoff_mutation_test.go index 0a04ea8e2..18a0acbbe 100644 --- a/apps/daemon/internal/dispatch/prepared_handoff_mutation_test.go +++ b/apps/daemon/internal/dispatch/prepared_handoff_mutation_test.go @@ -7,7 +7,6 @@ import ( "testing" "time" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -29,7 +28,7 @@ func (s *preparedMutationSession) SubmitFunctionResult(context.Context, proto.Fu return nil } -func (s *preparedMutationSession) Steer(context.Context, proto.PromptSteerPayload) error { +func (s *preparedMutationSession) SteerWithReceipt(context.Context, proto.PromptSteerPayload, func()) error { s.steers.Add(1) return nil } @@ -79,7 +78,7 @@ func TestPreparedHandoffReleaseWaitsForMutationReceipt(t *testing.T) { } session := &preparedMutationSession{fakeSession: &fakeSession{closeOutOnCancel: true}, cancelEntered: make(chan struct{})} p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out return session, nil } @@ -193,7 +192,7 @@ func TestPreparedHandoffRouterShutdownWaitsForReceiptAttempt(t *testing.T) { defer close(sender.release) session := &preparedMutationSession{fakeSession: &fakeSession{closeOutOnCancel: true}, cancelEntered: make(chan struct{})} p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out return session, nil } @@ -236,7 +235,7 @@ func TestPreparedHandoffEarlyDonePublishesAfterStarted(t *testing.T) { emitted, allowReturn := make(chan struct{}), make(chan struct{}) session := &fakeSession{closeOutOnCancel: true} p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(ctx context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(ctx context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out out <- mustEnv(t, proto.TypeDone, "run", proto.DonePayload{Content: "complete"}) close(emitted) @@ -270,6 +269,3 @@ func TestPreparedHandoffEarlyDonePublishesAfterStarted(t *testing.T) { t.Fatalf("started/Done order invalid: started=%d done=%d cancels=%d", started, done, session.cancels()) } } - -var _ agent.FunctionResultSubmitter = (*preparedMutationSession)(nil) -var _ agent.Steerer = (*preparedMutationSession)(nil) diff --git a/apps/daemon/internal/dispatch/prepared_handoff_test.go b/apps/daemon/internal/dispatch/prepared_handoff_test.go index 3a76fbac6..efb9ccb3d 100644 --- a/apps/daemon/internal/dispatch/prepared_handoff_test.go +++ b/apps/daemon/internal/dispatch/prepared_handoff_test.go @@ -8,7 +8,6 @@ import ( "testing" "time" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) @@ -19,7 +18,7 @@ func TestPreparedHandoffDrainsBurstBeforeStartReturns(t *testing.T) { sent, allowReturn := make(chan struct{}), make(chan struct{}) session := &fakeSession{closeOutOnCancel: true} p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(ctx context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(ctx context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out for sequence := uint64(1); sequence <= preparedBurstFrames; sequence++ { select { @@ -113,7 +112,7 @@ func (s *blockingStartingSender) Send(ctx context.Context, env proto.Envelope) e func TestPreparedHandoffAbortBeforeStartAdmissionSkipsNativeStart(t *testing.T) { sender := &blockingStartingSender{recSender: &recSender{}, entered: make(chan struct{}), release: make(chan struct{})} p := &cancellationPreparation{controlledPreparation: &controlledPreparation{closed: make(chan struct{})}} - p.start = func(context.Context, string, proto.MessageInput, chan<- proto.Envelope) (agent.Session, error) { + p.start = func(context.Context, string, proto.MessageInput, chan<- proto.Envelope) (fixtureSession, error) { t.Fatal("abort that won admission called native Start") return nil, errors.New("unexpected Start") } @@ -159,7 +158,7 @@ func TestPreparedHandoffDuplicateStartDoesNotReexecuteDuringPublication(t *testi sender := &blockingStartedSender{recSender: &recSender{}, entered: make(chan struct{}), release: make(chan struct{})} session := &preparedMutationSession{fakeSession: &fakeSession{closeOutOnCancel: true}, cancelEntered: make(chan struct{})} p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out return session, nil } @@ -210,7 +209,7 @@ func TestPreparedHandoffUnsupportedFunctionReleasesOperationBarrier(t *testing.T sender := &recSender{} session := &fakeSession{closeOutOnCancel: true} p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out return session, nil } @@ -235,7 +234,7 @@ func TestPreparedHandoffEarlyDoneStillAllowsExplicitAbort(t *testing.T) { var once sync.Once unblock := func() { once.Do(func() { close(cancelled) }) } p := &cancellationPreparation{controlledPreparation: &controlledPreparation{closed: make(chan struct{})}} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { out <- mustEnv(t, proto.TypeDone, "run", proto.DonePayload{}) <-cancelled return nil, context.Canceled @@ -261,7 +260,7 @@ func TestPreparedHandoffExpiresDuringStartedPublication(t *testing.T) { sender := &blockingStartedSender{recSender: &recSender{}, entered: make(chan struct{}), release: make(chan struct{})} session := &fakeSession{closeOutOnCancel: true} p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out return session, nil } @@ -293,7 +292,7 @@ func TestPreparedHandoffEarlyDoneDetachesPublishedPreparation(t *testing.T) { defer unblock() session := &fakeSession{closeOutOnCancel: true} p := &cancellationPreparation{controlledPreparation: &controlledPreparation{closed: make(chan struct{})}} - p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(_ context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { session.out = out out <- mustEnv(t, proto.TypeDone, "run", proto.DonePayload{}) <-allowReturn diff --git a/apps/daemon/internal/dispatch/receipt_order_test.go b/apps/daemon/internal/dispatch/receipt_order_test.go index f3eaeeb17..4ade9ce34 100644 --- a/apps/daemon/internal/dispatch/receipt_order_test.go +++ b/apps/daemon/internal/dispatch/receipt_order_test.go @@ -50,9 +50,9 @@ func TestDurableCompletionWaitsForSteeringReceiptSend(t *testing.T) { registry := agent.NewRegistry() var session *fakeSession var calls atomic.Int32 - registerSession(registry, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(registry, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { session = &fakeSession{out: out, closeOutOnCancel: true} - return &steeringSession{fakeSession: session, steer: func(context.Context, proto.PromptSteerPayload) error { + return &steeringSession{fakeSession: session, steer: func(context.Context, proto.PromptSteerPayload, func()) error { calls.Add(1) if mode == "unknown" { return errors.New("native outcome unknown") @@ -125,9 +125,9 @@ func TestShutdownReleasesSteeringWorkerAndReceiptJoin(t *testing.T) { registry := agent.NewRegistry() entered, exited := make(chan struct{}), make(chan struct{}) var session *fakeSession - registerSession(registry, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(registry, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { session = &fakeSession{out: out, closeOutOnCancel: true} - return &steeringSession{fakeSession: session, steer: func(ctx context.Context, _ proto.PromptSteerPayload) error { + return &steeringSession{fakeSession: session, steer: func(ctx context.Context, _ proto.PromptSteerPayload, _ func()) error { close(entered) defer close(exited) if phase == "native" { diff --git a/apps/daemon/internal/dispatch/receipt_shutdown_test.go b/apps/daemon/internal/dispatch/receipt_shutdown_test.go index f553d7581..dec26022a 100644 --- a/apps/daemon/internal/dispatch/receipt_shutdown_test.go +++ b/apps/daemon/internal/dispatch/receipt_shutdown_test.go @@ -45,9 +45,9 @@ func TestShutdownCancelsCompletionErrorSend(t *testing.T) { sender := &shutdownAllSendsBlockSender{entered: make(chan struct{}), terminal: make(chan context.Context, 1), rescue: make(chan struct{})} registry := agent.NewRegistry() var session *fakeSession - registerSession(registry, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(registry, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { session = &fakeSession{out: out, closeOutOnCancel: true} - return &steeringSession{fakeSession: session, steer: func(context.Context, proto.PromptSteerPayload) error { return nil }}, nil + return &steeringSession{fakeSession: session, steer: func(context.Context, proto.PromptSteerPayload, func()) error { return nil }}, nil }) router, err := dispatch.New(dispatch.Config{Registry: registry, Sender: sender}) if err != nil { diff --git a/apps/daemon/internal/dispatch/router_test.go b/apps/daemon/internal/dispatch/router_test.go index 5b19337a9..966c6d7b9 100644 --- a/apps/daemon/internal/dispatch/router_test.go +++ b/apps/daemon/internal/dispatch/router_test.go @@ -72,6 +72,15 @@ type fakeSession struct { func (s *fakeSession) CancellationOutcome() proto.DonePayload { return proto.DonePayload{} } +// fakeSession takes no active input and has no outstanding function calls. +func (s *fakeSession) SteerWithReceipt(context.Context, proto.PromptSteerPayload, func()) error { + return agent.ErrSteeringInactive +} + +func (s *fakeSession) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { + return agent.ErrUnknownFunctionCall +} + func (s *fakeSession) Cancel(context.Context) error { s.cancelMu.Lock() s.cancelCalls++ @@ -114,7 +123,7 @@ func newHarness(t *testing.T) *harness { gotReq: make(chan proto.PromptRequestPayload, 16), gotSess: make(chan *fakeSession, 16), } - registerSession(h.reg, proto.SupportedAgentKind{Kind: "fake_alpha", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(h.reg, proto.SupportedAgentKind{Kind: "fake_alpha", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { sess := &fakeSession{out: out, ctx: ctx, closeOutOnCancel: true} h.gotReq <- req h.gotSess <- sess @@ -130,11 +139,11 @@ func newHarness(t *testing.T) *harness { // registerSession declares info, without an Environment, and starts each // Turn of the kind with factory. -func registerSession(reg *agent.Registry, info proto.SupportedAgentKind, factory func(context.Context, proto.PromptRequestPayload, chan<- proto.Envelope) (agent.Session, error)) { +func registerSession(reg *agent.Registry, info proto.SupportedAgentKind, factory func(context.Context, proto.PromptRequestPayload, chan<- proto.Envelope) (fixtureSession, error)) { info.Capabilities.EnvironmentNone = proto.CapabilitySupported reg.RegisterKind(info, prototest.ModelConfiguration()) reg.RegisterExecutor(info.Kind, preparationExecutorFixture(func(_ context.Context, req proto.PromptRequestPayload) (preparedFixture, error) { - return &controlledPreparation{start: func(ctx context.Context, id string, input proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + return &controlledPreparation{start: func(ctx context.Context, id string, input proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { req.RunID, req.Input = id, input return factory(ctx, req, out) }}, nil diff --git a/apps/daemon/internal/dispatch/steering.go b/apps/daemon/internal/dispatch/steering.go index a60d6b8c7..d097ce4e1 100644 --- a/apps/daemon/internal/dispatch/steering.go +++ b/apps/daemon/internal/dispatch/steering.go @@ -23,7 +23,6 @@ const ( type steeringReceipt struct { fingerprint [32]byte ack proto.PromptSteerAckPayload - durable bool } func (r *Router) handlePromptSteer(ctx context.Context, env proto.Envelope) error { @@ -70,7 +69,7 @@ func (r *Router) queueSteering(ctx context.Context, env proto.Envelope, input pr encoded, _ := json.Marshal(input.Input) fingerprint := sha256.Sum256(encoded) if previous, ok := state.steering[input.InputID]; ok { - if previous.fingerprint != fingerprint || previous.durable != input.DurableReceipt { + if previous.fingerprint != fingerprint { ack.ErrorCode, ack.Error = "input_conflict", "This input ID was already used with different text." return &ack } @@ -91,19 +90,13 @@ func (r *Router) queueSteering(ctx context.Context, env proto.Envelope, input pr ack.ErrorCode, ack.Error = "not_ready", "The run is still starting." // Bind input identity before any retryable state so changed text cannot // slip through a startup or in-flight retry. - state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack, durable: input.DurableReceipt} + state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack} if state.session == nil { return &ack } - session := state.session - steerer, supportsSteering := session.(agent.Steerer) - if !input.DurableReceipt && !supportsSteering { - ack.ErrorCode, ack.Error = "unsupported", "This engine does not support active-turn input." - return &ack - } if state.steerBusy { ack.ErrorCode, ack.Error = "busy", "Another input is awaiting an engine receipt; retry this input later." - state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack, durable: input.DurableReceipt} + state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack} return &ack } session, finishOperation, ready := r.preparedOperationLocked(state) @@ -112,7 +105,7 @@ func (r *Router) queueSteering(ctx context.Context, env proto.Envelope, input pr return &ack } ack.ErrorCode, ack.Error = "in_flight", "This input is awaiting an engine receipt." - state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack, durable: input.DurableReceipt} + state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack} state.steerBusy = true r.shutdownWG.Add(1) go func() { @@ -125,14 +118,7 @@ func (r *Router) queueSteering(ctx context.Context, env proto.Envelope, input pr }() ctx, stop := r.shutdownContext(ctx) defer stop() - var err error - if input.DurableReceipt { - err = r.steerDurably(ctx, state, session, env, input, fingerprint) - } else { - callCtx, cancel := context.WithTimeout(ctx, steeringCallTimeout) - err = steerer.Steer(callCtx, input) - cancel() - } + err := r.steerDurably(ctx, state, session, env, input, fingerprint) r.publishSteeringReceipt(ctx, state, env, input, fingerprint, steeringResult(input.InputID, err)) }() return nil diff --git a/apps/daemon/internal/dispatch/steering_lifetime.go b/apps/daemon/internal/dispatch/steering_lifetime.go index adbfbe054..97ba37ebe 100644 --- a/apps/daemon/internal/dispatch/steering_lifetime.go +++ b/apps/daemon/internal/dispatch/steering_lifetime.go @@ -9,7 +9,7 @@ import ( "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) -func (r *Router) steerDurably(sendCtx context.Context, state *sessionState, session agent.DurableSteerer, env proto.Envelope, input proto.PromptSteerPayload, fingerprint [32]byte) error { +func (r *Router) steerDurably(sendCtx context.Context, state *sessionState, session agent.Turn, env proto.Envelope, input proto.PromptSteerPayload, fingerprint [32]byte) error { ctx, cancel := context.WithCancel(state.ctx) defer cancel() stopShutdown := context.AfterFunc(sendCtx, cancel) @@ -32,7 +32,7 @@ func (r *Router) steerDurably(sendCtx context.Context, state *sessionState, sess func (r *Router) publishSteeringReceipt(ctx context.Context, state *sessionState, env proto.Envelope, input proto.PromptSteerPayload, fingerprint [32]byte, ack proto.PromptSteerAckPayload) bool { r.mu.Lock() - state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack, durable: input.DurableReceipt} + state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack} r.mu.Unlock() sendCtx, cancel := context.WithTimeout(ctx, steeringSendTimeout) defer cancel() diff --git a/apps/daemon/internal/dispatch/steering_lifetime_test.go b/apps/daemon/internal/dispatch/steering_lifetime_test.go index c30af7c6a..067d030a3 100644 --- a/apps/daemon/internal/dispatch/steering_lifetime_test.go +++ b/apps/daemon/internal/dispatch/steering_lifetime_test.go @@ -6,29 +6,19 @@ import ( "testing" "time" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" ) -type durableSteeringSession struct { - *steeringSession - phased func(context.Context, proto.PromptSteerPayload, func()) error -} - -func (s *durableSteeringSession) SteerWithReceipt(ctx context.Context, input proto.PromptSteerPayload, written func()) error { - return s.phased(ctx, input, written) -} - func TestDurableSteeringWaitsBeyondTransportDeadline(t *testing.T) { h := newHarness(t) defer h.router.Shutdown(context.Background()) var session *fakeSession var calls atomic.Int32 release := make(chan struct{}) - registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { session = &fakeSession{out: out, closeOutOnCancel: true} - return &durableSteeringSession{steeringSession: &steeringSession{fakeSession: session}, phased: func(ctx context.Context, input proto.PromptSteerPayload, written func()) error { + return &steeringSession{fakeSession: session, steer: func(ctx context.Context, input proto.PromptSteerPayload, written func()) error { calls.Add(1) written() select { @@ -40,7 +30,7 @@ func TestDurableSteeringWaitsBeyondTransportDeadline(t *testing.T) { }}, nil }) startRun(t, h.router, h.sender, "codex", "durable") - input := proto.PromptSteerPayload{InputID: "extra", Input: proto.TextInput("additional"), DurableReceipt: true} + input := proto.PromptSteerPayload{InputID: "extra", Input: proto.TextInput("additional")} env := scoped(t, "durable", proto.TypePromptSteer, "durable", input) if err := handleSteeringAndWait(t, h, env); err != nil { t.Fatal(err) @@ -77,8 +67,8 @@ func TestDurableSteeringTransportTimeoutAndShutdown(t *testing.T) { h := newHarness(t) defer h.router.Shutdown(context.Background()) exited := make(chan struct{}) - registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { - return &durableSteeringSession{steeringSession: &steeringSession{fakeSession: &fakeSession{out: out, closeOutOnCancel: true}}, phased: func(ctx context.Context, _ proto.PromptSteerPayload, written func()) error { + registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { + return &steeringSession{fakeSession: &fakeSession{out: out, closeOutOnCancel: true}, steer: func(ctx context.Context, _ proto.PromptSteerPayload, written func()) error { defer close(exited) if phase == "written" { written() @@ -89,7 +79,7 @@ func TestDurableSteeringTransportTimeoutAndShutdown(t *testing.T) { }) ctx := context.Background() startRun(t, h.router, h.sender, "codex", "run") - if err := h.router.Handle(ctx, scoped(t, "run", proto.TypePromptSteer, "run", proto.PromptSteerPayload{InputID: "one", Input: proto.TextInput("text"), DurableReceipt: true})); err != nil { + if err := h.router.Handle(ctx, scoped(t, "run", proto.TypePromptSteer, "run", proto.PromptSteerPayload{InputID: "one", Input: proto.TextInput("text")})); err != nil { t.Fatal(err) } if phase == "blocked-write" { diff --git a/apps/daemon/internal/dispatch/steering_test.go b/apps/daemon/internal/dispatch/steering_test.go index 488438abc..d1eb361de 100644 --- a/apps/daemon/internal/dispatch/steering_test.go +++ b/apps/daemon/internal/dispatch/steering_test.go @@ -14,11 +14,11 @@ import ( type steeringSession struct { *fakeSession - steer func(context.Context, proto.PromptSteerPayload) error + steer func(context.Context, proto.PromptSteerPayload, func()) error } -func (s *steeringSession) Steer(ctx context.Context, input proto.PromptSteerPayload) error { - return s.steer(ctx, input) +func (s *steeringSession) SteerWithReceipt(ctx context.Context, input proto.PromptSteerPayload, written func()) error { + return s.steer(ctx, input, written) } func TestSteeringReceiptsAndRetries(t *testing.T) { @@ -31,15 +31,12 @@ func TestSteeringReceiptsAndRetries(t *testing.T) { h := newHarness(t) defer h.router.Shutdown(context.Background()) calls, starts := 0, 0 - registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { starts++ return &steeringSession{ fakeSession: &fakeSession{out: out, closeOutOnCancel: true}, - steer: func(ctx context.Context, input proto.PromptSteerPayload) error { + steer: func(_ context.Context, input proto.PromptSteerPayload, _ func()) error { calls++ - if _, ok := ctx.Deadline(); !ok { - t.Error("native request has no deadline") - } if input.InputID != "input-1" || *input.Input[0].Content[0].Text != "additional text" { t.Errorf("input lost: %+v", input) } @@ -90,10 +87,10 @@ func TestSteeringReadinessAndInputValidation(t *testing.T) { h := newHarness(t) defer h.router.Shutdown(context.Background()) calls := 0 - registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { return &steeringSession{ fakeSession: &fakeSession{out: out, closeOutOnCancel: true}, - steer: func(context.Context, proto.PromptSteerPayload) error { + steer: func(context.Context, proto.PromptSteerPayload, func()) error { calls++ if calls == 1 { return agent.ErrSteeringNotReady @@ -156,10 +153,10 @@ func TestSteeringDoesNotBlockOtherRunCancellation(t *testing.T) { defer h.router.Shutdown(context.Background()) entered, release := make(chan struct{}), make(chan struct{}) defer close(release) - registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { return &steeringSession{ fakeSession: &fakeSession{out: out, closeOutOnCancel: true}, - steer: func(context.Context, proto.PromptSteerPayload) error { + steer: func(context.Context, proto.PromptSteerPayload, func()) error { close(entered) <-release return agent.ErrSteeringRejected @@ -214,10 +211,10 @@ func TestSteeringCapacityPreservesExistingReceipts(t *testing.T) { h := newHarness(t) defer h.router.Shutdown(context.Background()) calls := 0 - registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) { + registerSession(h.reg, proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (fixtureSession, error) { return &steeringSession{ fakeSession: &fakeSession{out: out, closeOutOnCancel: true}, - steer: func(context.Context, proto.PromptSteerPayload) error { + steer: func(context.Context, proto.PromptSteerPayload, func()) error { calls++ return nil }, diff --git a/apps/daemon/internal/dispatch/workspace_directory_test.go b/apps/daemon/internal/dispatch/workspace_directory_test.go index bf09335c5..68b006490 100644 --- a/apps/daemon/internal/dispatch/workspace_directory_test.go +++ b/apps/daemon/internal/dispatch/workspace_directory_test.go @@ -8,14 +8,13 @@ import ( "testing" "time" - "github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent" "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" ) func TestWorkspaceDirectoryRetainsEnvironmentAndTransferredOwner(t *testing.T) { sender := &recSender{} p := &controlledPreparation{closed: make(chan struct{})} - p.start = func(ctx context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (agent.Session, error) { + p.start = func(ctx context.Context, _ string, _ proto.MessageInput, out chan<- proto.Envelope) (fixtureSession, error) { return &fakeSession{out: out, ctx: ctx, closeOutOnCancel: true}, nil } r := preparationRouter(t, sender, time.Minute, func(context.Context, proto.PromptRequestPayload) (preparedFixture, error) { return p, nil }) diff --git a/apps/daemon/internal/wireconformance/wire_test.go b/apps/daemon/internal/wireconformance/wire_test.go index 6e9bb9394..cfd2ff38e 100644 --- a/apps/daemon/internal/wireconformance/wire_test.go +++ b/apps/daemon/internal/wireconformance/wire_test.go @@ -361,7 +361,10 @@ func (turn *controlledTurn) AwaitSettlement(ctx context.Context) (agent.TurnSett } } -// These fixtures exercise settlement only; active input is deliberately rejected. +// These fixtures exercise settlement only; active input and results are deliberately rejected. func (*controlledTurn) SteerWithReceipt(context.Context, proto.PromptSteerPayload, func()) error { return agent.ErrSteeringRejected } +func (*controlledTurn) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { + return agent.ErrUnknownFunctionCall +} diff --git a/apps/daemon/testdata/onboarding/main.go b/apps/daemon/testdata/onboarding/main.go index 6c19b5e33..1edfede2a 100644 --- a/apps/daemon/testdata/onboarding/main.go +++ b/apps/daemon/testdata/onboarding/main.go @@ -127,9 +127,6 @@ func (s *session) Cancel(context.Context) error { func (s *session) CancellationOutcome() proto.DonePayload { return proto.DonePayload{Content: "cancelled", Metadata: map[string]any{proto.DoneMetaAgentSessionID: s.native}} } -func (s *session) Steer(ctx context.Context, p proto.PromptSteerPayload) error { - return s.SteerWithReceipt(ctx, p, func() {}) -} func (s *session) SteerWithReceipt(_ context.Context, p proto.PromptSteerPayload, written func()) error { text, err := p.Input.TextOnly() if err != nil { @@ -150,6 +147,10 @@ func (s *session) SteerWithReceipt(_ context.Context, p proto.PromptSteerPayload return nil } +func (s *session) SubmitFunctionResult(context.Context, proto.FunctionResultPayload) error { + return fmt.Errorf("%w: the fixture has no functions", agent.ErrUnsupportedOperation) +} + func (s *session) AwaitSettlement(ctx context.Context) (agent.TurnSettlement, error) { select { case <-s.settled: diff --git a/contracts/agents-api/harness-onboarding.md b/contracts/agents-api/harness-onboarding.md index 96170206a..522782de9 100644 --- a/contracts/agents-api/harness-onboarding.md +++ b/contracts/agents-api/harness-onboarding.md @@ -58,20 +58,19 @@ Implement the mandatory text lifecycle and handle every extension explicitly. Qu ## Required adapter interfaces -[`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) is the interface entry point. The required lifecycle is `ExecutorFactory`, `Executor`, `Turn` (including `DurableSteerer`) and `TurnSettlement`. Required methods perform their native obligations; returning Unsupported is not an implementation of cancellation, receipts, settlement or cleanup. Turn extension interfaces stay small and separate, but every public adapter implements each one explicitly. All use the neutral protocol types. +[`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) is the interface entry point. The required lifecycle is `ExecutorFactory`, `Executor`, `Turn` and `TurnSettlement`. `Turn` is one interface: `Cancel`, `CancellationOutcome`, `SteerWithReceipt`, `SubmitFunctionResult` and `AwaitSettlement`. Required methods perform their native obligations; returning Unsupported is not an implementation of cancellation, receipts, settlement or cleanup. An operation the adapter does not support returns Unsupported, and the capability declaration, not the method, decides whether the Runtime calls it. All use the neutral protocol types. -For example, the Codex adapter keeps its app-server and thread, the Claude adapter one streaming Query, and the MiniMax adapter its ACP connection and native session. All expose the same Executor and Turn contract. Native callbacks and resources stay inside the adapter; the Runtime owns admission, idle expiry and replacement. Cancellation targets the exact Turn through `agent.Session`, and the adapter supplies native completion evidence to the Runtime. +For example, the Codex adapter keeps its app-server and thread, the Claude adapter one streaming Query, and the MiniMax adapter its ACP connection and native session. All expose the same Executor and Turn contract. Native callbacks and resources stay inside the adapter; the Runtime owns admission, idle expiry and replacement. Cancellation targets the exact Turn through `Turn.Cancel`, and the adapter supplies native completion evidence to the Runtime. | Interface or contract | Required handling | Obligation | | --- | --- | --- | | `ExecutorFactory`, `Executor.StartTurn`, `Executor.Close` | Real implementation | Prepare without model input; keep ownership of failed or uncertain resources; confirm cleanup | -| `Session`, `Turn`, `CancellationOutcome`, `AwaitSettlement` | Real implementation | Cancel the exact Turn, keep observed results and confirm settlement independently of cancellation requests | -| `DurableSteerer` | Real implementation on every Turn | Distinguish a complete write from the native application receipt; keep retry identity | -| `Steerer` | Explicit implementation or Unsupported | Additional non-durable active-Turn input | -| `FunctionResultSubmitter` | Explicit implementation or Unsupported | Match native call and result identity and acknowledge application | +| `Turn.Cancel`, `CancellationOutcome`, `AwaitSettlement` | Real implementation | Cancel the exact Turn, keep observed results and confirm settlement independently of cancellation requests | +| `Turn.SteerWithReceipt` | Real implementation | Distinguish a complete write from the native application receipt; keep retry identity | +| `Turn.SubmitFunctionResult` | Real implementation or Unsupported | Match native call and result identity and acknowledge application | | Neutral messages, images, MCP, structured output and Subagent observations | Explicit capability decisions | Keep each operation's protocol semantics; reject unsupported input before submission | -Each adapter's `contracts.go` holds an individual compile-time assertion for each small interface. Do not embed a default implementation that makes future interfaces appear implemented. Adding a contract also requires a classification in the common completeness check and an explicit assertion in every public adapter; the check follows the authored Harness catalog. +Each adapter's `contracts.go` asserts at compile time that it implements `agent.Executor` and `agent.Turn`. Do not embed a default implementation that makes a new method appear implemented. The common completeness check follows the authored Harness catalog and rejects any other exported interface in `agent`. For a design-level refusal, implement the method directly: @@ -103,7 +102,7 @@ A Session owns one reusable Executor in its connected Runtime; a Turn owns one i **Binding.** The Runtime binds its Executor record to the Session, Environment, connection and immutable execution configuration. Resume identity and prior-Turn recovery flags are continuity assertions, not configuration changes. A supplied native identity must match the retained owner, and when existing history is required, recovery never starts a new root. A configuration conflict is an error, not a hot switch. A lost connection retires its owners and handles; old timers, output and cancellation cannot affect their replacements. -**Per-Turn state.** Each Turn gets a fresh wrapper, output channel and receipt state. Steering and function interfaces belong to that Turn. Native callbacks capture the originating Turn before asynchronous work, so a late event is never attributed to whichever Turn is active. Native processes, query or transport connections, fixed capability configuration and native session identity belong to the Executor. Do not reset completed `sync.Once` values or reuse an old Turn object. +**Per-Turn state.** Each Turn gets a fresh wrapper, output channel and receipt state. Steering and function results belong to that Turn. Native callbacks capture the originating Turn before asynchronous work, so a late event is never attributed to whichever Turn is active. Native processes, query or transport connections, fixed capability configuration and native session identity belong to the Executor. Do not reset completed `sync.Once` values or reuse an old Turn object. **Start.** A nil Turn from `StartTurn` guarantees that no native input was submitted and the output channel was not retained; the Runtime then closes the channel. Once input may have been submitted, return a non-nil Turn even with an error: that Turn owns exactly-once output closure and stays tracked until settlement. Unknown input is never replayed. A definite `executor_unavailable` Start rejection allows one common recovery attempt, only after the previous Executor has been closed and no input was submitted; the Runtime rechecks the same physical peer and the current authorization. @@ -113,8 +112,8 @@ A Session owns one reusable Executor in its connected Runtime; a Turn owns one i - An error means settlement is unconfirmed and frees neither ownership nor capacity. Caller deadlines stop the wait, not the tracked cleanup. Retry the same cleanup target serially; a failed cleanup blocks replacement and keeps its resource slot. - `Executor.Close` confirms resource retirement independently of the Turn outcome: an immutable Turn error must not prevent closing the native transport once its work and output have stopped. - Include owned background work in settlement and keep the exact native cleanup target after a failure. Native termination belongs to the adapter; a bulk cleanup acknowledgement alone does not establish quiescence. -- Every `Session` declares `CancellationOutcome`. `Turn` inherits it. The snapshot keeps observed native identity, Usage and output and remains readable after cancellation. Missing evidence stays unset; an empty `DonePayload` means nothing has been observed, not that cancellation succeeded or is unsupported. Reading the snapshot does not wait for settlement. -- `Session.Cancel` requests cancellation; output closure signals teardown. Turn settlement still requires `AwaitSettlement` and any required `Executor.Close`; neither a successful cancellation request nor its snapshot replaces those waits. +- Every Turn implements `CancellationOutcome`. The snapshot keeps observed native identity, Usage and output and remains readable after cancellation. Missing evidence stays unset; an empty `DonePayload` means nothing has been observed, not that cancellation succeeded or is unsupported. Reading the snapshot does not wait for settlement. +- `Turn.Cancel` requests cancellation; output closure signals teardown. Turn settlement still requires `AwaitSettlement` and any required `Executor.Close`; neither a successful cancellation request nor its snapshot replaces those waits. **What the Runtime does around a Turn.** One output consumer starts before native Start, drains the bounded 64-frame channel and keeps the terminal observation until Start publication, Turn settlement and admitted operation receipts finish. Natural completion never calls Cancel. Input and function admission close before settlement; operations already admitted hold their barrier through native receipts and outbound acknowledgement. The Runtime sends cancellation to the Turn before waiting on that barrier, because a written input may need a native interrupt to produce its receipt. It joins native settlement, any required confirmed Executor close, output drain and all admitted operations before an applied acknowledgement or reuse, and only then forwards Done or an applied cancellation receipt. A failed Close can report failure while keeping the same Run and outstanding operations for retry; a closed caller wait cannot manufacture an applied input receipt. The Runtime commits native continuity and releases the old Run's admission before publishing Done, since the receiver may start another Turn at once; a late terminal-send failure belongs to the old Run and cannot invalidate a successor that already owns the Executor. Connection shutdown owns transport-loss cleanup. The settlement wait is ten seconds and the receipt send budget five seconds; a timeout is not proof of quiescence. @@ -160,7 +159,7 @@ Registration is static and requires a build. Export one `agent.Declaration` from Every `proto.AgentKindCapabilities` field must be explicitly `proto.CapabilitySupported` or `proto.CapabilityUnsupported`, even for an unavailable Harness. `proto.CapabilityUnspecified` is invalid: zero values and omitted fields never mean Unsupported. An installation probe may set an individual field with `proto.CapabilityFromBool`; it must not populate unmentioned or future fields. Availability stays separate in `SupportedAgentKind.Available`. Registration validates the complete declaration before changing the registry, and the wire carries an explicit boolean for every field, so omitted and null fields are invalid. A new field requires a decision in every production declaration. Runtime consumers use `IsSupported()` and reject unsupported requests before native operations; an interface assertion verifies implementation, never support. Every declaration must match the behavior verified for that installation; the [Core–Runtime protocol](../../docs/runtime-protocol.md#capability-declarations) owns how declarations travel and are frozen. -Every available Harness implements, without a declaration, the shared Turn lifecycle, including `DurableSteerer` input and the Turn settlement contract that `contracttest.TextLifecycle` checks, typed `execution_controls` and tool observations. A Harness that cannot meet them on a platform reports `Available` false there. `WorkspaceReadPreparation` admits `execution_prepare` with `workspace_read_only`, which the Environment owner readies and serves without calling the Executor factory. Runtime registration does not grant Core qualification; the service profile does. +Every available Harness implements, without a declaration, the shared Turn lifecycle, including durable `SteerWithReceipt` input and the Turn settlement contract that `contracttest.TextLifecycle` checks, typed `execution_controls` and tool observations. A Harness that cannot meet them on a platform reports `Available` false there. `FunctionTools` admits `SubmitFunctionResult`. `WorkspaceReadPreparation` admits `execution_prepare` with `workspace_read_only`, which the Environment owner readies and serves without calling the Executor factory. Runtime registration does not grant Core qualification; the service profile does. The runnable test-only example [`testdata/onboarding/main.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/testdata/onboarding/main.go) registers a text-only synthetic Harness under the `mcode` kind, because Core admits only [catalog](./harness-catalog.md) Harnesses. It shows a Session-owned Executor, fresh Turns, durable steering, cancellation and history binding, and is never shipped. diff --git a/contracts/agents-api/zh/harness-onboarding.md b/contracts/agents-api/zh/harness-onboarding.md index 608b3c734..013461899 100644 --- a/contracts/agents-api/zh/harness-onboarding.md +++ b/contracts/agents-api/zh/harness-onboarding.md @@ -1,7 +1,7 @@ --- title: "将原生 Harness 添加到 OpenAgentCore" source: contracts/agents-api/harness-onboarding.md -source_hash: 749d48c7bcdf5dfbf2f01d21b4524ff8c7fd45943252d0ce5f2ca5ba9efae8eb +source_hash: f7c2ab839807374439ed19d9ba5a433e8756261fd7d74705aa3dcc2a4c4f8423 --- **Harness** 是一种运行模型和工具循环的原生代理引擎(Codex、Claude Code、MiniMax Code)。**Harness 适配器**将 Runtime 的 Executor 和 Turn 契约转换到该引擎的 SDK 或协议。本文档定义 Runtime–Harness 协议:适配器接口及其生命周期义务、注册、Core 资格认定和验收。[Harness capabilities](harness-capabilities.md) 记录了当前每个 Harness 支持的功能。 @@ -60,20 +60,19 @@ Environment 提供执行资源。受管 E2B、Docker 和 microsandbox 机器以 ## 必需的适配器接口 {#required-adapter-interfaces} -[`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) 是接口入口。必需的生命周期包括 `ExecutorFactory`、`Executor`、`Turn`(包括 `DurableSteerer`)和 `TurnSettlement`。必需方法必须履行其原生义务;返回 Unsupported 并不构成对取消、回执、结算或清理的实现。Turn 扩展接口应保持小而独立,但每个公共适配器都必须明确实现每一个接口。所有接口都使用中立协议类型。 +[`agent/harness.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/internal/agent/harness.go) 是接口入口。必需的生命周期包括 `ExecutorFactory`、`Executor`、`Turn` 和 `TurnSettlement`。`Turn` 是一个接口:`Cancel`、`CancellationOutcome`、`SteerWithReceipt`、`SubmitFunctionResult` 和 `AwaitSettlement`。必需方法必须履行其原生义务;返回 Unsupported 并不构成对取消、回执、结算或清理的实现。适配器不支持的操作返回 Unsupported,由能力声明而不是方法决定 Runtime 是否调用它。所有接口都使用中立协议类型。 -例如,Codex 适配器保留其 app-server 和 thread,Claude 适配器保留一个流式 Query,MiniMax 适配器保留其 ACP 连接和原生 session。它们都公开相同的 Executor 和 Turn 契约。原生回调和资源保留在适配器内部;Runtime 负责准入、空闲过期和替换。取消通过 `agent.Session` 精确定位到目标 Turn,适配器则向 Runtime 提供原生完成证据。 +例如,Codex 适配器保留其 app-server 和 thread,Claude 适配器保留一个流式 Query,MiniMax 适配器保留其 ACP 连接和原生 session。它们都公开相同的 Executor 和 Turn 契约。原生回调和资源保留在适配器内部;Runtime 负责准入、空闲过期和替换。取消通过 `Turn.Cancel` 精确定位到目标 Turn,适配器则向 Runtime 提供原生完成证据。 | 接口或契约 | 必需处理 | 义务 | | --- | --- | --- | | `ExecutorFactory`、`Executor.StartTurn`、`Executor.Close` | 真实实现 | 在没有模型输入的情况下准备;保留失败或不确定资源的所有权;确认清理 | -| `Session`、`Turn`、`CancellationOutcome`、`AwaitSettlement` | 真实实现 | 取消精确的 Turn,保留已观察结果,并独立于取消请求确认结算 | -| `DurableSteerer` | 每个 Turn 上真实实现 | 区分完整写入与原生应用回执;保留重试身份 | -| `Steerer` | 明确实现或 Unsupported | 额外的非持久化活动 Turn 输入 | -| `FunctionResultSubmitter` | 明确实现或 Unsupported | 匹配原生调用和结果身份,并确认应用 | +| `Turn.Cancel`、`CancellationOutcome`、`AwaitSettlement` | 真实实现 | 取消精确的 Turn,保留已观察结果,并独立于取消请求确认结算 | +| `Turn.SteerWithReceipt` | 真实实现 | 区分完整写入与原生应用回执;保留重试身份 | +| `Turn.SubmitFunctionResult` | 真实实现或 Unsupported | 匹配原生调用和结果身份,并确认应用 | | 中立消息、图像、MCP、结构化输出和 Subagent 观察 | 明确作出能力决策 | 保持每项操作的协议语义;在提交前拒绝不受支持的输入 | -每个适配器的 `contracts.go` 都包含针对每个小型接口的单项编译时断言。不要嵌入会让未来接口看起来已经实现的默认实现。添加契约时,还必须在通用完整性检查中进行分类,并在每个公共适配器中添加明确断言;该检查遵循已编写的 Harness 目录。 +每个适配器的 `contracts.go` 在编译时断言其实现了 `agent.Executor` 和 `agent.Turn`。不要嵌入会让新方法看起来已经实现的默认实现。通用完整性检查遵循已编写的 Harness 目录,并拒绝 `agent` 中的任何其他导出接口。 对于设计层面的拒绝,请直接实现该方法: @@ -105,7 +104,7 @@ Session 在其已连接的 Runtime 中拥有一个可复用的 Executor;Turn **绑定。** Runtime 将其 Executor 记录绑定到 Session、Environment、连接和不可变执行配置。恢复身份和先前 Turn 恢复标志是连续性断言,而不是配置更改。提供的原生身份必须与保留的所有者匹配;当需要现有历史时,恢复绝不能启动新的根。配置冲突属于错误,而不是热切换。连接丢失会让其所有者和句柄退役;旧计时器、输出和取消操作不能影响替代对象。 -**每 Turn 状态。** 每个 Turn 都会获得全新的包装器、输出通道和回执状态。引导和函数接口均属于该 Turn。原生回调必须在异步工作开始前捕获来源 Turn,因此迟到事件绝不会被归到当前活动的 Turn 上。原生进程、query 或传输连接、固定能力配置和原生 session 身份均属于 Executor。不要重置已完成的 `sync.Once` 值,也不要复用旧 Turn 对象。 +**每 Turn 状态。** 每个 Turn 都会获得全新的包装器、输出通道和回执状态。引导和函数结果均属于该 Turn。原生回调必须在异步工作开始前捕获来源 Turn,因此迟到事件绝不会被归到当前活动的 Turn 上。原生进程、query 或传输连接、固定能力配置和原生 session 身份均属于 Executor。不要重置已完成的 `sync.Once` 值,也不要复用旧 Turn 对象。 **开始。** `StartTurn` 返回 nil Turn,保证没有提交任何原生输入,也没有保留输出通道;随后由 Runtime 关闭该通道。一旦输入可能已经提交,即使同时返回错误,也必须返回非 nil Turn:该 Turn 拥有恰好一次的输出关闭权,并在结算前持续接受跟踪。未知输入绝不能重放。明确的 `executor_unavailable` Start 拒绝允许进行一次通用恢复尝试,但只能在此前 Executor 已关闭且未提交输入之后进行;Runtime 会重新检查同一物理对端和当前授权。 @@ -115,8 +114,8 @@ Session 在其已连接的 Runtime 中拥有一个可复用的 Executor;Turn - 错误表示结算尚未确认,既不释放所有权,也不释放容量。调用方截止时间只会停止等待,不会停止受跟踪的清理。必须串行重试同一个清理目标;清理失败会阻止替换并保留其资源槽位。 - `Executor.Close` 独立于 Turn 结果确认资源退役:不可变的 Turn 错误不得阻止在其工作和输出已经停止后关闭原生传输层。 - 结算必须包含所属的后台工作,并在失败后保留精确的原生清理目标。原生终止由适配器负责;仅有批量清理确认并不能证明已达到静默状态。 -- 每个 `Session` 都要声明 `CancellationOutcome`。`Turn` 继承该声明。快照保留已观察到的原生身份、Usage 和输出,并在取消后仍可读取。缺失的证据保持未设置;空的 `DonePayload` 表示未观察到任何内容,而不是表示取消成功或不受支持。读取快照不会等待结算。 -- `Session.Cancel` 请求取消;输出关闭表示拆卸开始。Turn 结算仍需要 `AwaitSettlement` 和所需的任何 `Executor.Close`;取消请求成功或其快照都不能替代这些等待。 +- 每个 Turn 都实现 `CancellationOutcome`。快照保留已观察到的原生身份、Usage 和输出,并在取消后仍可读取。缺失的证据保持未设置;空的 `DonePayload` 表示未观察到任何内容,而不是表示取消成功或不受支持。读取快照不会等待结算。 +- `Turn.Cancel` 请求取消;输出关闭表示拆卸开始。Turn 结算仍需要 `AwaitSettlement` 和所需的任何 `Executor.Close`;取消请求成功或其快照都不能替代这些等待。 **Runtime 在 Turn 前后执行的工作。** 一个输出消费者会在原生 Start 之前启动,耗尽有界的 64 帧通道,并将终态观察保留到 Start 发布、Turn 结算和已准入操作回执完成为止。正常完成绝不调用 Cancel。输入和函数准入会在结算前关闭;已准入的操作会持有其屏障,直至原生回执和出站确认完成。Runtime 会在等待该屏障之前向 Turn 发送取消,因为已写入的输入可能需要原生中断才能生成回执。Runtime 会汇合原生结算、所需的已确认 Executor 关闭、输出耗尽和所有已准入操作,然后应用确认或执行复用,之后才会转发 Done 或已应用的取消回执。Close 失败可以报告失败,同时保留同一 Run 和未完成操作以供重试;已关闭的调用方等待无法凭空生成已应用输入回执。Runtime 会在发布 Done 前提交原生连续性状态并释放旧 Run 的准入,因为接收方可能立即启动另一个 Turn;迟到的终态发送失败属于旧 Run,不能使已拥有 Executor 的后继对象失效。连接关闭负责传输丢失清理。结算等待时间为十秒,回执发送预算为五秒;超时不能证明已达到静默状态。 @@ -162,7 +161,7 @@ MCP、公共函数、延迟函数发现、结构化输出、图像输入、详 每个 `proto.AgentKindCapabilities` 字段都必须显式设为 `proto.CapabilitySupported` 或 `proto.CapabilityUnsupported`,即使 Harness 不可用也是如此。`proto.CapabilityUnspecified` 无效:零值和省略字段绝不表示 Unsupported。安装探测可以使用 `proto.CapabilityFromBool` 设置单个字段;但不得填充未提及字段或未来字段。可用性通过 `SupportedAgentKind.Available` 单独表示。注册会在更改 registry 之前验证完整声明;线协议会为每个字段携带显式布尔值,因此省略字段和 null 字段均无效。添加新字段时,每个生产声明都必须作出决定。Runtime 使用者应调用 `IsSupported()`,并在原生操作前拒绝不受支持的请求;接口断言用于验证实现,绝不表示支持。每个声明都必须与针对该安装验证的行为一致;[Core–Runtime protocol](../../../docs/zh/runtime-protocol.md#capability-declarations) 负责声明的传输方式和冻结方式。 -每个可用 Harness 都无需声明即实现共享 Turn 生命周期(包括 `DurableSteerer` 输入和由 `contracttest.TextLifecycle` 检查的 Turn 结算契约)、类型化的 `execution_controls` 和工具观测。在某个平台上无法满足这些要求的 Harness 在该平台报告 `Available` 为 false。`WorkspaceReadPreparation` 准入带 `workspace_read_only` 的 `execution_prepare`,由 Environment owner 就绪并提供读取,不调用 Executor 工厂。Runtime 注册不会授予 Core 资格;服务 profile 才会授予。 +每个可用 Harness 都无需声明即实现共享 Turn 生命周期(包括持久的 `SteerWithReceipt` 输入和由 `contracttest.TextLifecycle` 检查的 Turn 结算契约)、类型化的 `execution_controls` 和工具观测。在某个平台上无法满足这些要求的 Harness 在该平台报告 `Available` 为 false。`FunctionTools` 准入 `SubmitFunctionResult`。`WorkspaceReadPreparation` 准入带 `workspace_read_only` 的 `execution_prepare`,由 Environment owner 就绪并提供读取,不调用 Executor 工厂。Runtime 注册不会授予 Core 资格;服务 profile 才会授予。 可运行的仅测试示例 [`testdata/onboarding/main.go`](https://github.com/MiniMax-AI/OpenAgentCore/blob/main/apps/daemon/testdata/onboarding/main.go) 会以 `mcode` 类型注册一个仅支持文本的合成 Harness,因为 Core 只接纳[目录](harness-catalog.md)中的 Harness。它展示 Session 所有的 Executor、全新的 Turn、持久化引导、取消和历史绑定,并且绝不会发布。 diff --git a/docs/runtime-protocol.md b/docs/runtime-protocol.md index 9dc8cd249..a99d60d33 100644 --- a/docs/runtime-protocol.md +++ b/docs/runtime-protocol.md @@ -63,7 +63,6 @@ The `execution_prepare` configuration carries the Session's model configuration | `disable_execution_environment` | For an Environment of type `none` | | `local_environment` | For `openai_hosted` and `self_hosted`, with the exact Environment binding. The request carries no working directory; the Runtime checks `workspace_directory` against its binding | | `require_existing_native_session` | When a native Session must be recovered | -| `durable_receipt` on `prompt_steer` | For every active input Core delivers | An execution configuration requires exactly one of `local_environment` and `disable_execution_environment`; `execution_prepare` rejects neither or both with `unsupported_configuration`. @@ -141,7 +140,7 @@ Idle expiry of an Executor is a Runtime resource policy, separate from Core's ac ## Active input receipts -Core delivers active input as `prompt_steer` with `durable_receipt: true`, one input at a time per Run, and waits for its receipt before sending the next: +Core delivers active input as `prompt_steer`, one input at a time per Run, and waits for its receipt before sending the next. Every input is durable, with these phases: | Phase | Timer | | --- | --- | diff --git a/docs/zh/runtime-protocol.md b/docs/zh/runtime-protocol.md index 940f456da..1d593c4fb 100644 --- a/docs/zh/runtime-protocol.md +++ b/docs/zh/runtime-protocol.md @@ -1,7 +1,7 @@ --- title: "Core–Runtime 协议" source: docs/runtime-protocol.md -source_hash: 667a4f20a030e5af7787a4240c14dc648ec622c3de01c7f965813fbefac55f6d +source_hash: 30ea6710bb325ad4fb81a75ff870cb70c03cd099806f135eb92104c2616fd521 --- 此协议在 Runtime daemon 获取机器凭据后连接 Core 与 daemon,定义 daemon 连接上消息的含义和顺序。wire 类型、限制和验证器仅在 [`internal/agentdaemon/proto`](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/internal/agentdaemon/proto) 中定义一次;Core 的 [gateway](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/services/core/internal/runtimegateway) 与参考 Runtime 的 [dispatcher](https://github.com/MiniMax-AI/OpenAgentCore/tree/main/apps/daemon/internal/dispatch) 都使用它们,因此无需同步第二套 payload schema。签发凭据和打开连接的 HTTP 路由见[机器连接 API](../../contracts/agents-api/zh/machine-api.md)。 @@ -65,7 +65,6 @@ wire 上每个字段都是 JSON boolean,所有字段都必须出现,包括 ` | `disable_execution_environment` | Environment 类型为 `none` 时设置 | | `local_environment` | 为 `openai_hosted` 和 `self_hosted` 设置,包含精确的 Environment 绑定。请求不携带 working directory;Runtime 按自身绑定检查 `workspace_directory` | | `require_existing_native_session` | 需要恢复原生 Session 时设置 | -| `prompt_steer` 上的 `durable_receipt` | Core 交付的每个活动输入都设置 | 执行配置必须且只能包含 `local_environment` 和 `disable_execution_environment` 之一;两者都缺失或同时存在时,`execution_prepare` 以 `unsupported_configuration` 拒绝。 @@ -143,7 +142,7 @@ Executor 空闲到期属于 Runtime 资源策略,与 Core 的活动 Turn 并 ## 活动输入回执 {#active-input-receipts} -Core 通过 `prompt_steer` 交付活动输入,设置 `durable_receipt: true`,每个 Run 一次交付一个输入,并等待回执后再发送下一个: +Core 通过 `prompt_steer` 交付活动输入,每个 Run 一次交付一个输入,并等待回执后再发送下一个。每个输入都是持久化的,经历以下阶段: | 阶段 | 定时器 | | --- | --- | diff --git a/internal/agentdaemon/proto/steering.go b/internal/agentdaemon/proto/steering.go index 6658752cb..13e5b944b 100644 --- a/internal/agentdaemon/proto/steering.go +++ b/internal/agentdaemon/proto/steering.go @@ -8,9 +8,8 @@ const TypePromptSteerAck = "prompt_steer_ack" // PromptSteerPayload identifies one input batch within an active run. type PromptSteerPayload struct { - InputID string `json:"input_id"` - Input MessageInput `json:"input"` - DurableReceipt bool `json:"durable_receipt,omitempty"` + InputID string `json:"input_id"` + Input MessageInput `json:"input"` } // PromptSteerAckPayload distinguishes a completed write from native acceptance. diff --git a/services/core/internal/execution/delivery.go b/services/core/internal/execution/delivery.go index adc8daa5e..31844e020 100644 --- a/services/core/internal/execution/delivery.go +++ b/services/core/internal/execution/delivery.go @@ -332,7 +332,7 @@ func (d *Dispatcher) deliver(ctx context.Context, tenantID, sessionID string, pe result.ErrorCode = "message_input_unsupported" return } - if send(ctx, peer, request.Assignment, proto.TypePromptSteer, runID, proto.PromptSteerPayload{InputID: strconv.FormatInt(pending.sequence, 10), Input: pending.input, DurableReceipt: true}) != nil { + if send(ctx, peer, request.Assignment, proto.TypePromptSteer, runID, proto.PromptSteerPayload{InputID: strconv.FormatInt(pending.sequence, 10), Input: pending.input}) != nil { result.ErrorCode = "input_outcome_unknown" return } diff --git a/services/core/tests/integration/prepared_dispatch_test.go b/services/core/tests/integration/prepared_dispatch_test.go index f4bca9404..ddb6a1878 100644 --- a/services/core/tests/integration/prepared_dispatch_test.go +++ b/services/core/tests/integration/prepared_dispatch_test.go @@ -94,7 +94,7 @@ func TestPreparedDispatchPromotesOriginalBatchAndPersistsCompletion(t *testing.T h.write(frame.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 3, State: "started", RunID: start.RunID}) steering := h.read(proto.TypePromptSteer) var steer proto.PromptSteerPayload - if steering.ID != start.RunID || steering.DecodePayload(&steer) != nil || inputTextForTest(t, steer.Input) != "third" || !steer.DurableReceipt { + if steering.ID != start.RunID || steering.DecodePayload(&steer) != nil || inputTextForTest(t, steer.Input) != "third" { t.Fatal("later input bypassed ordinary steering", steer) } h.write(start.RunID, proto.TypePromptSteerAck, proto.PromptSteerAckPayload{InputID: steer.InputID, Accepted: true}) diff --git a/services/core/tests/integration/steering_receipts_test.go b/services/core/tests/integration/steering_receipts_test.go index 71c63dc04..ac020f4b5 100644 --- a/services/core/tests/integration/steering_receipts_test.go +++ b/services/core/tests/integration/steering_receipts_test.go @@ -22,8 +22,8 @@ func TestExecutionDurableInputReceiptLifetime(t *testing.T) { h.read(testExecutionRequest) extra := h.message("extra", "additional") var input proto.PromptSteerPayload - if err := h.read(proto.TypePromptSteer).DecodePayload(&input); err != nil || !input.DurableReceipt { - t.Fatal("durable receipt was not requested", err) + if err := h.read(proto.TypePromptSteer).DecodePayload(&input); err != nil { + t.Fatal(err) } h.write(first.TurnID, proto.TypePromptSteerAck, proto.PromptSteerAckPayload{InputID: input.InputID, Written: true}) status := sessions.TurnFailed