Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 2 additions & 6 deletions apps/daemon/internal/agent/claudesdk/cancellation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,11 +41,7 @@ func TestCancellationWaitsForDrainAndPublishesOutcome(t *testing.T) {
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("unsettled cancellation reported %v", err)
}
provider, ok := running.(interface{ CancellationOutcome() proto.DonePayload })
if !ok {
t.Fatal("missing cancellation outcome provider")
}
if got := provider.CancellationOutcome(); !reflect.DeepEqual(got, proto.DonePayload{}) {
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) {
Expand All @@ -62,7 +58,7 @@ func TestCancellationWaitsForDrainAndPublishesOutcome(t *testing.T) {
default:
t.Fatal("successful cancellation preceded owned process release")
}
got := provider.CancellationOutcome()
got := running.CancellationOutcome()
if got.Content != "partialtaildrained" || got.Metadata[proto.DoneMetaAgentSessionID] != "native-session" || got.Usage.Raw["claude_sdk_result"] == nil || got.Usage.Tokens != nil {
t.Fatalf("lost drained cancellation outcome: %+v", got)
}
Expand Down
11 changes: 8 additions & 3 deletions apps/daemon/internal/agent/harness.go
Original file line number Diff line number Diff line change
Expand Up @@ -110,7 +110,6 @@ type Executor interface {
type Turn interface {
Session
DurableSteerer
CancellationOutcome() proto.DonePayload
// Success confirms closed output and settled native input, function,
// interaction and child-work obligations. Errors cannot prove cancellation.
AwaitSettlement(context.Context) (TurnSettlement, error)
Expand All @@ -123,13 +122,19 @@ type TurnSettlement struct {
Reason string
}

// Session is the cancellation surface shared by direct prompt runs and Turns.
// Session is the cancellation and outcome surface shared by direct prompt runs
// and Turns. Every owner exposes observed state, including direct-call sessions.
// 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
// happens asynchronously and is signalled via the out channel close.
Cancel(ctx context.Context) error
// CancellationOutcome snapshots observed native identity, usage and output.
// It remains readable after Cancel; missing evidence stays unset. An empty
// result means no observed evidence, not unsupported cancellation or success.
// Reading it does not wait for or establish native settlement.
CancellationOutcome() proto.DonePayload
}

// Turn extension contracts. Every public Harness implements each interface;
Expand Down Expand Up @@ -215,10 +220,10 @@ type Prepared interface {
// same native resource across Start. Read-only preparations need only Prepared.
type PreparedCancellation interface {
Prepared
Session
// Cancel returns after local cleanup and all output writes have stopped.
// An error retains ownership so callers can retry this exact object serially.
Cancel(context.Context) error
CancellationOutcome() proto.DonePayload
}

// A factory may return both a resource and an error when construction failed but
Expand Down
2 changes: 2 additions & 0 deletions apps/daemon/internal/agent/registry_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@ func stubFactory(marker string) agent.Factory {

type stubSession struct{ marker string }

func (stubSession) CancellationOutcome() proto.DonePayload { return proto.DonePayload{} }

func (stubSession) Cancel(context.Context) error { return nil }
func (stubSession) SubmitPermission(context.Context, string, proto.PermissionDecisionPayload) error {
return nil
Expand Down
21 changes: 8 additions & 13 deletions apps/daemon/internal/dispatch/cancellation.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,10 @@ import (
"context"
"errors"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto"
)

type cancellationOutcomeProvider interface {
CancellationOutcome() proto.DonePayload
}

func (r *Router) releaseCompletedSession(state *sessionState) error {
r.mu.Lock()
state.retain = false
Expand All @@ -35,7 +32,7 @@ func (r *Router) handlePromptCancel(ctx context.Context, env proto.Envelope) err
if state != nil {
state.retain = false
}
var cancelSession func(context.Context) error
var owner agent.Session
if state != nil && state.preparedHandoff != nil {
handoff := state.preparedHandoff
release, attempt := r.claimPreparedReleaseLocked(state, true, "", true)
Expand All @@ -47,22 +44,20 @@ func (r *Router) handlePromptCancel(ctx context.Context, env proto.Envelope) err
return nil
}
if state != nil && state.session != nil {
cancelSession = state.session.Cancel
owner = state.session
}
r.mu.Unlock()
ack := proto.InteractionDecisionAckPayload{DeliveryID: request.DeliveryID, ErrorCode: "run_inactive"}
if state != nil && cancelSession == nil {
if state != nil && owner == nil {
ack.ErrorCode = "not_ready"
} else if cancelSession != nil {
if err := cancelSession(ctx); err != nil {
} else if owner != nil {
if err := owner.Cancel(ctx); err != nil {
r.log.WarnContext(ctx, "session.Cancel failed", "run_id", env.ID, "err", err)
ack.ErrorCode = "cancel_failed"
} else {
ack.Applied, ack.ErrorCode = true, ""
if provider, ok := state.session.(cancellationOutcomeProvider); ok {
outcome := provider.CancellationOutcome()
ack.Outcome = &outcome
}
outcome := owner.CancellationOutcome()
ack.Outcome = &outcome
}
state.ctxCancel()
}
Expand Down
34 changes: 24 additions & 10 deletions apps/daemon/internal/dispatch/cancellation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import "github.com/MiniMax-AI/OpenAgentCore/internal/harnessconfig"
import (
"context"
"errors"
"reflect"
"testing"

"github.com/MiniMax-AI/OpenAgentCore/apps/daemon/internal/agent"
Expand All @@ -17,10 +18,11 @@ type cancelReceiptSession struct {
entered chan struct{}
release chan struct{}
err error
outcome proto.DonePayload
}

func (s *cancelReceiptSession) CancellationOutcome() proto.DonePayload {
return proto.DonePayload{Metadata: map[string]any{proto.DoneMetaAgentSessionID: "native-cancelled"}}
return s.outcome
}

func TestCompletionWaitsForNativeWriterRelease(t *testing.T) {
Expand Down Expand Up @@ -58,14 +60,22 @@ func (s *cancelReceiptSession) Cancel(ctx context.Context) error {
}

func TestCancellationReceiptFollowsAdapterOutcome(t *testing.T) {
for _, fails := range []bool{false, true} {
t.Run(map[bool]string{false: "applied", true: "rejected"}[fails], func(t *testing.T) {
observed := proto.DonePayload{Content: "partial output", Metadata: map[string]any{proto.DoneMetaAgentSessionID: "native-cancelled"}}
for _, test := range []struct {
name string
outcome proto.DonePayload
err error
}{
{name: "observed", outcome: observed},
{name: "unknown"},
{name: "failed", outcome: observed, err: errors.New("adapter could not cancel")},
{name: "unsupported", outcome: observed, err: agent.ErrUnsupportedOperation},
{name: "deadline", outcome: observed, err: context.DeadlineExceeded},
} {
t.Run(test.name, func(t *testing.T) {
h := newHarness(t)
defer h.router.Shutdown(context.Background())
sess := &cancelReceiptSession{entered: make(chan struct{}), release: make(chan struct{})}
if fails {
sess.err = errors.New("adapter could not cancel")
}
sess := &cancelReceiptSession{entered: make(chan struct{}), release: make(chan struct{}), outcome: test.outcome, err: test.err}
h.reg.RegisterKind(proto.SupportedAgentKind{Kind: "codex", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, harnessconfig.Configuration{}, func(ctx context.Context, req proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) {
sess.fakeSession = &fakeSession{out: out, closeOutOnCancel: true}
return sess, nil
Expand Down Expand Up @@ -93,11 +103,15 @@ func TestCancellationReceiptFollowsAdapterOutcome(t *testing.T) {
found = true
var ack proto.InteractionDecisionAckPayload
_ = env.DecodePayload(&ack)
if ack.Applied == fails || ack.DeliveryID != "cancel-1" {
if ack.Applied != (test.err == nil) || ack.DeliveryID != "cancel-1" {
t.Fatalf("wrong receipt: %+v", ack)
}
if !fails && (ack.Outcome == nil || ack.Outcome.Metadata[proto.DoneMetaAgentSessionID] != "native-cancelled") {
t.Fatal("cancellation receipt lost native identity")
if test.err == nil {
if ack.ErrorCode != "" || ack.Outcome == nil || !reflect.DeepEqual(*ack.Outcome, test.outcome) {
t.Fatalf("cancellation receipt changed observed evidence: %+v", ack)
}
} else if ack.ErrorCode != "cancel_failed" || ack.Outcome != nil {
t.Fatalf("failed cancellation supplied a success outcome: %+v", ack)
}
}
}
Expand Down
8 changes: 3 additions & 5 deletions apps/daemon/internal/dispatch/optional_interactions_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,10 +11,8 @@ import (
"github.com/MiniMax-AI/OpenAgentCore/internal/agentdaemon/proto/prototest"
)

// Wrapping only Cancel proves that no responder stubs are required for a Session.
type lifecycleOnly struct{ cancel func(context.Context) error }

func (s lifecycleOnly) Cancel(ctx context.Context) error { return s.cancel(ctx) }
// Wrapping Session proves that no responder stubs are required for a Session.
type lifecycleOnly struct{ agent.Session }

func TestOptionalInteractionResponders(t *testing.T) {
for _, ask := range []bool{false, true} {
Expand All @@ -25,7 +23,7 @@ func TestOptionalInteractionResponders(t *testing.T) {
h.reg.RegisterKind(proto.SupportedAgentKind{Kind: "minimal", Available: true, Capabilities: prototest.Capabilities(proto.AgentKindCapabilities{})}, harnessconfig.Configuration{}, func(_ context.Context, _ proto.PromptRequestPayload, out chan<- proto.Envelope) (agent.Session, error) {
output = out
s := &fakeSession{out: out, closeOutOnCancel: true}
return lifecycleOnly{cancel: s.Cancel}, nil
return lifecycleOnly{Session: s}, nil
})
if err := h.router.Handle(t.Context(), mustEnv(t, proto.TypePromptRequest, "run", proto.PromptRequestPayload{AgentKind: "minimal"})); err != nil {
t.Fatal(err)
Expand Down
4 changes: 2 additions & 2 deletions apps/daemon/internal/dispatch/preparation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -65,8 +65,8 @@ func (p *controlledPreparation) CancellationOutcome() proto.DonePayload {
p.mu.Lock()
session := p.session
p.mu.Unlock()
if provider, ok := session.(interface{ CancellationOutcome() proto.DonePayload }); ok {
return provider.CancellationOutcome()
if session != nil {
return session.CancellationOutcome()
}
return proto.DonePayload{}
}
Expand Down
2 changes: 2 additions & 0 deletions apps/daemon/internal/dispatch/router_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,8 @@ type askCall struct {
decision proto.PromptForUserChoiceDecisionPayload
}

func (s *fakeSession) CancellationOutcome() proto.DonePayload { return proto.DonePayload{} }

func (s *fakeSession) Cancel(context.Context) error {
s.cancelMu.Lock()
s.cancelCalls++
Expand Down
2 changes: 2 additions & 0 deletions apps/daemon/internal/dispatch/suspend_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,8 @@ func (s suspendSender) Send(ctx context.Context, env proto.Envelope) error { ret

type suspendedSession struct{ cancelled atomic.Int32 }

func (s *suspendedSession) CancellationOutcome() proto.DonePayload { return proto.DonePayload{} }

func (s *suspendedSession) Cancel(context.Context) error { s.cancelled.Add(1); return nil }

func suspensionRouter(t *testing.T, sender Sender) *Router {
Expand Down
5 changes: 3 additions & 2 deletions contracts/agents-api/harness-onboarding.md
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,7 @@ For example, the Codex adapter keeps its app-server and thread, the Claude adapt
| 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 |
| `Turn`, `Session.Cancel`, `CancellationOutcome`, `AwaitSettlement` | Real implementation | Cancel the exact Turn, keep observed results and confirm settlement independently of cancellation requests |
| `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 |
Expand Down Expand Up @@ -112,7 +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.
- The observed cancellation outcome keeps native identity, Usage and output without fabricating missing evidence.
- Every `Session`, including a direct-call factory result, declares `CancellationOutcome`. `Turn` and `PreparedCancellation` inherit 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.
- Direct-call `Session.Cancel` requests cancellation; output closure signals teardown. Executable `PreparedCancellation.Cancel` waits for local cleanup and output writes to stop. 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, function and interaction 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.

Expand Down
Loading