Skip to content
Closed
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
9 changes: 5 additions & 4 deletions contracts/agents-api/node-generation-protocol.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,20 +42,21 @@ Each operation carries its own arguments and returns the following result on suc
| `command` | `RunCommand` | `command` | `command` |
| `observe` | `Observe` | `observation` | `sample` |
| `initial` | `Initial` | None | `compute` |
| `new_compute` | `NewCompute` | Positive compute `generation` and optional `snapshot` | `compute` |
| `new_compute` | `NewCompute` | Positive `generation` and optional `retained` | `compute` |
| `compute` | `GetCompute` | `compute` | `state` |
| `renew_compute` | `RenewCompute` | Exact current `compute` | `state` |
| `kill_compute` | `KillCompute` | `compute` | None |
| `resume_compute` | `ResumeCompute` | `compute` | `state` |
| `command_compute` | `RunCommandCompute` | `compute` and `command` | `command` |
| `suspend` | `Suspend` | `suspend` | `state` |
| `resume` | `Resume` | `resume` | `state` |
| `delete_snapshot` | `DeleteSnapshot` | `snapshot` | None |
| `delete_retained` | `DeleteRetained` | `retained` | None |

A request whose `connection_id`, `owner_epoch` or `sequence` does not match closes the connection. A malformed request gets an `invalid` response. A node without generation management accepts only its enrolled `deployment_generation`; a generation-managing node runs the request on that generation's provider and answers `unconfirmed` when it cannot. Core sends `create` and a `resume` that is not observe-only only to a generation that is ready on that node, and keeps at most 32 requests pending per connection.
A request whose `connection_id`, `owner_epoch` or `sequence` does not match closes the connection. A malformed request gets an `invalid` response. A node without generation management accepts only its enrolled `deployment_generation`; a generation-managing node runs the request on that generation's provider and answers `unconfirmed` when it cannot. Core sends `create` and a `resume` that is not reconciliation-only only to a generation that is ready on that node, and keeps at most 32 requests pending per connection.

The budget is relative: the node anchors `timeout_ms` to its own clock on receipt and consumes it while the request waits in its queue, so the hosts' clocks need not agree. Core still bounds its own wait. A full node queue closes the connection.

The `response` frame carries `id` and `connection_id`. A successful response carries the result named in the operation table, with no result field for `kill`, `kill_compute` or `delete_snapshot`. A failed response carries an `error_code`:
The `response` frame carries `id` and `connection_id`. A successful response carries the result named in the operation table, with no result field for `kill`, `kill_compute` or `delete_retained`. A failed response carries an `error_code`:

| `error_code` | Meaning |
| --- | --- |
Expand Down
1 change: 1 addition & 0 deletions deploy/install/test_node_install.py
Original file line number Diff line number Diff line change
Expand Up @@ -645,6 +645,7 @@ def run_as(account, function, *arguments):
service_account=lambda: self.account),
mock.patch.object(installer.shutil, "which", side_effect=lambda tool: None if tool == "docker" and not self.docker_installed else "/usr/bin/" + tool),
mock.patch.object(installer.grp, "getgrnam", side_effect=lambda name: SimpleNamespace(gr_mem=["oac-node"] if self.joined else [])),
mock.patch.object(installer.os, "getgrouplist", side_effect=lambda name, gid: [gid]),
mock.patch.object(installer.grp, "getgrgid", side_effect=lambda gid: SimpleNamespace(gr_name=self.device_group))):
patch.start()
self.addCleanup(patch.stop)
Expand Down
22 changes: 13 additions & 9 deletions docs/sandbox-provider.md

Large diffs are not rendered by default.

42 changes: 29 additions & 13 deletions services/core/cmd/server/managed_generation_operations.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,13 +15,13 @@ func (p *generationRouter) Initial(ctx context.Context, r sandbox.Reference) (sa
return sandbox.Compute{}, err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return sandbox.Compute{}, err
}
return cp.Initial(ctx, r)
}
func (p *generationRouter) NewCompute(ctx context.Context, r sandbox.Reference, g uint64, snapshot *sandbox.SnapshotIdentity) (sandbox.Compute, error) {
func (p *generationRouter) NewCompute(ctx context.Context, r sandbox.Reference, g uint64, snapshot *sandbox.RetainedState) (sandbox.Compute, error) {
if err := providercontract.Require(p, "NewCompute"); err != nil {
return sandbox.Compute{}, err
}
Expand All @@ -30,7 +30,7 @@ func (p *generationRouter) NewCompute(ctx context.Context, r sandbox.Reference,
return sandbox.Compute{}, err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return sandbox.Compute{}, err
}
Expand All @@ -45,7 +45,7 @@ func (p *generationRouter) GetCompute(ctx context.Context, r sandbox.Reference,
return sandbox.ComputeState{}, err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return sandbox.ComputeState{}, err
}
Expand All @@ -60,7 +60,7 @@ func (p *generationRouter) Suspend(ctx context.Context, q sandbox.SuspendRequest
return sandbox.ComputeState{}, err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return sandbox.ComputeState{}, err
}
Expand All @@ -75,7 +75,7 @@ func (p *generationRouter) Resume(ctx context.Context, q sandbox.ResumeRequest)
return sandbox.ComputeState{}, err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return sandbox.ComputeState{}, err
}
Expand All @@ -90,26 +90,26 @@ func (p *generationRouter) KillCompute(ctx context.Context, r sandbox.Reference,
return err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return err
}
return cp.KillCompute(ctx, r, c)
}
func (p *generationRouter) DeleteSnapshot(ctx context.Context, r sandbox.Reference, snapshot sandbox.SnapshotIdentity) error {
if err := providercontract.Require(p, "DeleteSnapshot"); err != nil {
func (p *generationRouter) DeleteRetained(ctx context.Context, r sandbox.Reference, snapshot sandbox.RetainedState) error {
if err := providercontract.Require(p, "DeleteRetained"); err != nil {
return err
}
v, done, err := p.route(ctx, r)
if err != nil {
return err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return err
}
return cp.DeleteSnapshot(ctx, r, snapshot)
return cp.DeleteRetained(ctx, r, snapshot)
}
func (p *generationRouter) RunCommandCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute, command sandbox.Command) (sandbox.CommandResult, error) {
if err := providercontract.Require(p, "RunCommandCompute"); err != nil {
Expand All @@ -120,7 +120,7 @@ func (p *generationRouter) RunCommandCompute(ctx context.Context, r sandbox.Refe
return sandbox.CommandResult{}, err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return sandbox.CommandResult{}, err
}
Expand All @@ -135,7 +135,7 @@ func (p *generationRouter) ResumeCompute(ctx context.Context, r sandbox.Referenc
return sandbox.ComputeState{}, err
}
defer done()
cp, err := sandbox.Checkpoint(v)
cp, err := sandbox.Suspension(v)
if err != nil {
return sandbox.ComputeState{}, err
}
Expand All @@ -147,3 +147,19 @@ func (*generationRouter) DiscoverSelection(context.Context, sandbox.Selection) (
func (*generationRouter) VerifyCredential(context.Context, []sandbox.Reference) error {
return &providercontract.UnsupportedError{Operation: "VerifyCredential", Reason: "generation_router_does_not_verify_configuration"}
}

func (p *generationRouter) RenewCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) (sandbox.ComputeState, error) {
if err := providercontract.Require(p, "RenewCompute"); err != nil {
return sandbox.ComputeState{}, err
}
v, done, err := p.route(ctx, r)
if err != nil {
return sandbox.ComputeState{}, err
}
defer done()
cp, err := sandbox.Suspension(v)
if err != nil {
return sandbox.ComputeState{}, err
}
return cp.RenewCompute(ctx, r, c)
}
2 changes: 1 addition & 1 deletion services/core/cmd/server/managed_setup.go
Original file line number Diff line number Diff line change
Expand Up @@ -126,7 +126,7 @@ func (s *managedSetup) configuration(setup store.SandboxSetup) (execution.Prepar
}
selected := &execution.RuntimeProvider{InstallationID: setup.InstallationID, ProviderKind: setup.Provider, Generation: setup.Generation, Mode: setup.Mode, AdmissionPaused: setup.AdmissionPaused,
CoreURL: s.publicURL + "/api/v1", BackendFingerprint: setup.BackendFingerprint, Provider: provider}
if sandbox.SupportsCheckpoint(provider) {
if sandbox.SupportsSuspension(provider) {
selected.Suspension = &execution.RuntimeSuspensionPolicy{IdleTimeout: time.Duration(setup.IdleSeconds) * time.Second,
Retention: time.Duration(setup.RetentionSeconds) * time.Second, MaxActive: 4, MaxRetained: 16}
}
Expand Down
4 changes: 2 additions & 2 deletions services/core/internal/db/queries/runtime_allocations.sql
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
-- name: CreateRuntimeAllocation :one
INSERT INTO runtime_allocations (id, environment_id, device_id, provider_key, node_id, deployment_generation)
VALUES ($1, $2, $3, $4, $5, $6) RETURNING *;
INSERT INTO runtime_allocations (id, environment_id, device_id, provider_key, node_id, deployment_generation, compute_state)
VALUES ($1, $2, $3, $4, $5, $6, jsonb_build_object('protocol_version', sqlc.arg(protocol_version)::text)) RETURNING *;

-- name: GetRuntimeAllocation :one
SELECT sqlc.embed(a), e.session_id, s.tenant_id, s.deleted_at, (CASE WHEN a.compute_phase NOT IN ('disabled', 'running') THEN a.compute_retained_until IS NOT NULL AND a.compute_retained_until <= clock_timestamp() ELSE a.node_id IS NULL AND (SELECT mode FROM runtime_deployment) <> 'direct' AND a.kept_at <= clock_timestamp() - interval '1 hour' END)::boolean AS expired
Expand Down
11 changes: 8 additions & 3 deletions services/core/internal/db/queries/runtime_suspension.sql
Original file line number Diff line number Diff line change
Expand Up @@ -27,8 +27,6 @@ WHERE id = $1 AND compute_phase = 'running' AND compute_activity_at <= $2;
-- name: GetRuntimeActivity :one
SELECT clock_timestamp()::timestamptz AS observed_at,
GREATEST(a.compute_activity_at,
CASE WHEN a.node_id IS NULL THEN COALESCE((SELECT max(t.completed_at) FROM turns t WHERE t.session_id = e.session_id), a.created_at) END,
CASE WHEN a.node_id IS NULL THEN (SELECT max(t.completed_at) FROM subagent_turns t WHERE t.session_id = e.session_id) END,
(SELECT max(f.settled_at) FROM environment_file_writes f WHERE f.environment_id = e.id))::timestamptz AS last_activity,
(EXISTS (SELECT 1 FROM turns t WHERE t.session_id = e.session_id AND t.status IN ('queued','in_progress','waiting'))
OR EXISTS (SELECT 1 FROM subagent_turns t WHERE t.session_id = e.session_id AND t.status IN ('queued','in_progress','waiting'))
Expand Down Expand Up @@ -57,10 +55,17 @@ SELECT EXISTS (
UPDATE runtime_allocations a SET compute_activity_at = clock_timestamp()
FROM environments e
WHERE a.environment_id = e.id AND e.session_id = $1
AND a.node_id IS NOT NULL AND a.state = 'running';
AND a.state = 'running';

-- name: SessionHasRuntimeNode :one
SELECT EXISTS (
SELECT 1 FROM runtime_allocations a JOIN environments e ON e.id = a.environment_id
WHERE e.session_id = $1 AND a.node_id IS NOT NULL
)::boolean;

-- name: HasIncompatibleRuntimeComputeState :one
SELECT EXISTS (
SELECT 1 FROM runtime_allocations
WHERE state <> 'released'
AND (compute_state->>'protocol_version') IS DISTINCT FROM sqlc.arg(protocol_version)::text
)::boolean;
6 changes: 4 additions & 2 deletions services/core/internal/db/sqlc/runtime_allocations.sql.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

19 changes: 16 additions & 3 deletions services/core/internal/db/sqlc/runtime_suspension.sql.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ func (waitingCleanupCheckpoint) ProviderOperations() providercontract.Operations
}

type waitingCleanupCheckpoint struct {
sandbox.CheckpointProvider
sandbox.SuspensionProvider
beforeKill func()
}

Expand Down Expand Up @@ -106,7 +106,7 @@ func TestArchiveWaitingCleanupReceiptBarrier(t *testing.T) {
}
currentCompute := sandbox.Compute{ID: uuid.NewString(), Name: owner.ID + "-g0"}
if checkpoint {
state, _ := json.Marshal(runtimeCompute{Current: currentCompute})
state, _ := json.Marshal(runtimeCompute{Version: sandbox.SuspensionStateVersion, Current: currentCompute})
owner, err = writer.SetRuntimeCompute(t.Context(), owner, "running", state, nil, 0)
if err != nil {
t.Fatal(err)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,11 +16,12 @@ func (*lifecycleOnlySandbox) ProviderOperations() providercontract.Operations {
"RunCommand": {State: providercontract.Supported},
"Initial": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"NewCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"RenewCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"GetCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"Suspend": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"Resume": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"KillCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"DeleteSnapshot": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"DeleteRetained": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"RunCommandCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"ResumeCompute": {State: providercontract.Unsupported, Reason: "fixture_operation_not_supported"},
"ObservationProviderType": {State: providercontract.Supported},
Expand All @@ -34,7 +35,7 @@ func (*lifecycleOnlySandbox) ProviderOperations() providercontract.Operations {
func (*lifecycleOnlySandbox) Initial(context.Context, sandbox.Reference) (sandbox.Compute, error) {
return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "Initial", Reason: "fixture_operation_not_supported"}
}
func (*lifecycleOnlySandbox) NewCompute(context.Context, sandbox.Reference, uint64, *sandbox.SnapshotIdentity) (sandbox.Compute, error) {
func (*lifecycleOnlySandbox) NewCompute(context.Context, sandbox.Reference, uint64, *sandbox.RetainedState) (sandbox.Compute, error) {
return sandbox.Compute{}, &providercontract.UnsupportedError{Operation: "NewCompute", Reason: "fixture_operation_not_supported"}
}
func (*lifecycleOnlySandbox) GetCompute(context.Context, sandbox.Reference, sandbox.Compute) (sandbox.ComputeState, error) {
Expand All @@ -49,8 +50,8 @@ func (*lifecycleOnlySandbox) Resume(context.Context, sandbox.ResumeRequest) (san
func (*lifecycleOnlySandbox) KillCompute(context.Context, sandbox.Reference, sandbox.Compute) error {
return &providercontract.UnsupportedError{Operation: "KillCompute", Reason: "fixture_operation_not_supported"}
}
func (*lifecycleOnlySandbox) DeleteSnapshot(context.Context, sandbox.Reference, sandbox.SnapshotIdentity) error {
return &providercontract.UnsupportedError{Operation: "DeleteSnapshot", Reason: "fixture_operation_not_supported"}
func (*lifecycleOnlySandbox) DeleteRetained(context.Context, sandbox.Reference, sandbox.RetainedState) error {
return &providercontract.UnsupportedError{Operation: "DeleteRetained", Reason: "fixture_operation_not_supported"}
}
func (*lifecycleOnlySandbox) RunCommandCompute(context.Context, sandbox.Reference, sandbox.Compute, sandbox.Command) (sandbox.CommandResult, error) {
return sandbox.CommandResult{}, &providercontract.UnsupportedError{Operation: "RunCommandCompute", Reason: "fixture_operation_not_supported"}
Expand All @@ -75,3 +76,7 @@ func (*lifecycleOnlySandbox) ObservationProviderType() string { return "fixture"
func (p *lifecycleOnlySandbox) ResolveObservationSource(context.Context) (runtimeobs.Source, error) {
return p, nil
}

func (p *lifecycleOnlySandbox) RenewCompute(ctx context.Context, r sandbox.Reference, c sandbox.Compute) (sandbox.ComputeState, error) {
return sandbox.ComputeState{}, &providercontract.UnsupportedError{Operation: "RenewCompute", Reason: "fixture_operation_not_supported"}
}
Loading
Loading