From 1cf1d5dbdf47c5989137464e90669a0739e7f0ca Mon Sep 17 00:00:00 2001 From: SaladDay <1203511142@qq.com> Date: Thu, 8 Oct 2026 13:48:02 +0000 Subject: [PATCH] Retain Turn admission declarations --- apps/daemon/internal/dispatch/steering.go | 5 + .../daemon/internal/dispatch/steering_test.go | 102 ++++++++ services/core/internal/execution/delivery.go | 8 +- .../core/internal/execution/dispatcher.go | 2 +- services/core/internal/execution/functions.go | 5 +- .../internal/execution/prepared_dispatch.go | 4 +- services/core/internal/execution/support.go | 10 - .../integration/turn_declaration_test.go | 229 ++++++++++++++++++ 8 files changed, 346 insertions(+), 19 deletions(-) create mode 100644 services/core/tests/integration/turn_declaration_test.go diff --git a/apps/daemon/internal/dispatch/steering.go b/apps/daemon/internal/dispatch/steering.go index d097ce4e1..bc2d1dea7 100644 --- a/apps/daemon/internal/dispatch/steering.go +++ b/apps/daemon/internal/dispatch/steering.go @@ -87,6 +87,11 @@ func (r *Router) queueSteering(ctx context.Context, env proto.Envelope, input pr if state.steering == nil { state.steering = make(map[string]steeringReceipt) } + if err := proto.ValidateSelection(state.declaration, proto.Selection{Messages: input.Input}); err != nil { + ack.ErrorCode, ack.Error = "unsupported", err.Error() + state.steering[input.InputID] = steeringReceipt{fingerprint: fingerprint, ack: ack} + return &ack + } 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. diff --git a/apps/daemon/internal/dispatch/steering_test.go b/apps/daemon/internal/dispatch/steering_test.go index 8bded4b25..25898a08c 100644 --- a/apps/daemon/internal/dispatch/steering_test.go +++ b/apps/daemon/internal/dispatch/steering_test.go @@ -4,6 +4,7 @@ import ( "context" "errors" "fmt" + "sync/atomic" "testing" "time" @@ -248,3 +249,104 @@ func handleSteeringAndWait(t *testing.T, h *harness, env proto.Envelope) error { waitFor(t, func() bool { return len(h.sender.snapshot()) > before }, "steering ack") return nil } + +func TestSteeringRetainsAdmittedDeclaration(t *testing.T) { + for _, test := range []struct { + name string + admitted, current proto.CapabilitySupport + available bool + }{ + {"narrowed", proto.CapabilitySupported, proto.CapabilityUnsupported, true}, + {"widened", proto.CapabilityUnsupported, proto.CapabilitySupported, true}, + {"unavailable", proto.CapabilitySupported, proto.CapabilitySupported, false}, + } { + t.Run(test.name, func(t *testing.T) { + h := newHarness(t) + defer h.router.Shutdown(context.Background()) + var calls atomic.Int32 + factory := func(_ context.Context, _ fixtureRun, out chan<- proto.Envelope) (fixtureSession, error) { + return &steeringSession{fakeSession: &fakeSession{out: out, closeOutOnCancel: true}, steer: func(_ context.Context, input proto.PromptSteerPayload, _ func()) error { + calls.Add(1) + if !input.Input.HasImages() { + t.Error("native input lost its image") + } + return nil + }}, nil + } + info := proto.SupportedAgentKind{Kind: "fixture", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{MessageImages: test.admitted})} + registerSession(h.reg, info, factory) + startRun(t, h.router, h.sender, "fixture", "active") + info.Capabilities.MessageImages, info.Available = test.current, test.available + registerSession(h.reg, info, factory) + image := "https://example.com/input.png" + input := proto.PromptSteerPayload{InputID: "image", Input: proto.MessageInput{{Content: []proto.InputContent{{Type: "input_image", ImageURL: &image}}}}} + env := scoped(t, "active", proto.TypePromptSteer, "active", input) + accepted := test.admitted.IsSupported() + code := "unsupported" + if accepted { + code = "" + } + for range 2 { + if err := handleSteeringAndWait(t, h, env); err != nil { + t.Fatal(err) + } + if ack := lastSteeringAck(t, h.sender, "active", "image"); ack.Accepted != accepted || ack.Written || ack.ErrorCode != code { + t.Fatalf("admitted declaration changed: %+v", ack) + } + } + wanted := int32(0) + if accepted { + wanted = 1 + } + if calls.Load() != wanted { + t.Fatalf("native calls=%d, want %d", calls.Load(), wanted) + } + // Rejected and accepted receipts both bind the original input identity. + changed := input + changed.Input = proto.TextInput("changed input") + if err := handleSteeringAndWait(t, h, scoped(t, "active", proto.TypePromptSteer, "active", changed)); err != nil { + t.Fatal(err) + } + if ack := lastSteeringAck(t, h.sender, "active", "image"); ack.ErrorCode != "input_conflict" { + t.Fatalf("receipt identity lost: %+v", ack) + } + foreign := env + foreign.Assignment.AssignmentID = "other" + if err := handleSteeringAndWait(t, h, foreign); err != nil { + t.Fatal(err) + } + if ack := lastSteeringAck(t, h.sender, "active", "image"); ack.ErrorCode != proto.AssignmentConflict { + t.Fatalf("receipt escaped assignment: %+v", ack) + } + if test.available { + // Another admitted Turn uses the current declaration, without changing + // the active Turn's permissions or replaying its receipt. + startRun(t, h.router, h.sender, "fixture", "new") + if err := handleSteeringAndWait(t, h, scoped(t, "new", proto.TypePromptSteer, "new", input)); err != nil { + t.Fatal(err) + } + accepted = test.current.IsSupported() + code = "unsupported" + if accepted { + code = "" + wanted++ + } + if ack := lastSteeringAck(t, h.sender, "new", "image"); ack.Accepted != accepted || ack.ErrorCode != code { + t.Fatalf("new Turn ignored current declaration: %+v", ack) + } + } else { + assign(t, h.router, "new", "") + prepare := scoped(t, "new", proto.TypeExecutionPrepare, "prepare-new", noEnvironmentPreparation("new", proto.PromptRequestPayload{AgentKind: "fixture"})) + if err := h.router.Handle(t.Context(), prepare); err == nil { + t.Fatal("new admission accepted an unavailable kind") + } + if status := waitPreparationStatus(t, h.sender, prepare.ID, "rejected", ""); status.ErrorCode != "resource_unavailable" { + t.Fatalf("new admission ignored unavailability: %+v", status) + } + } + if calls.Load() != wanted { + t.Fatalf("native calls=%d, want %d", calls.Load(), wanted) + } + }) + } +} diff --git a/services/core/internal/execution/delivery.go b/services/core/internal/execution/delivery.go index cf17072c0..bd6a63684 100644 --- a/services/core/internal/execution/delivery.go +++ b/services/core/internal/execution/delivery.go @@ -53,7 +53,7 @@ func abort(peer *runtimegateway.Session, ref proto.AssignmentRef, runID string) _ = send(context.Background(), peer, ref, proto.TypePromptCancel, runID, proto.PromptCancelPayload{}) } -func (d *Dispatcher) deliver(ctx context.Context, tenantID, sessionID string, peer *runtimegateway.Session, request proto.PromptRequestPayload, runID string, input proto.MessageInput, first int64, prepared *preparedStart) (result Result, status string) { +func (d *Dispatcher) deliver(ctx context.Context, tenantID, sessionID string, peer *runtimegateway.Session, request proto.PromptRequestPayload, declaration proto.Declaration, runID string, input proto.MessageInput, first int64, prepared *preparedStart) (result Result, status string) { changed, unsubscribeChanges := d.notifications.subscribe(tenantID, sessionID) defer unsubscribeChanges() status = sessions.TurnFailed @@ -102,7 +102,7 @@ func (d *Dispatcher) deliver(ctx context.Context, tenantID, sessionID string, pe var pending *pendingInput var cancelSent time.Time var cancelReply <-chan cancellationResult - functions := &functionExchange{assignment: prepared.assignment, kind: request.AgentKind, turns: d.SessionsReader, sessions: d.sessionExecution, tenant: tenantID, session: sessionID, turn: runID, tools: request.FunctionTools} + functions := &functionExchange{assignment: prepared.assignment, turns: d.SessionsReader, sessions: d.sessionExecution, tenant: tenantID, session: sessionID, turn: runID, tools: request.FunctionTools} done := false cancelCtx, stopCancellation := context.WithCancel(ctx) defer stopCancellation() @@ -289,7 +289,7 @@ func (d *Dispatcher) deliver(ctx context.Context, tenantID, sessionID string, pe cancelSent = time.Now() continue } - if err := functions.start(cancelCtx, peer); err != nil { + if err := functions.start(cancelCtx, peer, declaration); err != nil { result.ErrorCode = "function_result_invalid" return } @@ -328,7 +328,7 @@ func (d *Dispatcher) deliver(ctx context.Context, tenantID, sessionID string, pe return } if !pending.waiting && !pending.written { - if validateDelivery(peer, request.AgentKind, proto.Selection{Messages: pending.input}) != nil { + if proto.ValidateSelection(declaration, proto.Selection{Messages: pending.input}) != nil { result.ErrorCode = "message_input_unsupported" return } diff --git a/services/core/internal/execution/dispatcher.go b/services/core/internal/execution/dispatcher.go index d9b6304d0..2cf721769 100644 --- a/services/core/internal/execution/dispatcher.go +++ b/services/core/internal/execution/dispatcher.go @@ -120,7 +120,7 @@ func (d *Dispatcher) Run(ctx context.Context, tenantID, sessionID, turnID string if _, err := d.sessionExecution.TransitionTurn(ctx, tenantID, sessionID, turnID, sessions.TurnTransition{ExpectedStatus: sessions.TurnQueued, Status: sessions.TurnInProgress}); err != nil { return sessions.Turn{}, err } - result, status := d.deliver(ctx, tenantID, sessionID, peer, req, turnID, text, through, prepared) + result, status := d.deliver(ctx, tenantID, sessionID, peer, req, declaration, turnID, text, through, prepared) return d.finishRun(tenantID, sessionID, turnID, snapshot.Agent.Model, result, status) } diff --git a/services/core/internal/execution/functions.go b/services/core/internal/execution/functions.go index b6eda7adf..8150184c0 100644 --- a/services/core/internal/execution/functions.go +++ b/services/core/internal/execution/functions.go @@ -53,7 +53,6 @@ type functionExchange struct { sessions *sessions.ExecutionOperations tenant, session, turn string assignment proto.AssignmentRef - kind string tools []proto.FunctionTool callID string reply <-chan functionReply @@ -77,7 +76,7 @@ func (f *functionExchange) record(ctx context.Context, env proto.Envelope) error return f.unlessCancelling(ctx, err) } -func (f *functionExchange) start(ctx context.Context, peer *runtimegateway.Session) error { +func (f *functionExchange) start(ctx context.Context, peer *runtimegateway.Session, declaration proto.Declaration) error { if f.reply != nil || len(f.tools) == 0 { return nil } @@ -93,7 +92,7 @@ func (f *functionExchange) start(ctx context.Context, peer *runtimegateway.Sessi if err != nil { return err } - if err := validateDelivery(peer, f.kind, proto.Selection{FunctionResult: &result}); err != nil { + if err := proto.ValidateSelection(declaration, proto.Selection{FunctionResult: &result}); err != nil { return err } env, err := proto.NewEnvelope(proto.TypeFunctionResult, f.turn, result) diff --git a/services/core/internal/execution/prepared_dispatch.go b/services/core/internal/execution/prepared_dispatch.go index 6c7a19683..427d14daf 100644 --- a/services/core/internal/execution/prepared_dispatch.go +++ b/services/core/internal/execution/prepared_dispatch.go @@ -88,6 +88,8 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t if err != nil || run.Reservation.State != sessions.EnvironmentInputPending { return run, err } + // Recheck before claiming without replacing the declaration that admitted + // preparation; later heartbeats cannot expand this Turn's operations. if _, err := admitSession(peer, session.Engine, snapshot, messages); err != nil { return run, err } @@ -108,7 +110,7 @@ func (d *Dispatcher) RunEnvironmentInput(ctx context.Context, lease Ownership, t } turnID := run.Reservation.Receipts[0].TurnID through := run.Reservation.Receipts[len(run.Reservation.Receipts)-1].Sequence - result, status := d.deliver(owner, tenantID, sessionID, peer, req, turnID, messages, through, prepared) + result, status := d.deliver(owner, tenantID, sessionID, peer, req, declaration, turnID, messages, through, prepared) result, status = d.captureCompletedArtifacts(owner, peer, session, environment, bound.Device, turnID, result, status) run.Turn, err = d.finishRun(tenantID, sessionID, turnID, snapshot.Agent.Model, result, status) return run, err diff --git a/services/core/internal/execution/support.go b/services/core/internal/execution/support.go index 62c80618d..e88ebb09e 100644 --- a/services/core/internal/execution/support.go +++ b/services/core/internal/execution/support.go @@ -141,13 +141,3 @@ func admitSession(peer *runtimegateway.Session, engine string, snapshot Snapshot selection.Messages = messages return declaration, proto.ValidateSelection(declaration, selection) } - -// validateDelivery checks a message or function result delivered to a running -// Turn against the peer's declaration. -func validateDelivery(peer *runtimegateway.Session, engine string, selection proto.Selection) error { - declaration, err := runtimeDeclaration(peer, engine) - if err != nil { - return err - } - return proto.ValidateSelection(declaration, selection) -} diff --git a/services/core/tests/integration/turn_declaration_test.go b/services/core/tests/integration/turn_declaration_test.go new file mode 100644 index 000000000..3bcf540c5 --- /dev/null +++ b/services/core/tests/integration/turn_declaration_test.go @@ -0,0 +1,229 @@ +package integration + +import ( + "context" + "encoding/json" + "strings" + "testing" + "time" + + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto" + "github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/execution" + "github.com/MiniMax-AI/OpenAgentCore/services/core/internal/sessions" +) + +func awaitTurnDeclaration(t *testing.T, h *dispatchHarness, declaration proto.SupportedAgentKind) { + t.Helper() + h.write("", proto.TypeHeartbeat, proto.HeartbeatPayload{HomeRemoval: proto.CapabilityUnsupported, SupportedAgentKinds: []proto.SupportedAgentKind{declaration}}) + awaitDaemonRemoteCondition(t, t.Context(), 3*time.Second, "same-connection declaration update", func() bool { + peer, err := h.registry.LookupDevice(h.device.ID) + if err != nil { + return false + } + current, found, known := peer.AgentKindStatus(declaration.Kind) + return found && known && current == declaration + }) +} + +func assertFreshTurnRejected(t *testing.T, h *dispatchHarness, turn, reason string) { + t.Helper() + ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second) + defer cancel() + if _, err := h.bound().Run(ctx, h.tenant, h.session.ID, turn); err == nil || !strings.Contains(err.Error(), reason) { + t.Fatalf("fresh Turn did not use current declaration: %v", err) + } + stored, err := sessionAdapter(h.s).GetTurn(t.Context(), h.tenant, h.session.ID, turn) + if err != nil || stored.Status != sessions.TurnQueued { + t.Fatal("rejected fresh Turn was claimed", stored, err) + } +} + +func TestTurnRetainsMessageDeclaration(t *testing.T) { + for _, change := range []string{"narrow", "unavailable", "widen", "reprepare"} { + t.Run(change, func(t *testing.T) { + h := newDispatchHarness(t) + declaration := proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{ + EnvironmentNone: proto.CapabilitySupported, NativeSessionRecovery: proto.CapabilitySupported, + MessageImages: proto.CapabilitySupported, + })} + widen := change == "widen" || change == "reprepare" + if widen { + declaration.Capabilities.MessageImages = proto.CapabilityUnsupported + } + awaitTurnDeclaration(t, h, declaration) + first := h.message("first", "Start") + running := h.run(t.Context(), first.TurnID) + if change == "reprepare" { + frame, original := readyExecutorAttempt(t, h) + declaration.Capabilities.MessageImages = proto.CapabilitySupported + awaitTurnDeclaration(t, h, declaration) + h.write(frame.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{State: "rejected", Operation: proto.TypeExecutionStart, ErrorCode: "executor_unavailable"}) + next, replacement := readyExecutorAttempt(t, h) + if next.ID == frame.ID || replacement.ExecutorID == original.ExecutorID || replacement.RunID != first.TurnID { + t.Fatal("repreparation changed Turn or reused Executor", replacement) + } + h.write(next.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: replacement.Handle, ExecutorID: replacement.ExecutorID, Revision: 2, State: "started", RunID: replacement.RunID}) + } else { + h.read(testExecutionRequest) + switch change { + case "narrow": + declaration.Capabilities.MessageImages = proto.CapabilityUnsupported + case "unavailable": + declaration.Available = false + case "widen": + declaration.Capabilities.MessageImages = proto.CapabilitySupported + } + awaitTurnDeclaration(t, h, declaration) + } + image := imageAdmissionBatch()[1] + receipts, err := submitInputs(t.Context(), h.s, h.tenant, h.session.ID, "image", []sessions.Input{image}) + if err != nil || receipts[0].TurnID != first.TurnID { + t.Fatal(receipts, err) + } + if widen { + turn := h.finished(running, sessions.TurnFailed) + var outcome execution.Result + if json.Unmarshal(turn.Outcome, &outcome) != nil || outcome.ErrorCode != "message_input_unsupported" || outcome.AppliedThrough != first.Sequence { + t.Fatal("heartbeat granted unadmitted image steering", string(turn.Outcome)) + } + } else { + var steer proto.PromptSteerPayload + if h.read(proto.TypePromptSteer).DecodePayload(&steer) != nil || !steer.Input.HasImages() { + t.Fatal("admitted image did not reach Runtime", steer) + } + h.write(first.TurnID, proto.TypePromptSteerAck, proto.PromptSteerAckPayload{InputID: steer.InputID, Accepted: true}) + h.write(first.TurnID, proto.TypeDone, proto.DonePayload{Metadata: map[string]any{proto.DoneMetaAgentSessionID: "retained-message-native"}}) + turn := h.finished(running, sessions.TurnCompleted) + var outcome execution.Result + if json.Unmarshal(turn.Outcome, &outcome) != nil || outcome.AppliedThrough != receipts[0].Sequence { + t.Fatal("image receipt did not advance the Turn", string(turn.Outcome)) + } + } + next, err := submitInputs(t.Context(), h.s, h.tenant, h.session.ID, "fresh-image", []sessions.Input{image}) + if err != nil || next[0].TurnID == first.TurnID { + t.Fatal("new input did not create a fresh Turn", next, err) + } + if !widen { + reason := "message images" + if change == "unavailable" { + reason = "available" + } + assertFreshTurnRejected(t, h, next[0].TurnID, reason) + return + } + running = h.run(t.Context(), next[0].TurnID) + var start testExecution + if h.read(testExecutionRequest).DecodePayload(&start) != nil || !start.Input.HasImages() { + t.Fatal("fresh Turn did not admit newly supported image", start) + } + h.write(next[0].TurnID, proto.TypeDone, proto.DonePayload{}) + h.finished(running, sessions.TurnCompleted) + }) + } +} + +func TestTurnRetainsFunctionDeclaration(t *testing.T) { + for _, change := range []string{"narrow", "unavailable", "widen"} { + for _, image := range []bool{false, true} { + name := change + "/text" + if image { + name = change + "/image" + } + t.Run(name, func(t *testing.T) { + h := newFunctionHarness(t) + declaration := proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{ + EnvironmentNone: proto.CapabilitySupported, NativeSessionRecovery: proto.CapabilitySupported, + FunctionTools: proto.CapabilitySupported, FunctionResultImages: proto.CapabilitySupported, + })} + if change == "widen" { + declaration.Capabilities.FunctionResultImages = proto.CapabilityUnsupported + } + awaitTurnDeclaration(t, h, declaration) + first := h.message("first", "Call the function") + running := h.run(t.Context(), first.TurnID) + h.read(testExecutionRequest) + switch change { + case "narrow": + declaration.Capabilities.FunctionTools = proto.CapabilityUnsupported + declaration.Capabilities.FunctionResultImages = proto.CapabilityUnsupported + case "unavailable": + declaration.Available = false + case "widen": + declaration.Capabilities.FunctionResultImages = proto.CapabilitySupported + } + awaitTurnDeclaration(t, h, declaration) + for attempt := range 2 { + h.write(first.TurnID, proto.TypeFunctionCall, proto.FunctionCallPayload{CallID: "lookup", Name: "lookup_ticket", Arguments: json.RawMessage(`{}`)}) + state := functionState(t, h, 1) + content := `{"success":true,"output":"answer"}` + if image { + content = `{"success":true,"output":[{"type":"input_image","image_url":"data:image/png;base64,AA=="}]}` + } + id := state.RequiredActions[0].CallID + if err := SubmitFixtureFunctionResult(t.Context(), h.s, h.tenant, h.session.ID, first.TurnID, id, json.RawMessage(content)); err != nil { + t.Fatal(err) + } + denied := attempt == 0 && change == "widen" && image + if denied { + turn := h.finished(running, sessions.TurnFailed) + var outcome execution.Result + if json.Unmarshal(turn.Outcome, &outcome) != nil || outcome.ErrorCode != "function_result_invalid" { + t.Fatal("heartbeat granted unadmitted function image", string(turn.Outcome)) + } + } else { + var result proto.FunctionResultPayload + if h.read(proto.TypeFunctionResult).DecodePayload(&result) != nil || result.CallID != "lookup" || !result.Success || len(result.Content) != 1 { + t.Fatal("function result did not reach Runtime", result) + } + if image && (result.Content[0].ImageURL == nil || *result.Content[0].ImageURL != "data:image/png;base64,AA==") || !image && (result.Content[0].Text == nil || *result.Content[0].Text != "answer") { + t.Fatal("function content changed", result.Content) + } + h.write(first.TurnID, proto.TypeInteractionDecisionAck, proto.InteractionDecisionAckPayload{DeliveryID: result.DeliveryID, Applied: true}) + h.write(first.TurnID, proto.TypeDone, proto.DonePayload{Metadata: map[string]any{proto.DoneMetaAgentSessionID: "retained-functions-native"}}) + h.finished(running, sessions.TurnCompleted) + } + saved, err := FixtureFunctionCall(t.Context(), h.s.pool, h.tenant, h.session.ID, first.TurnID, id) + if err != nil || saved.Applied == denied { + t.Fatal("function receipt disagrees with admitted delivery", saved.Applied, err) + } + if attempt == 1 { + break + } + first = h.message("fresh", "Call again") + if change != "widen" { + reason := "function tools" + if change == "unavailable" { + reason = "available" + } + assertFreshTurnRejected(t, h, first.TurnID, reason) + break + } + running = h.run(t.Context(), first.TurnID) + h.read(testExecutionRequest) + } + }) + } + } +} + +func TestPreparedTurnKeepsPreclaimDeclaration(t *testing.T) { + h, pending := preparedDispatchHarness(t) + result := runPreparedDispatch(h, t.Context(), pending) + frame := h.read(proto.TypeExecutionPrepare) + handle := acknowledgePreparation(h, frame.ID) + caps := workerEnvironmentCapabilities() + caps.MessageImages = proto.CapabilitySupported + awaitFixtureCapabilities(t, h, caps) + start := readyPreparedDispatch(t, h, frame.ID, handle) + h.write(frame.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 3, State: "started", RunID: start.RunID}) + if _, err := submitInputs(t.Context(), h.s, h.tenant, h.session.ID, "image", imageAdmissionBatch()[1:]); err != nil { + t.Fatal(err) + } + got := awaitPreparedDispatch(t, result) + var outcome execution.Result + if got.err != nil || got.run.Turn.Status != sessions.TurnFailed || json.Unmarshal(got.run.Turn.Outcome, &outcome) != nil || outcome.ErrorCode != "message_input_unsupported" { + t.Fatal("final preclaim widened the preparation declaration", got.err, string(got.run.Turn.Outcome)) + } + assertPreparationReleased(t, h, frame.ID, handle) +}