diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 90a743cf7..559a326ef 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -344,8 +344,22 @@ settlement distinct. Advance at most one bounded initialization operation per fu maintenance scan. At allocation EOF, begin the next page in the same call rather than consume an observation interval on an empty page. Refill at most once, retain the 32-allocation per-call bound and the five-second ticker, and never loop on an -empty store. Use process-local progress and the existing lifecycle gate. A recovered or uncertain -running installation fails and uses existing cleanup, without replaying writes. +empty store. Use process-local progress and the existing lifecycle gate. + +After a next-Turn input is durably pending, a completed managed allocation in a +suspension/recovery phase may hint this loop. Initial inputs, cold creation, +running/disabled compute, terminal receipts, cancellation/tool-result events, +history and file operations do not use this hint. Eligibility lookup and delivery +are best effort; persisted work and the normal ticker remain authoritative. +Coalesce hints without blocking, and allow at most one extra scan per normal +five-second cycle. Keep the ticker independent of requests. A normal tick consumes +already queued hints before scanning; simultaneous tick/hint readiness is one +normal scan. Preserve hints arriving during a scan, the allocation cursor and all +ownership checks. Never close the hint channel while handlers may still send. +This bounds extra maintenance work but does not bypass capacity, a busy lifecycle +gate or multi-page scheduling, and does not guarantee a resume deadline. + +A recovered or uncertain running installation fails and uses existing cleanup, without replaying writes. Completed environments never reinstall initial files on reconnect or native recovery. Provider RunCommand carries bounded stdin, not confidential argv. Only fixed trusted initializers may run with Runtime authority. User setup and package install hooks diff --git a/services/agents-api/internal/execution/environment_admission.go b/services/agents-api/internal/execution/environment_admission.go index eadda58d7..2a712aca3 100644 --- a/services/agents-api/internal/execution/environment_admission.go +++ b/services/agents-api/internal/execution/environment_admission.go @@ -86,6 +86,9 @@ func (w *Worker) submitEnvironmentInputs(ctx context.Context, session store.Sess if err != nil { return nil, err } + if reservation.State == store.EnvironmentInputPending && !reservation.IsInitial { + w.hintRuntimeWake(ctx, session) + } ticker := time.NewTicker(250 * time.Millisecond) defer ticker.Stop() for { diff --git a/services/agents-api/internal/execution/runtime_lifecycle.go b/services/agents-api/internal/execution/runtime_lifecycle.go index 6b81261aa..e098ee640 100644 --- a/services/agents-api/internal/execution/runtime_lifecycle.go +++ b/services/agents-api/internal/execution/runtime_lifecycle.go @@ -39,6 +39,7 @@ type runtimeLifecycle struct { pendingCursor string connections map[string]*runtimeConnection initializing *runtimeInitialization + wakeHints chan struct{} } func newRuntimeLifecycle(s *store.Store, registry *gateway.Registry, config *RuntimeProvider) (*runtimeLifecycle, error) { @@ -63,7 +64,7 @@ func newRuntimeLifecycle(s *store.Store, registry *gateway.Registry, config *Run copied.Suspension = &policy } ctx, stop := context.WithCancel(context.Background()) - return &runtimeLifecycle{store: s, registry: registry, config: copied, gate: make(chan struct{}, 1), ctx: ctx, stop: stop, connections: make(map[string]*runtimeConnection)}, nil + return &runtimeLifecycle{store: s, registry: registry, config: copied, gate: make(chan struct{}, 1), ctx: ctx, stop: stop, connections: make(map[string]*runtimeConnection), wakeHints: make(chan struct{}, 1)}, nil } func (r *runtimeLifecycle) lock(ctx context.Context) error { @@ -331,14 +332,5 @@ func runtimeReference(owner store.RuntimeAllocation) sandbox.Reference { func (w *Worker) runManagedRuntimes(ctx context.Context) error { ticker := time.NewTicker(5 * time.Second) defer ticker.Stop() - for { - if err := w.ReconcileManagedRuntimes(ctx); err != nil { - return err - } - select { - case <-ctx.Done(): - return ctx.Err() - case <-ticker.C: - } - } + return runRuntimeMaintenance(ctx, ticker.C, w.runtimes.wakeHints, w.ReconcileManagedRuntimes) } diff --git a/services/agents-api/internal/execution/runtime_wake_hint.go b/services/agents-api/internal/execution/runtime_wake_hint.go new file mode 100644 index 000000000..364a5524b --- /dev/null +++ b/services/agents-api/internal/execution/runtime_wake_hint.go @@ -0,0 +1,87 @@ +package execution + +import ( + "context" + "time" + + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" +) + +// A hint only accelerates observation of an already committed input. Lookup or +// delivery failure leaves that input for the normal maintenance scan. +func (w *Worker) hintRuntimeWake(ctx context.Context, session store.Session) { + r := w.runtimes + if r == nil || r.config.Suspension == nil || r.ctx.Err() != nil { + return + } + lookup, cancel := context.WithTimeout(ctx, time.Second) + defer cancel() + environment, err := w.admission.GetSessionEnvironment(lookup, session.TenantID, session.ID) + if err != nil { + return + } + owner, err := w.admission.GetRuntimeAllocation(lookup, session.TenantID, environment.ID) + if err != nil || owner.ProviderKey != r.config.InstallationID || owner.State != "running" || + !owner.CreateSettled || owner.Initialization != "complete" || owner.SessionDeleted || owner.Expired { + return + } + switch owner.ComputePhase { + case "quiescing", "suspending", "suspended", "restoring", "waking": + select { + case r.wakeHints <- struct{}{}: + default: + } + } +} + +// Keep the ordinary ticker as the recovery guarantee. Hints allow at most one +// extra scan per normal cycle, including when clients repeatedly retry an input. +func runRuntimeMaintenance(ctx context.Context, ticks <-chan time.Time, hints <-chan struct{}, reconcile func(context.Context) error) error { + drainRuntimeWakeHint(hints) + if err := ctx.Err(); err != nil { + return err + } + if err := reconcile(ctx); err != nil { + return err + } + extraAllowed := true + for { + readyHints := hints + if !extraAllowed { + readyHints = nil + } + normal := false + select { + case <-ctx.Done(): + return ctx.Err() + case <-ticks: + normal = true + case <-readyHints: + // A simultaneously due tick owns this scan; do not scan twice. + select { + case <-ticks: + normal = true + default: + } + } + extraAllowed = normal + if normal { + // Pending hints are covered by the scan about to start. Hints that + // arrive during it stay queued for the new cycle's extra scan. + drainRuntimeWakeHint(hints) + } + if err := ctx.Err(); err != nil { + return err + } + if err := reconcile(ctx); err != nil { + return err + } + } +} + +func drainRuntimeWakeHint(hints <-chan struct{}) { + select { + case <-hints: + default: + } +} diff --git a/services/agents-api/internal/execution/runtime_wake_hint_test.go b/services/agents-api/internal/execution/runtime_wake_hint_test.go new file mode 100644 index 000000000..955afc33c --- /dev/null +++ b/services/agents-api/internal/execution/runtime_wake_hint_test.go @@ -0,0 +1,357 @@ +package execution + +import ( + "context" + "errors" + "runtime" + "sync" + "sync/atomic" + "testing" + "time" +) + +type maintenanceTestLoop struct { + t *testing.T + ticks chan time.Time + hints chan struct{} + entered chan int + release chan struct{} + done chan error + finished chan struct{} + cancel context.CancelFunc + calls atomic.Int32 +} + +func newMaintenanceTestLoop(t *testing.T, queuedHint bool, result func(int) error) *maintenanceTestLoop { + t.Helper() + ctx, cancel := context.WithCancel(t.Context()) + loop := &maintenanceTestLoop{ + t: t, ticks: make(chan time.Time, 1), hints: make(chan struct{}, 1), + entered: make(chan int, 1), release: make(chan struct{}), + done: make(chan error, 1), finished: make(chan struct{}), cancel: cancel, + } + if queuedHint { + loop.hints <- struct{}{} + } + go func() { + defer close(loop.finished) + loop.done <- runRuntimeMaintenance(ctx, loop.ticks, loop.hints, func(ctx context.Context) error { + call := int(loop.calls.Add(1)) + loop.entered <- call + select { + case <-ctx.Done(): + return ctx.Err() + case <-loop.release: + } + if result != nil { + return result(call) + } + return nil + }) + }() + t.Cleanup(func() { + cancel() + select { + case <-loop.finished: + case <-time.After(2 * time.Second): + t.Error("maintenance loop did not stop") + } + }) + return loop +} + +func (loop *maintenanceTestLoop) expectScan(want int) { + loop.t.Helper() + select { + case got := <-loop.entered: + if got != want { + loop.t.Fatalf("scan = %d, want %d", got, want) + } + case err := <-loop.done: + loop.t.Fatalf("maintenance stopped before scan %d: %v", want, err) + case <-time.After(2 * time.Second): + loop.t.Fatalf("maintenance did not enter scan %d", want) + } +} + +func (loop *maintenanceTestLoop) finishScan() { + loop.t.Helper() + select { + case loop.release <- struct{}{}: + case <-time.After(2 * time.Second): + loop.t.Fatal("maintenance did not release its current scan") + } +} + +func (loop *maintenanceTestLoop) expectIdle() { + loop.t.Helper() + // A short negative assertion catches immediate extra scans. Positive ordering + // is controlled by callback barriers and manual ticks, never wall-clock ticks. + select { + case got := <-loop.entered: + loop.t.Fatalf("unexpected scan %d without an eligible trigger", got) + case err := <-loop.done: + loop.t.Fatalf("maintenance stopped while idle: %v", err) + case <-time.After(25 * time.Millisecond): + } +} + +func (loop *maintenanceTestLoop) hintBurst(count int) int { + accepted := 0 + for range count { + select { + case loop.hints <- struct{}{}: + accepted++ + default: + } + } + return accepted +} + +func (loop *maintenanceTestLoop) tick() { + loop.t.Helper() + select { + case loop.ticks <- time.Now(): + default: + loop.t.Fatal("previous manual tick was not consumed") + } +} + +func (loop *maintenanceTestLoop) expectResult(want error) { + loop.t.Helper() + select { + case got := <-loop.done: + if !errors.Is(got, want) { + loop.t.Fatalf("maintenance error = %v, want %v", got, want) + } + case <-time.After(2 * time.Second): + loop.t.Fatal("maintenance did not return") + } +} + +func TestRuntimeMaintenanceIdleKeepsOnlyNormalScans(t *testing.T) { + loop := newMaintenanceTestLoop(t, false, nil) + loop.expectScan(1) + loop.finishScan() + loop.expectIdle() + for want := 2; want <= 4; want++ { + loop.tick() + loop.expectScan(want) + loop.finishScan() + loop.expectIdle() + } + loop.cancel() + loop.expectResult(context.Canceled) +} + +func TestRuntimeMaintenanceStartupAbsorbsQueuedHint(t *testing.T) { + loop := newMaintenanceTestLoop(t, true, nil) + loop.expectScan(1) + if len(loop.hints) != 0 { + t.Fatal("startup did not absorb the hint before scanning") + } + loop.finishScan() + loop.expectIdle() + if loop.hintBurst(1) != 1 { + t.Fatal("fresh hint was not queued") + } + loop.expectScan(2) + loop.finishScan() + loop.expectIdle() +} + +func TestRuntimeMaintenanceHintBurstHasOneExtraPerPeriod(t *testing.T) { + loop := newMaintenanceTestLoop(t, false, nil) + loop.expectScan(1) + for period := 0; period < 4; period++ { + // The normal scan is blocked while a burst arrives. Its hint must survive + // completion and authorize exactly one additional scan in this period. + if accepted := loop.hintBurst(10000); accepted != 1 { + t.Fatalf("burst accepted %d queued hints, want 1", accepted) + } + loop.finishScan() + loop.expectScan(2 + 2*period) + if accepted := loop.hintBurst(10000); accepted != 1 { + t.Fatalf("busy extra scan accepted %d queued hints, want 1", accepted) + } + loop.finishScan() + loop.expectIdle() + // Another burst after the extra scan cannot create another allowance. + if accepted := loop.hintBurst(10000); accepted != 0 { + t.Fatalf("spent-period hint was consumed early: accepted %d", accepted) + } + loop.expectIdle() + loop.tick() + loop.expectScan(3 + 2*period) + if len(loop.hints) != 0 { + t.Fatal("normal tick did not absorb the prior queued hint") + } + } + loop.finishScan() + loop.expectIdle() + if got := loop.calls.Load(); got != 9 { + t.Fatalf("four periods produced %d scans, want startup + 4 extra + 4 normal", got) + } +} + +func TestRuntimeMaintenanceSimultaneousTickAndHintUseOneNormalScan(t *testing.T) { + // Both triggers are ready before the callback barrier opens. Repeat to cover + // selection variability without making assertions depend on either outcome. + for attempt := 0; attempt < 20; attempt++ { + loop := newMaintenanceTestLoop(t, false, nil) + loop.expectScan(1) + loop.hintBurst(1) + loop.tick() + loop.finishScan() + loop.expectScan(2) + if len(loop.ticks) != 0 || len(loop.hints) != 0 { + t.Fatal("simultaneous triggers were not merged before the normal scan") + } + loop.finishScan() + loop.expectIdle() + // The merged scan was normal, so this period still permits one extra. + loop.hintBurst(1) + loop.expectScan(3) + loop.finishScan() + loop.cancel() + loop.expectResult(context.Canceled) + } +} + +func TestRuntimeMaintenancePreservesHintArrivingDuringNormalScan(t *testing.T) { + loop := newMaintenanceTestLoop(t, false, nil) + loop.expectScan(1) + loop.finishScan() + loop.expectIdle() + loop.tick() + loop.expectScan(2) + loop.hintBurst(1) + loop.finishScan() + loop.expectScan(3) + loop.finishScan() + loop.expectIdle() +} + +func TestRuntimeMaintenanceCancellationPrecedesNewScan(t *testing.T) { + t.Run("before startup", func(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + cancel() + ticks := make(chan time.Time, 1) + hints := make(chan struct{}, 1) + ticks <- time.Now() + hints <- struct{}{} + calls := 0 + err := runRuntimeMaintenance(ctx, ticks, hints, func(context.Context) error { + calls++ + return nil + }) + if !errors.Is(err, context.Canceled) || calls != 0 { + t.Fatalf("cancelled startup: calls=%d, error=%v", calls, err) + } + }) + t.Run("busy scan with queued triggers", func(t *testing.T) { + loop := newMaintenanceTestLoop(t, false, nil) + loop.expectScan(1) + loop.hintBurst(1) + loop.tick() + loop.cancel() + loop.expectResult(context.Canceled) + if got := loop.calls.Load(); got != 1 { + t.Fatalf("cancellation admitted another scan: %d", got) + } + }) +} + +func TestRuntimeMaintenanceReturnsReconcileError(t *testing.T) { + failure := errors.New("controlled reconciliation failure") + for _, failAt := range []int{1, 2} { + t.Run(map[int]string{1: "startup", 2: "hint"}[failAt], func(t *testing.T) { + loop := newMaintenanceTestLoop(t, false, func(call int) error { + if call == failAt { + return failure + } + return nil + }) + loop.expectScan(1) + if failAt == 2 { + loop.finishScan() + loop.hintBurst(1) + loop.expectScan(2) + } + // An already pending normal tick cannot retry a failed callback. + loop.tick() + loop.finishScan() + loop.expectResult(failure) + if got := loop.calls.Load(); got != int32(failAt) { + t.Fatalf("failure triggered another scan: %d", got) + } + }) + } +} + +func TestRuntimeMaintenanceContinuousHintStormKeepsPeriodLimit(t *testing.T) { + loop := newMaintenanceTestLoop(t, false, nil) + loop.expectScan(1) + stopStorm := make(chan struct{}) + stormDone := make(chan struct{}) + stormStarted := make(chan struct{}) + go func() { + defer close(stormDone) + loop.hintBurst(1) + close(stormStarted) + for { + select { + case <-stopStorm: + return + default: + } + loop.hintBurst(1) + runtime.Gosched() + } + }() + var stopOnce sync.Once + stop := func() { + stopOnce.Do(func() { close(stopStorm) }) + <-stormDone + } + t.Cleanup(stop) + <-stormStarted + for period := 0; period < 3; period++ { + loop.finishScan() + loop.expectScan(2 + 2*period) + loop.finishScan() + loop.expectIdle() + if period < 2 { + loop.tick() + loop.expectScan(3 + 2*period) + } + } + stop() + if got := loop.calls.Load(); got != 6 { + t.Fatalf("continuous hint storm produced %d scans, want 3 normal and 3 extra", got) + } + loop.cancel() + loop.expectResult(context.Canceled) +} + +func TestRuntimeMaintenanceCancellationAfterSuccessfulScanWinsReadyTriggers(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + ticks := make(chan time.Time, 1) + hints := make(chan struct{}, 1) + calls := 0 + unwantedScan := errors.New("scan started after cancellation") + err := runRuntimeMaintenance(ctx, ticks, hints, func(context.Context) error { + calls++ + if calls > 1 { + return unwantedScan + } + ticks <- time.Now() + hints <- struct{}{} + cancel() + return nil + }) + if !errors.Is(err, context.Canceled) || calls != 1 { + t.Fatalf("ready triggers bypassed cancellation: calls=%d, error=%v", calls, err) + } +} diff --git a/services/agents-api/internal/store/runtime_wake_hint_integration_test.go b/services/agents-api/internal/store/runtime_wake_hint_integration_test.go new file mode 100644 index 000000000..d64570177 --- /dev/null +++ b/services/agents-api/internal/store/runtime_wake_hint_integration_test.go @@ -0,0 +1,252 @@ +package store_test + +import ( + "context" + "encoding/json" + "errors" + "sync" + "sync/atomic" + "testing" + "time" + + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/execution" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/sandbox" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" +) + +type wakeHintScanProvider struct { + *fakeCheckpointProvider + sentinel string + release chan struct{} + scans chan int + count atomic.Int32 +} + +func (p *wakeHintScanProvider) GetCompute(ctx context.Context, reference sandbox.Reference, compute sandbox.Compute) (sandbox.ComputeState, error) { + if reference.AllocationID == p.sentinel { + call := int(p.count.Add(1)) + p.scans <- call + if call == 1 { + select { + case <-p.release: + case <-ctx.Done(): + return sandbox.ComputeState{}, ctx.Err() + } + } + } + return p.fakeCheckpointProvider.GetCompute(ctx, reference, compute) +} + +type wakeHintIntegrationTarget struct { + tenant string + session store.Session + environment store.Environment + owner store.RuntimeAllocation +} + +type wakeHintIntegration struct { + fixture *computeLifecycleFixture + provider *wakeHintScanProvider + worker *execution.Worker + target wakeHintIntegrationTarget + sentinel wakeHintIntegrationTarget + started time.Time + release func() +} + +func newWakeHintIntegration(t *testing.T) *wakeHintIntegration { + t.Helper() + f := newComputeLifecycleFixture(t, 2, 4) + create := func() wakeHintIntegrationTarget { + tenant, session, environment, owner := f.create() + return wakeHintIntegrationTarget{tenant, session, environment, owner} + } + target, sentinel := create(), create() + if target.owner.ID > sentinel.owner.ID { + target, sentinel = sentinel, target + } + f.complete(target.owner) + target.owner = f.phase(target.tenant, target.environment.ID, "suspended") + f.stop() + provider := &wakeHintScanProvider{ + fakeCheckpointProvider: f.provider, sentinel: sentinel.owner.ID, + release: make(chan struct{}), scans: make(chan int, 16), + } + worker, err := execution.StartWorker(t.Context(), &execution.Dispatcher{ + Store: f.store, Registry: f.provider.registry, + ManagedRuntimes: &execution.RuntimeProvider{ + CoreURL: "http://core.invalid/api/v1", InstallationID: f.key, + BackendFingerprint: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", + Provider: provider, Suspension: &f.policy, + }, + }) + if err != nil { + t.Fatal(err) + } + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + var once sync.Once + release := func() { once.Do(func() { close(provider.release) }) } + result := &wakeHintIntegration{f, provider, worker, target, sentinel, time.Now(), release} + go func() { done <- worker.Run(ctx) }() + t.Cleanup(func() { + cancel() + release() + select { + case err := <-done: + if err != nil && !errors.Is(err, context.Canceled) { + t.Error("maintenance worker stopped unexpectedly", err) + } + case <-time.After(3 * time.Second): + t.Error("maintenance worker did not stop") + } + }) + select { + case call := <-provider.scans: + if call != 1 { + t.Fatal("unexpected sentinel scan", call) + } + case <-time.After(3 * time.Second): + t.Fatal("initial scan did not reach the sentinel") + } + // UUID ordering places the suspended target before the blocked sentinel. + // No input exists yet, so the initial scan cannot have restored the target. + return result +} + +func wakeHintInput(text string) []store.Input { + payload, _ := json.Marshal(map[string]string{"text": text}) + return []store.Input{{Kind: "message", Payload: payload}} +} + +func (f *wakeHintIntegration) pending(t *testing.T, target wakeHintIntegrationTarget, key string) store.EnvironmentInputReservation { + t.Helper() + var id string + awaitDaemonRemoteCondition(t, t.Context(), 2*time.Second, "committed wake input", func() bool { + return f.fixture.pool.QueryRow(t.Context(), + "SELECT id::text FROM environment_input_reservations WHERE session_id=$1 AND idempotency_key=$2", + target.session.ID, key).Scan(&id) == nil + }) + pending, err := f.fixture.store.GetEnvironmentInputReservation(t.Context(), target.tenant, target.session.ID, id) + if err != nil || pending.State != store.EnvironmentInputPending { + t.Fatal("input was not durably pending", pending.State, err) + } + return pending +} + +func TestRuntimeWakeHintCommittedSubmitResumesBeforeNormalTick(t *testing.T) { + f := newWakeHintIntegration(t) + ctx, cancel := context.WithCancel(t.Context()) + done := make(chan error, 1) + go func() { + _, err := f.worker.SubmitInputs(ctx, f.target.tenant, f.target.session.ID, "wake", wakeHintInput("next turn")) + done <- err + }() + t.Cleanup(func() { + cancel() + select { + case err := <-done: + if !errors.Is(err, context.Canceled) { + t.Error("pending input waiter returned unexpectedly", err) + } + case <-time.After(2 * time.Second): + t.Error("pending input waiter did not stop") + } + }) + pending := f.pending(t, f.target, "wake") + // A same-key caller may stop waiting without removing the committed input. + // Its retry must retain the reservation and cannot restore a second target. + retryCtx, stopRetry := context.WithTimeout(t.Context(), 250*time.Millisecond) + _, err := f.worker.SubmitInputs(retryCtx, f.target.tenant, f.target.session.ID, "wake", wakeHintInput("next turn")) + stopRetry() + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatal("retry did not retain pending admission", err) + } + retried := f.pending(t, f.target, "wake") + if retried.ID != pending.ID || !retried.Deadline.Equal(pending.Deadline) { + t.Fatal("same-key retry replaced the reservation or its deadline") + } + f.release() + awaitDaemonRemoteCondition(t, t.Context(), 2*time.Second, "hint restored retained compute", func() bool { + owner, err := f.fixture.store.GetRuntimeAllocation(t.Context(), f.target.tenant, f.target.environment.ID) + return err == nil && owner.ComputePhase == "running" + }) + if elapsed := time.Since(f.started); elapsed >= 3*time.Second { + t.Fatalf("restore did not precede the normal five-second tick: %s", elapsed) + } + f.provider.mu.Lock() + restores, creates := f.provider.restores, f.provider.creates + f.provider.mu.Unlock() + if restores != 1 || creates != 2 || f.provider.promptFrames.Load() != 0 { + t.Fatal("wake replayed creation/restoration or sent native input", restores, creates, f.provider.promptFrames.Load()) + } + var reservations, turns int + if err := f.fixture.pool.QueryRow(t.Context(), + "SELECT (SELECT count(*) FROM environment_input_reservations WHERE session_id=$1), (SELECT count(*) FROM turns WHERE session_id=$1)", + f.target.session.ID).Scan(&reservations, &turns); err != nil || reservations != 1 || turns != 1 { + t.Fatal("retry duplicated input or started a Turn before preparation", reservations, turns, err) + } +} + +func TestRuntimeWakeHintRejectedSubmitDoesNotAccelerateScan(t *testing.T) { + for _, name := range []string{"invalid", "idempotency conflict", "competing batch"} { + t.Run(name, func(t *testing.T) { + f := newWakeHintIntegration(t) + key, inputs, want := "wake", wakeHintInput("next turn"), store.ErrInvalidInput + if name == "invalid" { + inputs = []store.Input{{Kind: "unsupported", Payload: json.RawMessage("{}")}} + } else { + // Persist directly while the sentinel is blocked. Only the failing + // Worker submission could emit a hint; Store persistence cannot. + if _, err := f.fixture.store.ReserveEnvironmentInput(t.Context(), f.target.tenant, f.target.session.ID, key, inputs); err != nil { + t.Fatal(err) + } + if name == "idempotency conflict" { + inputs, want = wakeHintInput("different input"), store.ErrIdempotencyConflict + } else { + key, want = "different-key", store.ErrTurnConflict + } + } + if _, err := f.worker.SubmitInputs(t.Context(), f.target.tenant, f.target.session.ID, key, inputs); !errors.Is(err, want) { + t.Fatal("unexpected rejected submission result", err, want) + } + f.release() + // The provider barrier fixes ordering; this short negative window + // verifies that rejected admission did not queue an immediate scan. + select { + case call := <-f.provider.scans: + t.Fatal("rejected submission accelerated lifecycle scan", call) + case <-time.After(250 * time.Millisecond): + } + f.provider.mu.Lock() + restores := f.provider.restores + f.provider.mu.Unlock() + if restores != 0 { + t.Fatal("rejected submission triggered restoration", restores) + } + }) + } +} + +func TestRuntimeWakeHintRunningSubmitDoesNotAccelerateScan(t *testing.T) { + f := newWakeHintIntegration(t) + // A completed Turn makes this a subsequent input, so the running compute + // filter must suppress the hint independently of the initial-input filter. + f.fixture.complete(f.sentinel.owner) + ctx, cancel := context.WithTimeout(t.Context(), 250*time.Millisecond) + _, err := f.worker.SubmitInputs(ctx, f.sentinel.tenant, f.sentinel.session.ID, "running", wakeHintInput("next turn")) + cancel() + if !errors.Is(err, context.DeadlineExceeded) { + t.Fatal("running input did not remain pending", err) + } + pending := f.pending(t, f.sentinel, "running") + if pending.IsInitial { + t.Fatal("running input unexpectedly exercised initial admission") + } + f.release() + select { + case call := <-f.provider.scans: + t.Fatal("running input accelerated lifecycle scan", call) + case <-time.After(250 * time.Millisecond): + } +} diff --git a/services/agents-api/tests/container_server.py b/services/agents-api/tests/container_server.py index 240900917..884e651ac 100755 --- a/services/agents-api/tests/container_server.py +++ b/services/agents-api/tests/container_server.py @@ -20,7 +20,7 @@ "--mount", f"type=bind,source={credential_key},target=/run/credential.key,readonly", "--env", "AGENTS_API_CREDENTIAL_KEY_FILE=/run/credential.key", ]) -for name in ("AGENTS_API_DATABASE_URL", "AGENTS_API_ADDR", "AGENTS_API_ENGINE"): +for name in ("AGENTS_API_DATABASE_URL", "AGENTS_API_ADDR", "AGENTS_API_ENGINE", "AGENTS_API_DAEMON_WS_URL"): args.extend(["--env", name]) args.append(os.environ["AGENTS_API_IMAGE"]) # Docker forwards termination to the API and --rm removes the stopped container.