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
18 changes: 16 additions & 2 deletions CONTRIBUTING.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
14 changes: 3 additions & 11 deletions services/agents-api/internal/execution/runtime_lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -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) {
Expand All @@ -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 {
Expand Down Expand Up @@ -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)
}
87 changes: 87 additions & 0 deletions services/agents-api/internal/execution/runtime_wake_hint.go
Original file line number Diff line number Diff line change
@@ -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:
}
}
Loading
Loading