Skip to content
Open
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 @@ -44,20 +44,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
4 changes: 4 additions & 0 deletions docs/runtime-bootstrap.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,10 @@ The Runtime validates the input and owns authentication and connection. A succes

Self-hosted executors and operator-provisioned devices get their daemon identity in other ways; the [machine connection API](../contracts/agents-api/machine-api.md#credentials) lists every credential source. All of them enter the same Runtime execution loop.

## Hosted suspension control

The private hosted park/wake control-file path is authored as `SuspendControlFile` in `internal/runtimebootstrap/bootstrap.go`. Core recovery and native Go adapters read that value; the E2B helper contract generator projects it into the template builder. Managed startup prepares its private directory and supplies `OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE` to enable Runtime suspension. This is a packaged protocol setting. The shared Sandbox Provider registration owns idle and retention defaults; the adapter owns its native lease timeout.

## Verification

`go test ./internal/runtimebootstrap ./apps/daemon/internal/cli` covers the input contract, the exclusivity of credential sources and restart behavior. Provider tests verify delivery and file permissions without relying on the Runtime's private storage.
22 changes: 13 additions & 9 deletions docs/sandbox-provider.md

Large diffs are not rendered by default.

3 changes: 3 additions & 0 deletions internal/runtimebootstrap/bootstrap.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,9 @@ import (
"github.com/google/uuid"
)

// SuspendControlFile is the packaged private hosted Runtime park/wake location.
const SuspendControlFile = "/run/oac/daemon-suspend.json"

const Version = 1
const MaxBytes = 16 * 1024

Expand Down
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)
}
4 changes: 2 additions & 2 deletions services/core/cmd/server/managed_generations_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,7 @@ q=json.load(sys.stdin)
with (pathlib.Path(q['Config']['StateDir'])/'requests').open('a') as f: f.write(json.dumps(q)+'\n')
info=dict(q['Reference'],State='running',ProviderID='owned',CreateSettled=True)
if q['Operation']=='kill': info['State']='absent'
print(json.dumps({'Version':1,'Info':info}))
print(json.dumps({'Version':q['Version'],'Info':info}))
`
if err := os.WriteFile(helper, []byte(script), 0700); err != nil {
t.Fatal(err)
Expand Down Expand Up @@ -132,7 +132,7 @@ code = ''
if k == 'revoked': code = 'unauthorized'
elif k == 'unconfirmed': code = 'unconfirmed'
elif not (k.startswith('team-a') and template in ('owned-a', 'new-a') or k == 'team-b' and template == 'public-b'): code = 'team_mismatch'
result = {'Version': 1, 'ErrorCode': code}
result = {'Version': q['Version'], 'ErrorCode': code}
if not code:
result['DeploymentValid'] = True
if q['Operation'] == 'validate_deployment':
Expand Down
11 changes: 6 additions & 5 deletions services/core/cmd/server/managed_setup_preflight_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@ import (
"errors"
"os"
"path/filepath"
"strconv"
"strings"
"testing"
"time"
Expand All @@ -19,7 +20,7 @@ import (

func TestE2BRejectedSpecificationHasSafeActionableDiagnostic(t *testing.T) {
helper := filepath.Join(t.TempDir(), "provider")
if err := os.WriteFile(helper, []byte("#!/bin/sh\ncat >/dev/null\nprintf '%s' '{\"Version\":1,\"ErrorCode\":\"invalid\"}'\n"), 0700); err != nil {
if err := os.WriteFile(helper, []byte("#!/bin/sh\ncat >/dev/null\nprintf '%s' '{\"Version\":"+strconv.Itoa(e2b.ProtocolVersion)+",\"ErrorCode\":\"invalid\"}'\n"), 0700); err != nil {
t.Fatal(err)
}
state := filepath.Join(t.TempDir(), "e2b")
Expand Down Expand Up @@ -49,12 +50,12 @@ op = q['Operation']
with (pathlib.Path(q['Config']['StateDir']) / 'operations').open('a') as log:
log.write(op + '\n')
if op == 'validate_deployment':
print(json.dumps({'Version': 1, 'ErrorCode': 'invalid'}))
print(json.dumps({'Version': q['Version'], 'ErrorCode': 'invalid'}))
else:
info = dict(q['Reference'], ProviderID='owned-compute', State='stopped', CreateSettled=True)
if op == 'kill':
info.update(ProviderID='', State='absent')
print(json.dumps({'Version': 1, 'Info': info}))
print(json.dumps({'Version': q['Version'], 'Info': info}))
`
if err := os.WriteFile(helper, []byte(script), 0700); err != nil {
t.Fatal(err)
Expand Down Expand Up @@ -101,7 +102,7 @@ func TestE2BCandidateAdoptsTemplateBuildForOmittedResources(t *testing.T) {
}
helper := filepath.Join(t.TempDir(), "provider")
requests := filepath.Join(state, "requests")
script := "#!/bin/sh\ncat >>" + requests + "\necho >>" + requests + "\nprintf '%s' '{\"Version\":1,\"ErrorCode\":\"\",\"DeploymentValid\":true,\"TemplateBuild\":{\"Status\":\"ready\",\"CPUs\":4,\"MemoryMiB\":4096,\"RootDiskMiB\":24063}}'\n"
script := "#!/bin/sh\ncat >>" + requests + "\necho >>" + requests + "\nprintf '%s' '{\"Version\":" + strconv.Itoa(e2b.ProtocolVersion) + ",\"ErrorCode\":\"\",\"DeploymentValid\":true,\"TemplateBuild\":{\"Status\":\"ready\",\"CPUs\":4,\"MemoryMiB\":4096,\"RootDiskMiB\":24063}}'\n"
if err := os.WriteFile(helper, []byte(script), 0700); err != nil {
t.Fatal(err)
}
Expand Down Expand Up @@ -133,7 +134,7 @@ func TestE2BCandidateAdoptsTemplateBuildForOmittedResources(t *testing.T) {

func TestInitialE2BPublicTemplateOutsideTeamIsRejected(t *testing.T) {
helper := filepath.Join(t.TempDir(), "provider")
if err := os.WriteFile(helper, []byte("#!/bin/sh\ncat >/dev/null\nprintf '%s' '{\"Version\":1,\"ErrorCode\":\"team_mismatch\"}'\n"), 0700); err != nil {
if err := os.WriteFile(helper, []byte("#!/bin/sh\ncat >/dev/null\nprintf '%s' '{\"Version\":"+strconv.Itoa(e2b.ProtocolVersion)+",\"ErrorCode\":\"team_mismatch\"}'\n"), 0700); err != nil {
t.Fatal(err)
}
state := filepath.Join(t.TempDir(), "e2b")
Expand Down
2 changes: 2 additions & 0 deletions services/core/deploy/e2b/build-template.py
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@
import tempfile

from e2b import Template
from helper_contract_generated import SUSPEND_CONTROL_FILE

BASE = 'node:22.23.1-bookworm-slim@sha256:8607a9064d4a571140998ae9e52a3b3fcf9cff361d04642d5971e6cd76d39e27'
parser = argparse.ArgumentParser()
Expand All @@ -27,6 +28,7 @@
if value.startswith(('HOME=', 'OAC_')))
if environment.get('OAC_RUNTIME_WORKSPACE') != '/environment/workspace':
parser.error('Image does not use the colocated Runtime layout')
environment['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'] = SUSPEND_CONTROL_FILE
dev_home = Path(os.environ.get('OAC_DEV_HOME') or Path.home() / '.oac')
if not dev_home.is_absolute():
parser.error('OAC_DEV_HOME must be absolute')
Expand Down
3 changes: 3 additions & 0 deletions services/core/deploy/e2b/build_template_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -64,6 +64,9 @@ def build(instance, **kwargs):
self.assertIs(instance, template)
context = Path(factory.call_args.kwargs['file_context_path'])
self.assertEqual(stat.S_IMODE(context.stat().st_mode), 0o700)
from helper_contract_generated import SUSPEND_CONTROL_FILE
runtime_env = json.loads((context / 'runtime-env.json').read_text())
self.assertEqual(runtime_env['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'], SUSPEND_CONTROL_FILE)
with tarfile.open(context / 'runtime.tar.gz') as archive:
modes = {m.name: stat.S_IMODE(m.mode) for m in archive.getmembers()}
for parent in ('usr', 'usr/local', 'etc'):
Expand Down
13 changes: 9 additions & 4 deletions services/core/deploy/e2b/helper_contract_generated.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
# Code generated by contractgen; DO NOT EDIT.
"""Adapter-private wire declarations. No SDK or repository dependency."""
COMPUTE_FIELDS = ["Generation","Name","ID","RestoredFrom"]
ERROR_CODES = ["","invalid","ownership","exists","not_found","command_unconfirmed","unconfirmed","template_invalid","team_mismatch","unauthorized"]
MANAGED_BOOTSTRAP_FIELDS = ["TenantID","EnvironmentID","AllocationID","SessionID","DeviceID","NetworkAccess","AllowedDomains","InstallationID","RuntimeBootstrap"]
MANAGED_IDENTITY_FIELDS = ["TenantID","EnvironmentID","AllocationID","SessionID","DeviceID","InstallationID"]
Expand All @@ -10,9 +11,13 @@
MAX_REQUEST = 75497472
MAX_RESPONSE = 16777216
NETWORK_ACCESS = ["enabled","disabled","restricted"]
OPERATIONS = ["create","inspect","renew","kill","command","validate_deployment","observe","list_templates","list_builds","verify_credential"]
PROTOCOL_VERSION = 1
OPERATIONS = ["create","inspect","renew","kill","command","validate_deployment","observe","list_templates","list_builds","verify_credential","compute_info","compute_renew","suspend","resume","compute_kill","delete_retained","compute_command","resume_compute"]
PROTOCOL_VERSION = 2
REFERENCE_FIELDS = ["TenantID","EnvironmentID","AllocationID"]
REQUEST_FIELDS = ["Version","Operation","Config","Reference","References","Bootstrap","RuntimeBootstrap","Command","Deadline"]
RESPONSE_FIELDS = ["Version","Info","Command","ErrorCode","DeploymentValid","TemplateBuild","Templates","Builds","Observations"]
REQUEST_FIELDS = ["Version","Operation","Config","Reference","References","Bootstrap","RuntimeBootstrap","Command","Compute","Suspend","Resume","Retained","Deadline"]
RESPONSE_FIELDS = ["Version","State","Info","Command","ErrorCode","DeploymentValid","TemplateBuild","Templates","Builds","Observations"]
RESUME_FIELDS = ["Reference","OperationID","Retained","Target","ReconcileOnly"]
RETAINED_FIELDS = ["Reference","ID","Data","OperationID","SourceGeneration","SourceName","SourceID"]
SDK_VERSION = "2.51.0"
SUSPEND_CONTROL_FILE = "/run/oac/daemon-suspend.json"
SUSPEND_FIELDS = ["Reference","OperationID","Source","Retained","ReconcileOnly"]
2 changes: 2 additions & 0 deletions services/core/deploy/e2b/init.py
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,8 @@ def initialize():
# Claim before any side effect. An interrupted attempt must never start twice.
write_private(ROOT / 'launch.json', identity)
environment = prepare_runtime()
# Application-owned startup does not participate in Core-managed suspension.
environment.pop("OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE", None)
credential = PROFILE.parent / 'executor-key.json'
write_private(credential, payload['executor_key'], owner=1000)
source.unlink()
Expand Down
4 changes: 3 additions & 1 deletion services/core/deploy/e2b/init_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,8 @@ def test_launch_handoff_is_private_and_cannot_replay(self):
profile = Path(temporary) / 'private/default'
environment_file = Path(temporary) / 'image.json'
environment_file.write_text(json.dumps({'HOME': '/home/runtime',
'OAC_RUNTIME_HOME': '/home/runtime/.oac', 'OAC_RUNTIME_WORKSPACE': '/environment/workspace'}))
'OAC_RUNTIME_HOME': '/home/runtime/.oac', 'OAC_RUNTIME_WORKSPACE': '/environment/workspace',
'OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE': '/private/control/fixture.json'}))
(root / 'bootstrap.json').write_text(json.dumps(PAYLOAD))
real_chmod = Path.chmod

Expand All @@ -57,6 +58,7 @@ def chmod(path, mode):
self.assertEqual(argv[argv.index('--remote') + 1], PAYLOAD['remote_url'])
self.assertNotIn('test-private-key', repr(popen.call_args))
self.assertNotIn('OAC_RUNTIME_SESSION_ID', options['env'])
self.assertNotIn('OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE', options['env'])
self.assertEqual((options['user'], options['group'], options['extra_groups']), (1000, 1000, []))
self.assertEqual(options['umask'], 0o077)
key = profile.parent / 'executor-key.json'
Expand Down
7 changes: 7 additions & 0 deletions services/core/deploy/e2b/managed_init.py
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,13 @@ def initialize():
binding = identity(payload)
shared.write_private(root / 'managed-launch.json', binding)
environment = shared.prepare_runtime()
control_file = Path(environment['OAC_RUNTIME_DAEMON_SUSPEND_PID_FILE'])
if not control_file.is_absolute() or '..' in control_file.parts:
raise ValueError('Absolute private suspend control file required')
control_directory = control_file.parent
control_directory.mkdir(mode=0o700, parents=True, exist_ok=True)
os.chown(control_directory, 1000, 1000)
control_directory.chmod(0o700)
environment.update(OAC_RUNTIME_ENVIRONMENT_ID=payload['EnvironmentID'],
OAC_RUNTIME_SESSION_ID=payload['SessionID'],
OAC_RUNTIME_NETWORK_ACCESS=payload['NetworkAccess'],
Expand Down
Loading
Loading