diff --git a/CONTRIBUTING.md b/CONTRIBUTING.md index 0676274db..82f547fc3 100644 --- a/CONTRIBUTING.md +++ b/CONTRIBUTING.md @@ -1171,7 +1171,12 @@ The connection owner spans preparation and the transferred Run without a reserva deadline; every exit releases it. The existing Worker scans pending inputs using the same bounded scheduling slots, -Session locks, durable deadlines and engine capability checks. A self-hosted Session +Session locks, durable deadlines and engine capability checks. Keep the five-second +Environment-input scan cadence and at most 100 candidates per scan. At EOF after a +nonempty cursor, refill the first page once in the same scan; an empty queue must +not spin. Advance the cursor before readiness checks so an unavailable Runtime +cannot starve later candidates. Preserve the four active slots and alternation +between ordinary Turns and Environment inputs. A self-hosted Session waits for its dedicated enrolled device; it cannot select an arbitrary same-tenant device or migrate an existing binding. Preparation failure can retry while still pending without extending the deadline. Unknown promotion results or errors after diff --git a/services/agents-api/internal/execution/worker_schedule.go b/services/agents-api/internal/execution/worker_schedule.go index fce848474..51ca11939 100644 --- a/services/agents-api/internal/execution/worker_schedule.go +++ b/services/agents-api/internal/execution/worker_schedule.go @@ -35,8 +35,14 @@ func (s *workerSchedule) selectWork(ctx context.Context, w *Worker, devices []st return nil, err } s.nextEnvironmentScan = time.Now().Add(5 * time.Second) - if len(environments) == 0 { + if len(environments) == 0 && s.environmentCursor != "" { s.environmentCursor = "" + // Retry the first page now instead of spending a scan interval on EOF. + // A single refill preserves the candidate bound and cannot spin when empty. + environments, err = w.dispatcher.Store.ListEnvironmentInputWork(ctx, "", devices) + if err != nil { + return nil, err + } } } var selected []scheduledWork diff --git a/services/agents-api/internal/store/environment_worker_scan_test.go b/services/agents-api/internal/store/environment_worker_scan_test.go new file mode 100644 index 000000000..462120807 --- /dev/null +++ b/services/agents-api/internal/store/environment_worker_scan_test.go @@ -0,0 +1,92 @@ +package store_test + +import ( + "testing" + "time" + + "github.com/MiniMax-AI-Dev/parsar/internal/agentdaemon/proto" + "github.com/MiniMax-AI-Dev/parsar/services/agents-api/internal/store" +) + +func TestWorkerEnvironmentRetriesNewlyReadyAtNextScan(t *testing.T) { + h := newDispatchHarness(t) + enableWorkerEnvironment(t, h) + pending := workerEnvironmentReservation(t, h) + runtime := h.environments[pending.SessionID] + caps := workerEnvironmentCapabilities() + caps.Preparation = false + awaitFixtureCapabilities(t, runtime, caps) + frames := workerFrames(t, h, runtime) + + h.session = publicSession(t, h, "scan-barrier") + receipt := h.message("barrier", "ordinary work") + scanned := time.Now() + _, stop := startEnvironmentExpiryWorker(t, h.d) + // Dispatch starts only after selectWork has examined the pending input on + // the same pass, while its exact Runtime is still incapable of preparation. + barrier := nextWorkerFrame(t, frames, proto.TypePromptRequest) + if barrier.ID != receipt.TurnID { + t.Fatal("unexpected scan barrier") + } + awaitFixtureCapabilities(t, runtime, workerEnvironmentCapabilities()) + h.write(barrier.ID, proto.TypeDone, proto.DonePayload{Content: "complete"}) + waitTurn(t, h, barrier.ID, store.TurnCompleted) + + prepare := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) + if elapsed := time.Since(scanned); elapsed < 4*time.Second || elapsed > 8*time.Second { + t.Fatal("readiness retry must use the next scan, without an empty scan interval", elapsed) + } + if workerRuntimeForPreparation(t, h, prepare) != runtime { + t.Fatal("readiness retry moved Runtime ownership") + } + handle := acknowledgePreparation(runtime, prepare.ID) + runtime.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "failed"}) + nextWorkerFrame(t, frames, proto.TypeExecutionRelease) + stop() + stored, err := h.s.GetEnvironmentInputReservation(t.Context(), h.tenant, pending.SessionID, pending.ID) + if err != nil || stored.State != store.EnvironmentInputPending || !stored.Deadline.Equal(pending.Deadline) || len(stored.Receipts) != 0 { + t.Fatal("readiness retry changed pending identity or admitted work", stored, err) + } +} + +func TestWorkerEnvironmentPaginationReachesReadyTail(t *testing.T) { + h := newDispatchHarness(t) + enableWorkerEnvironment(t, h) + var last store.EnvironmentInputReservation + for range 101 { + pending := unboundWorkerEnvironmentReservation(t, h) + if pending.ID > last.ID { + last = pending + } + } + session, err := h.s.GetSession(t.Context(), h.tenant, last.SessionID) + if err != nil { + t.Fatal(err) + } + runtime := connectFixtureRuntime(t, h, session) + h.environments[last.SessionID] = runtime + frames := workerFrames(t, h, runtime) + h.session = publicSession(t, h, "page-barrier") + receipt := h.message("barrier", "ordinary work") + scanned := time.Now() + _, stop := startEnvironmentExpiryWorker(t, h.d) + barrier := nextWorkerFrame(t, frames, proto.TypePromptRequest) + if barrier.ID != receipt.TurnID { + t.Fatal("unexpected page barrier") + } + h.write(barrier.ID, proto.TypeDone, proto.DonePayload{Content: "complete"}) + waitTurn(t, h, barrier.ID, store.TurnCompleted) + // The first 100 unbound inputs must not pin the cursor, and the ready + // tail must wait for its own bounded page rather than an unbounded drain. + prepare := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) + if elapsed := time.Since(scanned); elapsed < 4*time.Second || elapsed > 8*time.Second { + t.Fatal("pagination lost the scan bound or starved the ready tail", elapsed) + } + if workerRuntimeForPreparation(t, h, prepare) != runtime { + t.Fatal("pagination selected the wrong Runtime") + } + handle := acknowledgePreparation(runtime, prepare.ID) + runtime.write(prepare.ID, proto.TypePreparationStatus, proto.PreparationStatusPayload{Handle: handle, Revision: 2, State: "failed"}) + nextWorkerFrame(t, frames, proto.TypeExecutionRelease) + stop() +} diff --git a/services/agents-api/internal/store/environment_worker_test.go b/services/agents-api/internal/store/environment_worker_test.go index 7b872a49f..2484ae9b9 100644 --- a/services/agents-api/internal/store/environment_worker_test.go +++ b/services/agents-api/internal/store/environment_worker_test.go @@ -147,8 +147,8 @@ func TestWorkerEnvironmentRetriesPendingWithoutExtendingDeadline(t *testing.T) { h.write(request.ID, proto.TypeDone, proto.DonePayload{Content: "complete"}) waitTurn(t, h, request.ID, store.TurnCompleted) second := nextWorkerFrame(t, frames, proto.TypeExecutionPrepare) - if time.Since(started) < 4*time.Second || first.ID == second.ID { - t.Fatal("pending preparation retried rapidly or reused a released owner") + if elapsed := time.Since(started); elapsed < 4*time.Second || elapsed > 8*time.Second || first.ID == second.ID { + t.Fatal("pending preparation missed its next scan or reused a released owner") } acknowledgePreparation(runtime, second.ID) select {