diff --git a/apps/docs/content/docs/environments-and-files.mdx b/apps/docs/content/docs/environments-and-files.mdx index 361abedf7..f74e1439f 100644 --- a/apps/docs/content/docs/environments-and-files.mdx +++ b/apps/docs/content/docs/environments-and-files.mdx @@ -239,7 +239,10 @@ affect a later reservation or Turn. Evaluate deadlines after acquiring the Sessi lock, and return terminal storage outcomes without rolling their transaction back. A validated Runtime `failed` response with `preparation_failed` and no Run settles -the pending input immediately with `runtime_preparation_failed`. It records a +the pending input immediately with `runtime_preparation_failed`. Before admission, +a rejection of the current `execution_prepare` request with no Run and code +`invalid_configuration`, `unsupported_configuration` or `unsupported_preparation` +settles through that same failure path. It records a safe Session failure before any Turn exists and releases the input gate; fixing the local cause allows new input. Transport loss, capacity rejection and unconfirmed cleanup remain retryable within the original deadline. Core uses diff --git a/apps/docs/content/guide-sources.json b/apps/docs/content/guide-sources.json index ec05b3614..6c9c5fc66 100644 --- a/apps/docs/content/guide-sources.json +++ b/apps/docs/content/guide-sources.json @@ -13,7 +13,7 @@ "contracts/agents-api/execution-tools.md": "8cc0dbe207e8e80ac104bf482c37ea297cd64250b51d288ed4baf553756424d2", "docs/api/public-agent-api.md": "00979732412a971013e8f01b4c74820ff51aafdded6b0f105d25a78af86a6627", "docs/examples.md": "0e1bdaeff51c9c36779f817be31ea9816b7d8cb2290cf8350d3c4801f436f4f9", - "contracts/agents-api/environments.md": "424dac4bfd63418afc314896dd6a3e4c5c4375dd0a8bb7579920c5dd9743da22", + "contracts/agents-api/environments.md": "29be005ded0ad9fae226bd590b5becb880cc7de0380ce847f16a153e7c3c3195", "docs/getting-started/nodes.md": "7c1b7e364ce85b5b5916358cafc618f29042b2125ac45653be665ba07e772fdb", "docs/getting-started/self-hosted.md": "5ade597e09cbfa2e321f341693b47730b3696fe7bd4d6834eafbe6b279a55e6b", "docs/self-hosted-native.md": "b3cf736f88792e6925c50b81a4c33eba6b9ee86e195f9ed4e2187db7bed71d88", @@ -44,7 +44,7 @@ "content/docs/agents-and-tools.mdx": "44dfde4e3b3906b30323c2e75a89650ae4837210c7ba425be2266edf87ce8ff8", "content/docs/sessions.mdx": "bb63d799ed90038836d652f3f866d0309822425114c6b2aa36654b8a991947f6", "content/docs/examples.mdx": "587061e65ab2e841d14980539ba94e2216d4e8135c948a852b9a8bb0359c85d7", - "content/docs/environments-and-files.mdx": "6d8b2a5d5e95e3a5eff47ce93f0239c7cc38d4678bac53c58f5c8df016c23dc5", + "content/docs/environments-and-files.mdx": "88d1fc4129270302a1f4bebd7febfaac622f42dcb1f3f5524b4fa49e10106479", "content/docs/hosted-providers.mdx": "9c241090709a9929ab6a34615db1e20a94c1f36649026281836060e81ac40b4c", "content/docs/self-hosted-execution.mdx": "eeed4c6b3646927ccc3c7ac2d20b2c0f9c8a65e42d734310fe3959e016324e57", "content/docs/self-hosted-native.mdx": "671451580dfb4f5c134221991f68292a4f992008174d23190b1da3742046571f", diff --git a/contracts/agents-api/environments.md b/contracts/agents-api/environments.md index 9308ae026..d14d387cf 100644 --- a/contracts/agents-api/environments.md +++ b/contracts/agents-api/environments.md @@ -236,7 +236,10 @@ affect a later reservation or Turn. Evaluate deadlines after acquiring the Sessi lock, and return terminal storage outcomes without rolling their transaction back. A validated Runtime `failed` response with `preparation_failed` and no Run settles -the pending input immediately with `runtime_preparation_failed`. It records a +the pending input immediately with `runtime_preparation_failed`. Before admission, +a rejection of the current `execution_prepare` request with no Run and code +`invalid_configuration`, `unsupported_configuration` or `unsupported_preparation` +settles through that same failure path. It records a safe Session failure before any Turn exists and releases the input gate; fixing the local cause allows new input. Transport loss, capacity rejection and unconfirmed cleanup remain retryable within the original deadline. Core uses diff --git a/services/agents-api/internal/execution/preparation.go b/services/agents-api/internal/execution/preparation.go index 1f82bf31c..fdf6a753c 100644 --- a/services/agents-api/internal/execution/preparation.go +++ b/services/agents-api/internal/execution/preparation.go @@ -107,6 +107,13 @@ func (d *Dispatcher) awaitPreparation(ctx context.Context, tenant, session strin } status, err := prepared.observation(env) if err != nil { + var rejection *preparationRejection + if status.RunID == "" && errors.As(err, &rejection) && rejection.operation == proto.TypeExecutionPrepare { + switch rejection.code { + case "invalid_configuration", "unsupported_configuration", "unsupported_preparation": + return pending, errPreparationFailed + } + } return pending, err } if status.State == "ready" && status.RunID == "" && status.ExecutorID != "" { diff --git a/services/agents-api/internal/store/environment_input_settlement_test.go b/services/agents-api/internal/store/environment_input_settlement_test.go index d163447e2..00d176de6 100644 --- a/services/agents-api/internal/store/environment_input_settlement_test.go +++ b/services/agents-api/internal/store/environment_input_settlement_test.go @@ -112,62 +112,75 @@ func TestEnvironmentInputPromotionRollsBackHistoryAndSettlement(t *testing.T) { } func TestEnvironmentInputDeadlineIsCheckedAfterSessionLock(t *testing.T) { - s, pool := testStore(t) - lease := executionLease(t, s) - writer := lease.Store() - tenant, session := environmentInputSession(t, s) - pending := reserveEnvironmentInput(t, s, tenant, session.ID, "pending") - ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) - defer cancel() - tx, err := pool.Begin(ctx) - if err != nil { - t.Fatal(err) - } - defer func() { _ = tx.Rollback(context.Background()) }() - var blocker int32 - if err := tx.QueryRow(ctx, "SELECT pg_backend_pid() FROM sessions WHERE id=$1 FOR UPDATE", session.ID).Scan(&blocker); err != nil { - t.Fatal(err) - } - type outcome struct { - value EnvironmentInputReservation - err error - } - done := make(chan outcome, 1) - go func() { - got, err := writer.PromoteEnvironmentInput(ctx, tenant, session.ID, pending.ID) - done <- outcome{got, err} - }() - for { - var blocked bool - if err := pool.QueryRow(ctx, "SELECT EXISTS (SELECT 1 FROM pg_stat_activity WHERE $1=ANY(pg_blocking_pids(pid)))", blocker).Scan(&blocked); err != nil { - t.Fatal(err) - } - if blocked { - break - } - select { - case result := <-done: - t.Fatal("promotion bypassed Session lock", result) - case <-ctx.Done(): - t.Fatal("promotion lock wait not observed") - case <-time.After(5 * time.Millisecond): - } - } - // Transaction-start time is now older than the controlled deadline. - if _, err := tx.Exec(ctx, "UPDATE environment_input_reservations SET deadline=clock_timestamp() WHERE id=$1", pending.ID); err != nil { - t.Fatal(err) - } - if err := tx.Commit(ctx); err != nil { - t.Fatal(err) - } - result := <-done - if result.err != nil || result.value.State != EnvironmentInputExpired || result.value.SettledAt == nil { - t.Fatal("lock wait extended input lifetime", result) - } - environmentInputHistory(t, pool, session.ID, 0, 0) - stored, err := s.GetEnvironmentInputReservation(ctx, tenant, session.ID, pending.ID) - if err != nil || stored.State != EnvironmentInputExpired { - t.Fatal("expiry was rolled back", stored, err) + for _, action := range []string{"promote", "fail"} { + t.Run(action, func(t *testing.T) { + s, pool := testStore(t) + lease := executionLease(t, s) + writer := lease.Store() + tenant, session := environmentInputSession(t, s) + pending := reserveEnvironmentInput(t, s, tenant, session.ID, "pending") + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + tx, err := pool.Begin(ctx) + if err != nil { + t.Fatal(err) + } + defer func() { _ = tx.Rollback(context.Background()) }() + var blocker int32 + if err := tx.QueryRow(ctx, "SELECT pg_backend_pid() FROM sessions WHERE id=$1 FOR UPDATE", session.ID).Scan(&blocker); err != nil { + t.Fatal(err) + } + type outcome struct { + value EnvironmentInputReservation + err error + } + done := make(chan outcome, 1) + go func() { + var got EnvironmentInputReservation + var err error + if action == "promote" { + got, err = writer.PromoteEnvironmentInput(ctx, tenant, session.ID, pending.ID) + } else { + err = writer.FailEnvironmentInput(ctx, tenant, session.ID, pending.ID, "runtime_preparation_failed") + if err == nil { + got, err = s.GetEnvironmentInputReservation(ctx, tenant, session.ID, pending.ID) + } + } + done <- outcome{got, err} + }() + for { + var blocked bool + if err := pool.QueryRow(ctx, "SELECT EXISTS (SELECT 1 FROM pg_stat_activity WHERE $1=ANY(pg_blocking_pids(pid)))", blocker).Scan(&blocked); err != nil { + t.Fatal(err) + } + if blocked { + break + } + select { + case result := <-done: + t.Fatal("settlement bypassed Session lock", result) + case <-ctx.Done(): + t.Fatal("settlement lock wait not observed") + case <-time.After(5 * time.Millisecond): + } + } + // Transaction-start time is now older than the controlled deadline. + if _, err := tx.Exec(ctx, "UPDATE environment_input_reservations SET deadline=clock_timestamp() WHERE id=$1", pending.ID); err != nil { + t.Fatal(err) + } + if err := tx.Commit(ctx); err != nil { + t.Fatal(err) + } + result := <-done + if result.err != nil || result.value.State != EnvironmentInputExpired || result.value.SettledAt == nil { + t.Fatal("lock wait extended input lifetime", result) + } + environmentInputHistory(t, pool, session.ID, 0, 0) + stored, err := s.GetEnvironmentInputReservation(ctx, tenant, session.ID, pending.ID) + if err != nil || stored.State != EnvironmentInputExpired { + t.Fatal("expiry was rolled back", stored, err) + } + }) } } diff --git a/services/agents-api/internal/store/environment_inputs.go b/services/agents-api/internal/store/environment_inputs.go index d2ddd489d..e3008c7c0 100644 --- a/services/agents-api/internal/store/environment_inputs.go +++ b/services/agents-api/internal/store/environment_inputs.go @@ -185,6 +185,9 @@ func (s *Store) FailEnvironmentInput(ctx context.Context, tenantID, sessionID, r return err } return s.withEnvironmentInputSession(ctx, tenantID, sessionID, func(ctx context.Context, q *sqlc.Queries, session pgtype.UUID) error { + if err := q.ExpireEnvironmentInputReservation(ctx, sqlc.ExpireEnvironmentInputReservationParams{SessionID: session, ID: id}); err != nil { + return err + } _, err := q.FailEnvironmentInput(ctx, sqlc.FailEnvironmentInputParams{SessionID: session, ID: id, FailureCode: pgtype.Text{String: code, Valid: true}}) return err }) diff --git a/services/agents-api/internal/store/worker_preparation_failure_test.go b/services/agents-api/internal/store/worker_preparation_failure_test.go index 7040cb630..905ac3849 100644 --- a/services/agents-api/internal/store/worker_preparation_failure_test.go +++ b/services/agents-api/internal/store/worker_preparation_failure_test.go @@ -10,70 +10,91 @@ import ( ) func TestWorkerSettlesConfirmedPreparationFailureAndAcceptsNewInput(t *testing.T) { - h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), false) - enableWorkerEnvironment(t, h) - frames := workerFrames(t, h) - pending, err := h.s.ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "first", []store.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) - if err != nil { - t.Fatal(err) - } - _, stop := startEnvironmentExpiryWorker(t, h.d) - defer stop() - prepare := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) - handle := acknowledgePreparation(h, prepare.ID) - h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "failed", ErrorCode: "preparation_failed"}) - nextWorkerFrame(t, frames, proto.TypeExecutionRelease) - awaitDaemonRemoteCondition(t, t.Context(), 3*time.Second, "failed reservation settlement", func() bool { - current, err := h.s.GetEnvironmentInputReservation(t.Context(), h.tenant, h.session.ID, pending.ID) - return err == nil && current.State == store.EnvironmentInputFailed - }) - session, err := h.s.GetSession(t.Context(), h.tenant, h.session.ID) - if err != nil || session.PendingInput || session.LastTurn != nil || session.EnvironmentInputActivity == nil || session.EnvironmentInputActivity.Failure != "runtime_preparation_failed" { - t.Fatal("preparation did not release input with a safe failure", err) - } - changes, err := h.s.ListSessionEvents(t.Context(), h.tenant, h.session.ID, 0) - if err != nil { - t.Fatal(err) - } - failures := 0 - for _, change := range changes { - if change.Event.Type == "agent.session.failed" { - failures++ - if !change.Settled || change.EnvironmentInputActivity.Failure != "runtime_preparation_failed" { - t.Fatal("failure event lost settlement") + for _, code := range []string{"preparation_failed", "invalid_configuration", "unsupported_configuration", "unsupported_preparation"} { + t.Run(code, func(t *testing.T) { + h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), false) + enableWorkerEnvironment(t, h) + frames := workerFrames(t, h) + pending, err := h.s.ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "first", []store.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) + if err != nil { + t.Fatal(err) } - } - } - if failures != 1 { - t.Fatal("failure events", failures) - } - select { - case frame := <-frames: - t.Fatal("failed input retried", frame.Type) - case <-time.After(1200 * time.Millisecond): - } - next, err := h.s.ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "next", []store.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"next"}`)}}) - if err != nil { - t.Fatal("new input remained blocked", err) - } - prepare = nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) - handle = acknowledgePreparation(h, prepare.ID) - h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "ready"}) - startFrame := nextWorkerFrame(t, frames, proto.TypeExecutionStart) - var start proto.ExecutionStartPayload - if startFrame.DecodePayload(&start) != nil || start.RunID == "" || inputTextForTest(t, start.Input) != "next" { - t.Fatal("new input was not admitted") - } - current, err := h.s.GetEnvironmentInputReservation(t.Context(), h.tenant, h.session.ID, next.ID) - if err != nil || current.State != store.EnvironmentInputAdmitted { - t.Fatal("new input state", err) + _, stop := startEnvironmentExpiryWorker(t, h.d) + defer stop() + prepare := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) + if code == "preparation_failed" { + handle := acknowledgePreparation(h, prepare.ID) + h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "failed", ErrorCode: code}) + nextWorkerFrame(t, frames, proto.TypeExecutionRelease) + } else { + h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{State: "rejected", Operation: proto.TypeExecutionPrepare, ErrorCode: code}) + } + awaitDaemonRemoteCondition(t, t.Context(), 3*time.Second, "failed reservation settlement", func() bool { + current, err := h.s.GetEnvironmentInputReservation(t.Context(), h.tenant, h.session.ID, pending.ID) + return err == nil && current.State == store.EnvironmentInputFailed + }) + session, err := h.s.GetSession(t.Context(), h.tenant, h.session.ID) + if err != nil || session.PendingInput || session.LastTurn != nil || session.EnvironmentInputActivity == nil || session.EnvironmentInputActivity.Failure != "runtime_preparation_failed" { + t.Fatal("preparation did not release input with a safe failure", err) + } + changes, err := h.s.ListSessionEvents(t.Context(), h.tenant, h.session.ID, 0) + if err != nil { + t.Fatal(err) + } + failures := 0 + for _, change := range changes { + if change.Event.Type == "agent.session.failed" { + failures++ + if !change.Settled || change.EnvironmentInputActivity.Failure != "runtime_preparation_failed" { + t.Fatal("failure event lost settlement") + } + } + } + if failures != 1 { + t.Fatal("failure events", failures) + } + select { + case frame := <-frames: + t.Fatal("failed input retried", frame.Type) + case <-time.After(1200 * time.Millisecond): + } + next, err := h.s.ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "next", []store.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"next"}`)}}) + if err != nil { + t.Fatal("new input remained blocked", err) + } + prepare = nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) + handle := acknowledgePreparation(h, prepare.ID) + h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "ready"}) + startFrame := nextWorkerFrame(t, frames, proto.TypeExecutionStart) + var start proto.ExecutionStartPayload + if startFrame.DecodePayload(&start) != nil || start.RunID == "" || inputTextForTest(t, start.Input) != "next" { + t.Fatal("new input was not admitted") + } + current, err := h.s.GetEnvironmentInputReservation(t.Context(), h.tenant, h.session.ID, next.ID) + if err != nil || current.State != store.EnvironmentInputAdmitted { + t.Fatal("new input state", err) + } + h.write(start.RunID, proto.TypeDone, proto.DonePayload{Content: "complete"}) + }) } - h.write(start.RunID, proto.TypeDone, proto.DonePayload{Content: "complete"}) } func TestWorkerRetriesUncertainPreparationFailure(t *testing.T) { - for _, code := range []string{"connection_closed", "executor_cleanup_unconfirmed"} { - t.Run(code, func(t *testing.T) { + for _, response := range []struct { + name, state, operation, code, runID string + }{ + {"connection_closed", "failed", "", "connection_closed", ""}, + {"cleanup_failed", "failed", "", "executor_cleanup_unconfirmed", ""}, + {"capacity", "rejected", proto.TypeExecutionPrepare, "preparation_capacity", ""}, + {"executor_capacity", "rejected", proto.TypeExecutionPrepare, "executor_capacity", ""}, + {"busy", "rejected", proto.TypeExecutionPrepare, "resource_unavailable", ""}, + {"cleanup_rejected", "rejected", proto.TypeExecutionPrepare, "executor_cleanup_unconfirmed", ""}, + {"unknown", "rejected", proto.TypeExecutionPrepare, "unknown_rejection", ""}, + {"wrong_operation", "rejected", proto.TypeExecutionStart, "unsupported_configuration", ""}, + {"unknown_operation", "rejected", "", "unsupported_configuration", ""}, + {"run_present", "rejected", proto.TypeExecutionPrepare, "unsupported_configuration", "unconfirmed-run"}, + } { + t.Run(response.name, func(t *testing.T) { h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), false) enableWorkerEnvironment(t, h) frames := workerFrames(t, h) @@ -84,9 +105,15 @@ func TestWorkerRetriesUncertainPreparationFailure(t *testing.T) { _, stop := startEnvironmentExpiryWorker(t, h.d) defer stop() prepare := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) - handle := acknowledgePreparation(h, prepare.ID) - h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "failed", ErrorCode: code}) - nextWorkerFrame(t, frames, proto.TypeExecutionRelease) + status := proto.PreparationStatusPayload{State: response.state, Operation: response.operation, ErrorCode: response.code, RunID: response.runID} + if response.state == "failed" { + status.Handle = acknowledgePreparation(h, prepare.ID) + status.Revision = 2 + } + h.write(prepare.ID, proto.TypePreparationStatus, status) + if response.state == "failed" { + nextWorkerFrame(t, frames, proto.TypeExecutionRelease) + } nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) current, err := h.s.GetEnvironmentInputReservation(t.Context(), h.tenant, h.session.ID, pending.ID) if err != nil || current.State != store.EnvironmentInputPending || !current.Deadline.Equal(pending.Deadline) { @@ -98,3 +125,45 @@ func TestWorkerRetriesUncertainPreparationFailure(t *testing.T) { }) } } + +func TestWorkerPreparationRejectionPreservesCancellationAndNewerInput(t *testing.T) { + h := newDispatchHarnessForSession(t, []byte(`{"agent":{"model":"test-model"},"environment":{"type":"self_hosted","workspace_directory":"/workspace"}}`), false) + enableWorkerEnvironment(t, h) + frames := workerFrames(t, h) + first, err := h.s.ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "first", []store.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"first"}`)}}) + if err != nil { + t.Fatal(err) + } + _, stop := startEnvironmentExpiryWorker(t, h.d) + defer stop() + old := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) + if _, err := h.s.CancelEnvironmentInput(t.Context(), h.tenant, h.session.ID, first.ID); err != nil { + t.Fatal(err) + } + next, err := h.s.ReserveEnvironmentInput(t.Context(), h.tenant, h.session.ID, "next", []store.Input{{Kind: "message", Payload: json.RawMessage(`{"text":"next"}`)}}) + if err != nil { + t.Fatal(err) + } + rejection := proto.PreparationStatusPayload{State: "rejected", Operation: proto.TypeExecutionPrepare, ErrorCode: "unsupported_configuration"} + h.write(old.ID, proto.TypePreparationStatus, rejection) + prepare := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) + for id, state := range map[string]string{first.ID: store.EnvironmentInputCancelled, next.ID: store.EnvironmentInputPending} { + current, err := h.s.GetEnvironmentInputReservation(t.Context(), h.tenant, h.session.ID, id) + if err != nil || current.State != state { + t.Fatal("late rejection changed cancellation or newer input", err) + } + if id == next.ID && !current.Deadline.Equal(next.Deadline) { + t.Fatal("late rejection changed the newer input deadline") + } + } + // A delayed duplicate belongs to the old request, even on the same connection. + h.write(old.ID, proto.TypePreparationStatus, rejection) + handle := acknowledgePreparation(h, prepare.ID) + h.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "ready"}) + startFrame := nextWorkerFrame(t, frames, proto.TypeExecutionStart) + var start proto.ExecutionStartPayload + if startFrame.DecodePayload(&start) != nil || start.RunID == "" || inputTextForTest(t, start.Input) != "next" { + t.Fatal("stale rejection prevented the newer input from starting") + } + h.write(start.RunID, proto.TypeDone, proto.DonePayload{Content: "complete"}) +}