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
7 changes: 6 additions & 1 deletion CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 7 additions & 1 deletion services/agents-api/internal/execution/worker_schedule.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
92 changes: 92 additions & 0 deletions services/agents-api/internal/store/environment_worker_scan_test.go
Original file line number Diff line number Diff line change
@@ -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()
}
4 changes: 2 additions & 2 deletions services/agents-api/internal/store/environment_worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Loading