Documentation
¶
Index ¶
- Variables
- func GenerateSandboxID() (string, error)
- func ImageRequiresLocalPlacement(req models.CreateSandboxRequest) bool
- func IngressInstalledVersion() uint64
- func NormalizeCreateFailover(req *models.CreateSandboxRequest) error
- func RecordCreateReservationState(state string)
- func RecordExpiredReservations(count int)
- func RecordFacadeIdempotencyConflict(scope string)
- func RecordFacadeIdempotencyReplay(scope string)
- func RecordRouteMiss(reasons ...string)
- func RedactClusterSecrets(req models.CreateSandboxRequest) models.CreateSandboxRequest
- func SetIngressRouteLag(fsmVersion uint64)
- func SnapshotNeedsPush(snapshot *models.SandboxSnapshot) bool
- type AutoImportConfig
- type AutoImportPendingStore
- type AutoImportReconcileStats
- type AutoImportReconciler
- type AutoImportRequest
- type AutoImportResult
- type AutoImportSpecMutator
- type AutoImportSpecResolver
- type AutoImporter
- type ImageDistributionProvider
- type PortEndpoint
- type RouteKind
- type RouteShape
- type Service
- func (s *Service) AttachCluster(c cluster.Client)
- func (s *Service) AttachSnapshotPusher(pusher *SnapshotPusher, reconciler *SnapshotPushReconciler)
- func (s *Service) Capacity() capacity.Snapshot
- func (s *Service) ClaimIdempotentRequest(ctx context.Context, scope, fingerprint string, now time.Time, ...) (*models.IdempotentRequestRecord, bool, error)
- func (s *Service) Cluster() cluster.Client
- func (s *Service) ClusterTopologyError() error
- func (s *Service) CompleteIdempotentRequest(ctx context.Context, scope, fingerprint, targetID string, now time.Time, ...) error
- func (s *Service) CreateSandbox(ctx context.Context, req models.CreateSandboxRequest) (*models.CreateSandboxResponse, error)
- func (s *Service) CreateSandboxWithID(ctx context.Context, req models.CreateSandboxRequest, id string) (*models.CreateSandboxResponse, error)
- func (s *Service) CreateSnapshot(ctx context.Context, sandboxID string, req models.CreateSandboxSnapshotRequest) (*models.SandboxSnapshot, error)
- func (s *Service) CreateSnapshotWithOwnership(ctx context.Context, sandboxID string, req models.CreateSandboxSnapshotRequest) (*models.SandboxSnapshot, bool, error)
- func (s *Service) DeleteClusterSecrets(ctx context.Context, sandboxID string) error
- func (s *Service) DeleteIdempotentRequest(ctx context.Context, scope, fingerprint string) error
- func (s *Service) DeleteSnapshot(ctx context.Context, idOrName string) error
- func (s *Service) DeleteSnapshotAlias(ctx context.Context, alias string) error
- func (s *Service) DestroySandbox(ctx context.Context, id string) error
- func (s *Service) EnsureClusterReady(ctx context.Context) error
- func (s *Service) EnsureLayer4Ready(ctx context.Context) error
- func (s *Service) EnsureNetstatsReady(ctx context.Context) error
- func (s *Service) EnsureSandboxAwakeForHTTP(ctx context.Context, id string) (*models.Sandbox, error)
- func (s *Service) ExposePort(ctx context.Context, id string, port int, protocol string) (models.ExposePortResponse, error)
- func (s *Service) ForceReconcileHTTPWakeShape(ctx context.Context) error
- func (s *Service) GetCompatState(ctx context.Context, sandboxID, facade string) (*models.SandboxCompatState, error)
- func (s *Service) GetIdempotentRequest(ctx context.Context, scope, fingerprint string) (*models.IdempotentRequestRecord, error)
- func (s *Service) GetNetworkUsage(ctx context.Context, id string) (*models.NetworkUsage, error)
- func (s *Service) GetSandbox(ctx context.Context, id string) (*models.Sandbox, error)
- func (s *Service) GetSnapshot(ctx context.Context, idOrName string) (*models.SandboxSnapshot, error)
- func (s *Service) GetSnapshotAlias(ctx context.Context, alias string) (*models.SnapshotAlias, error)
- func (s *Service) Health(ctx context.Context) (models.HealthStatus, error)
- func (s *Service) IsSandboxStarted(ctx context.Context, id string) (bool, error)
- func (s *Service) ListCompatState(ctx context.Context, facade string) (map[string]models.SandboxCompatState, error)
- func (s *Service) ListMounts(ctx context.Context, sandboxID string) ([]models.MountSpecRedacted, error)
- func (s *Service) ListSandboxes(ctx context.Context, tagFilter map[string]string) ([]*models.Sandbox, error)
- func (s *Service) ListSnapshotAliases(ctx context.Context, facade string) (map[string]models.SnapshotAlias, error)
- func (s *Service) ListSnapshots(ctx context.Context) ([]*models.SandboxSnapshot, error)
- func (s *Service) MarkImported(ctx context.Context, sandboxID string, registryRef string)
- func (s *Service) NormalizeCreateImageDistribution(ctx context.Context, req *models.CreateSandboxRequest) error
- func (s *Service) OpenClusterSecrets(ctx context.Context, redacted models.CreateSandboxRequest, ...) (models.CreateSandboxRequest, error)
- func (s *Service) OpenClusterSecretsForNode(ctx context.Context, redacted models.CreateSandboxRequest, ...) (out models.CreateSandboxRequest, err error)
- func (s *Service) PutClusterSecretsForRecipient(ctx context.Context, sandboxID string, req models.CreateSandboxRequest, ...) (cluster.PlacementSecrets, error)
- func (s *Service) Reconcile(ctx context.Context) error
- func (s *Service) ReconcileClusterIngress(ctx context.Context) error
- func (s *Service) ReconstructWakeArmedIfNeeded(ctx context.Context, sandbox *models.Sandbox)
- func (s *Service) RecreateSandbox(ctx context.Context, id string, spec models.CreateSandboxRequest, ...) error
- func (s *Service) RegisterSnapshot(ctx context.Context, snapshot *models.SandboxSnapshot) (*models.SandboxSnapshot, error)
- func (s *Service) ReplayClusterOwnership(ctx context.Context) (int, error)
- func (s *Service) ReplayReservations(ctx context.Context)
- func (s *Service) ResizeSandbox(ctx context.Context, id string, req models.ResizeSandboxRequest) (*models.Sandbox, error)
- func (s *Service) ResolveSandboxIDByName(ctx context.Context, name string) (string, error)
- func (s *Service) SealClusterSecrets(req models.CreateSandboxRequest) ([]byte, error)
- func (s *Service) SealClusterSecretsForRecipient(req models.CreateSandboxRequest, recipient string) ([]byte, error)
- func (s *Service) SetNetworkLimits(ctx context.Context, id string, bytesInLimit, bytesOutLimit int64) (*models.NetworkUsage, error)
- func (s *Service) SnapshotPushReconciler() *SnapshotPushReconciler
- func (s *Service) StartBuiltImageGC(ctx context.Context)
- func (s *Service) StartCaddyCoalescer(ctx context.Context)
- func (s *Service) StartClusterIngressReconcile(ctx context.Context)
- func (s *Service) StartEventMonitor(ctx context.Context)
- func (s *Service) StartL4WakeProxy(ctx context.Context) error
- func (s *Service) StartLifecycleSweep(ctx context.Context)
- func (s *Service) StartReconcileLoop(ctx context.Context)
- func (s *Service) StartSandbox(ctx context.Context, id string) (*models.Sandbox, error)
- func (s *Service) StopCaddyCoalescer()
- func (s *Service) StopSandbox(ctx context.Context, id string) (*models.Sandbox, error)
- func (s *Service) ToolboxTarget(ctx context.Context, id string) (ToolboxEndpoint, error)
- func (s *Service) TouchSandbox(ctx context.Context, id string) error
- func (s *Service) UnexposePort(ctx context.Context, id string, port int) error
- func (s *Service) UnsealClusterSecrets(redacted models.CreateSandboxRequest, sealed []byte) (models.CreateSandboxRequest, error)
- func (s *Service) UnsealClusterSecretsForNode(redacted models.CreateSandboxRequest, sealed []byte, nodeID string) (models.CreateSandboxRequest, error)
- func (s *Service) UnsealRegistry(sealed []byte) (*models.RegistryAuth, error)
- func (s *Service) UpdateLifecycle(ctx context.Context, id string, l models.Lifecycle) (*models.Sandbox, error)
- func (s *Service) UpdateTags(ctx context.Context, sandboxID string, tags map[string]string) error
- func (s *Service) UpsertCompatState(ctx context.Context, sandboxID, facade, stateJSON string) error
- func (s *Service) UpsertSnapshotAlias(ctx context.Context, alias models.SnapshotAlias) error
- func (s *Service) WakeAwareL4PortTarget(ctx context.Context, id string, port int) (string, error)
- func (s *Service) WakeAwarePortTarget(ctx context.Context, id string, port int) (PortEndpoint, error)
- func (s *Service) WakeAwareToolboxTarget(ctx context.Context, id string) (ToolboxEndpoint, error)
- type SnapshotPushConfig
- type SnapshotPushDocker
- type SnapshotPushReconcileStats
- type SnapshotPushReconciler
- type SnapshotPushResult
- type SnapshotPushStore
- type SnapshotPusher
- type ToolboxEndpoint
Constants ¶
This section is empty.
Variables ¶
var AutoImportDisabledError = errors.New("auto-import disabled")
AutoImportDisabledError is returned by callsites that want to distinguish "feature off" from "feature on, but failed". main() can switch on it to suppress alerts.
ErrPreferredHostPortUnavailable is returned by exposePort when a TCP replay supplied a specific preferredHostPort that's already reserved (cluster-wide or on this node) and so cannot be re-bound. The allocator deliberately does NOT silently fall through to a fresh random port: cluster-stable TCP endpoints are the entire point of B6 — clients addressing host:40123 must not be invisibly rerouted to host:55555 after a failover-recreate. Park is the policy; the FSM record (with the original HostPort) stays intact and the watcher / operator surfaces the parked state instead of mutating the contract behind the client's back.
var ErrSandboxManuallyStopped = errors.New("sandbox manually stopped (do not wake)")
ErrSandboxManuallyStopped is returned by the wake helper when the sandbox was stopped via the manual StopSandbox path. Surfaces 409 at the control-plane proxy so the caller knows to call StartSandbox explicitly instead of triggering an implicit resume.
var ErrWakeCircuitOpen = errors.New("wake circuit open")
ErrWakeCircuitOpen is returned when the per-sandbox wake breaker is tripped (D3). Surfaces 503 + Retry-After: 60 at the control-plane proxy so a sandbox stuck in a cold-start failure loop does not thrash StartSandbox on every request for the full open window.
Functions ¶
func GenerateSandboxID ¶ added in v0.2.1
GenerateSandboxID is the exported entry point for the cluster handler's reservation-first create path: the router (Node A) needs to mint a sandbox ID before opReserve so the chosen target (Node T) can accept the forward with X-Cluster-Create-ID and run CreateSandboxWithID against the same reservation. Routes through the package-private generateSandboxID so the format stays in lockstep with the local create path.
func ImageRequiresLocalPlacement ¶ added in v0.2.1
func ImageRequiresLocalPlacement(req models.CreateSandboxRequest) bool
func IngressInstalledVersion ¶ added in v0.2.1
func IngressInstalledVersion() uint64
IngressInstalledVersion returns the highest placement.Version this node's ingress reconciler has finished installing routes for. The /v1/cluster/ placements/{id} convergence-status read compares this against the per- placement Version to answer "has this node bound the cluster-stable TCP route yet?" without reaching into Caddy state. Zero means the reconciler has not yet run a successful pass (fresh boot) — treat as "not converged" for any non-zero placement.Version.
func NormalizeCreateFailover ¶ added in v0.2.2
func NormalizeCreateFailover(req *models.CreateSandboxRequest) error
func RecordCreateReservationState ¶ added in v0.2.1
func RecordCreateReservationState(state string)
func RecordExpiredReservations ¶ added in v0.2.1
func RecordExpiredReservations(count int)
func RecordFacadeIdempotencyConflict ¶ added in v0.2.1
func RecordFacadeIdempotencyConflict(scope string)
func RecordFacadeIdempotencyReplay ¶ added in v0.2.1
func RecordFacadeIdempotencyReplay(scope string)
func RecordRouteMiss ¶ added in v0.2.1
func RecordRouteMiss(reasons ...string)
RecordRouteMiss bumps the route-miss counter from the API layer. Exported so pkg/api/v1 can wire it in without exposing the expvar directly (keeps the metric name owned by this package). Service has no logical role here — this is a package-level counter the v1 wrap layer pokes when it observes the no-usable-URL case.
func RedactClusterSecrets ¶ added in v0.2.1
func RedactClusterSecrets(req models.CreateSandboxRequest) models.CreateSandboxRequest
RedactClusterSecrets returns a copy of req with credentials stripped — safe to replicate via raft. The Registry field's Server/Username are preserved (not secret) but Password is cleared; mount Credentials maps are dropped per-entry. Maps and slices that the caller might mutate are deep-copied so the original req is left untouched.
Lives next to SealClusterSecrets because the two are always called as a pair: seal returns the encrypted bag, redact returns the safe-to-replicate spec, and writing one without the other would either leak secrets (no redact) or lose them on failover (no seal).
func SetIngressRouteLag ¶ added in v0.2.1
func SetIngressRouteLag(fsmVersion uint64)
SetIngressRouteLag is the post-tick hook the reconciler calls with the FSM's current PlacementVersion. Lag is computed as max(0, fsmVersion - ingressPlacementVersionMax). Computed here (not at recordIngressReconcile) so callers that only have the FSM version can still publish the lag without needing the reconciler's maxVersion.
func SnapshotNeedsPush ¶ added in v0.2.4
func SnapshotNeedsPush(snapshot *models.SandboxSnapshot) bool
SnapshotNeedsPush reports whether a snapshot row is a candidate for background push. Only local_only snapshots benefit — anything already in a remote registry is already distributable, and pushing it would just double-store the image.
Types ¶
type AutoImportConfig ¶ added in v0.2.3
type AutoImportConfig struct {
// Enabled gates the whole feature. When false the post-pull path
// short-circuits and the auto-importer is never built.
Enabled bool
// HooksBaseURL is the root URL of the AOCR hooks service (no trailing
// slash, e.g. `https://aocr.aerol.ai`). The importer appends
// `/v1/internal/imports`.
HooksBaseURL string
// ClusterID identifies this sandboxd cluster to AOCR. Goes into the
// import request body and constrains the cluster PAT's scope.
ClusterID string
// ClusterPAT is the bearer token presented to the ImportAPI. Treat as
// secret; do not log. Sourced from a file path at startup, not env.
ClusterPAT string
// RetentionSuffix is appended to the target tag in the cluster
// namespace (e.g. `--idle-90d`). Empty means use the ImportAPI's
// server-side default (`--idle-90d`).
RetentionSuffix string
// RequestTimeout caps a single ImportAPI call. The mount-from-repo
// flow is normally sub-second (no byte transfer) but the timeout
// guards against an upstream that hangs on layer enumeration.
RequestTimeout time.Duration
}
AutoImportConfig wires the post-pull auto-import path to a specific AOCR hooks deployment. A zero value disables auto-import; the rest of the service runs exactly as before. main() is expected to validate that the PAT + cluster ID are set when Enabled=true.
See `plans/cluster-mirror-and-snapshot-distribution.md` F21 for the full motivation. Short version: the first time a private upstream image is pulled on a sandbox with `failover.policy: recreate`, sandboxd asks AOCR to re-mount the just-cached bytes under `cluster/<id>/_imported/<host>/<repo>`. Future recreates pull from the cluster namespace with a cluster PAT — surviving an upstream credential rotation that would otherwise break the F17 wrap-creds path.
func (AutoImportConfig) Validate ¶ added in v0.2.3
func (c AutoImportConfig) Validate() error
Validate enforces the consistency rules the rest of the service relies on. Returns nil for a disabled config so callers can validate unconditionally.
type AutoImportPendingStore ¶ added in v0.2.3
type AutoImportPendingStore interface {
ListAutoImportPendingIDs(ctx context.Context) ([]string, error)
Get(ctx context.Context, id string) (*models.Sandbox, error)
SetAutoImportPending(ctx context.Context, id string, pending bool) error
}
AutoImportPendingStore is the narrow store surface the reconciler needs. Pulled out as an interface so the reconciler can be tested without instantiating a full SQLite store — and so future store backends can satisfy it without growing the reconciler's dependency surface.
type AutoImportReconcileStats ¶ added in v0.2.3
AutoImportReconcileStats is what RunOnce reports back. Used by tests and by the metrics-collection wrapper in main.go.
type AutoImportReconciler ¶ added in v0.2.3
type AutoImportReconciler struct {
// contains filtered or unexported fields
}
AutoImportReconciler walks the rows the post-pull path flagged `auto_import_pending = true` and re-tries them. One ticker per node; fan-out inside a single tick is capped to keep a recovery storm from saturating the local Docker daemon or AOCR.
The reconciler does not own a goroutine — `RunOnce` is the unit of work, and a thin wrapper in main.go schedules it on a ticker. That keeps the type fully testable (table-driven, no time.Tick) and lets the operator trigger an immediate retry from a debug endpoint if they want.
func NewAutoImportReconciler ¶ added in v0.2.3
func NewAutoImportReconciler(importer *AutoImporter, store AutoImportPendingStore, specs AutoImportSpecResolver, logger *slog.Logger, maxInFlight int) *AutoImportReconciler
NewAutoImportReconciler builds the reconciler. Returns nil if the importer is nil (auto-import disabled) so callers can `if r == nil { return }` cheaply. maxInFlight clamps fan-out per tick; pass 0 for the safe default (4).
func (*AutoImportReconciler) RunOnce ¶ added in v0.2.3
func (r *AutoImportReconciler) RunOnce(ctx context.Context) (AutoImportReconcileStats, error)
RunOnce processes one sweep of the auto-import pending queue. Blocks until all per-sandbox attempts have returned (success, failure, or context cancel). The bounded semaphore caps in-flight HTTP calls; the AOCR side is happy with the load but the local Docker network stack is not infinitely parallel.
func (*AutoImportReconciler) SetSpecMutator ¶ added in v0.2.3
func (r *AutoImportReconciler) SetSpecMutator(m AutoImportSpecMutator)
SetSpecMutator installs the optional write-back hook. Calling again replaces the previous one; pass nil to disable. Separated from the constructor so existing call sites (and tests) keep working unchanged — the mutator is a strict upgrade, not a required dependency.
type AutoImportRequest ¶ added in v0.2.3
type AutoImportRequest struct {
UpstreamHost string
UpstreamRepo string
UpstreamTag string
UpstreamDigest string
}
AutoImportRequest is the per-sandbox import payload. Fields mirror what AOCR `ImportAPI` (hooks/src/controllers/ImportAPI.ts) expects on the wire; the JSON marshaller below pins the contract.
type AutoImportResult ¶ added in v0.2.3
AutoImportResult is what comes back. RegistryRef is the *cluster-side* pull reference the caller should write onto the sandbox spec (replacing the user's upstream ref for future recreates). AlreadyPresent=true means the destination tag already existed at the right digest — the importer treats that as success but the caller may want to suppress the spec mutation to avoid an unnecessary Raft write.
type AutoImportSpecMutator ¶ added in v0.2.3
type AutoImportSpecMutator interface {
MarkImported(ctx context.Context, sandboxID string, registryRef string)
}
AutoImportSpecMutator is the optional write-back hook. After a successful import, the reconciler calls MarkImported so the replicated spec carries the new cluster registry ref and `aocr_imported` distribution mode. Future failovers then pull from the cluster namespace with the cluster PAT, fully decoupled from the original upstream credential.
Implementations must be safe to call on a non-leader node: the leader- forwarding lives below this surface (cluster.UpsertSpec proposes through the Raft FSM regardless of caller node). Errors are absorbed inside the implementation — the reconciler always treats import success as a clean retryOutcome even if the mutator can't reach the leader, because the imported bytes are already in AOCR.
A nil mutator disables write-through cleanly; the reconciler logs the new ref and moves on. This is the standalone-mode default.
type AutoImportSpecResolver ¶ added in v0.2.3
type AutoImportSpecResolver interface {
GetSandboxSpec(sandboxID string) (*models.CreateSandboxRequest, bool)
}
AutoImportSpecResolver lets the reconciler look up the replicated CreateSandboxRequest that owns a given sandbox. The reconciler needs this because the spec carries the original upstream ref + the failover policy (the gate for whether import should happen at all), but the sandbox row only carries the runtime view. In standalone mode this is backed by an in-memory map; in cluster mode by the Raft FSM.
Returns nil, false when no replicated spec exists — the reconciler treats that as "drop the flag, nothing to retry" rather than as an error, since a sandbox without a spec cannot meaningfully be recreated elsewhere anyway.
type AutoImporter ¶ added in v0.2.3
type AutoImporter struct {
// contains filtered or unexported fields
}
AutoImporter is the leaf dependency that actually talks to AOCR. It holds no per-sandbox state; instantiate once at boot and reuse. Safe for concurrent use (it's just an http.Client wrapper).
func NewAutoImporter ¶ added in v0.2.3
func NewAutoImporter(cfg AutoImportConfig) (*AutoImporter, error)
NewAutoImporter builds an importer. Returns nil if cfg.Enabled is false so callers can do `if importer == nil { return }` on the post-pull path without an extra config check.
func (*AutoImporter) Endpoint ¶ added in v0.2.3
func (a *AutoImporter) Endpoint() string
Endpoint is exposed for log/debug use. Returns "" for a nil importer so callers can log it unconditionally.
func (*AutoImporter) Import ¶ added in v0.2.3
func (a *AutoImporter) Import(ctx context.Context, req AutoImportRequest) (AutoImportResult, error)
Import calls the AOCR ImportAPI once. Idempotent on (upstream_digest, cluster_id, target_tag): re-posting an already-imported sandbox returns AlreadyPresent=true without an error. Network/HTTP failures return a wrapped error so the caller can decide whether to log+flag or surface.
type ImageDistributionProvider ¶ added in v0.2.1
type ImageDistributionProvider interface {
ClassifyImage(ctx context.Context, image string) (models.ImageDistributionMetadata, error)
}
ImageDistributionProvider is the control-plane contract for image availability. Providers may verify digests or talk to an external registry cache, but the core daemon only needs the resulting placement metadata.
type PortEndpoint ¶ added in v0.2.4
type PortEndpoint struct {
URL string
}
PortEndpoint is the upstream the ingress proxy dials for a wake-aware exposed-port request. URL has the form http://{containerIP}:{port}.
type RouteKind ¶ added in v0.2.4
type RouteKind int
RouteKind is the protocol surface a callsite is choosing a shape for. HTTP and L4 (TCP/TLS) consult different bypass flags so each can be rolled out independently — Phase 1 ships HTTP behind SB_HTTP_WAKE_DIRECT_BYPASS_ENABLED; Phase 2 ships L4 behind SB_L4_WAKE_DIRECT_BYPASS_ENABLED. Within each kind the decision tree is identical (Status + WakeArmed + ContainerIP) — only the "is bypass enabled?" predicate differs.
type RouteShape ¶ added in v0.2.4
type RouteShape int
RouteShape is the published Caddy shape for a (sandbox, exposed-port). One shape is live at a time per exposure; transitions are driven by status changes (Start/Stop/Die/Reconcile), never by per-request work.
See plans/warm-direct-route-bypass.md D7/D8 — this enum + chooseRouteShape are the single source of truth so HTTP, TCP, TLS, reconcile, and event callsites cannot disagree about what shape should be live for a given sandbox row.
const ( // RouteShapeNone means no route should be published. Destroyed // sandboxes and disarmed-stopped serverless sandboxes fall through // to Caddy's fallback (404). RouteShapeNone RouteShape = iota // RouteShapeDirect publishes the upstream as ContainerIP:port. // Used for non-serverless sandboxes always, and for serverless // sandboxes only while warm (HTTPWakeDirectBypassEnabled=true). RouteShapeDirect // RouteShapeWake publishes the upstream as the sandboxd loopback // ingress proxy (or per-exposure unix socket for TLS). Used for // serverless sandboxes while Stopped+armed, and for serverless // sandboxes always when HTTPWakeDirectBypassEnabled=false. RouteShapeWake )
func (RouteShape) String ¶ added in v0.2.4
func (r RouteShape) String() string
type Service ¶
type Service struct {
// contains filtered or unexported fields
}
func (*Service) AttachCluster ¶ added in v0.2.1
AttachCluster swaps in a cluster.Client. Called from cmd/sandboxd/main after service.New when SB_ENABLE_CLUSTER=true. Idempotent.
func (*Service) AttachSnapshotPusher ¶ added in v0.2.4
func (s *Service) AttachSnapshotPusher(pusher *SnapshotPusher, reconciler *SnapshotPushReconciler)
AttachSnapshotPusher wires in the optional AOCR snapshot-push pipeline. pusher must be non-nil to activate the feature; a nil pusher is a no-op and leaves the service in legacy local-only snapshot mode. reconciler may be nil — when nil, kick-after-create is disabled (useful in tests that assert on the initial persisted state before any reconcile runs). Called once from main() after cfg.SnapshotPushEnabled validation.
func (*Service) Capacity ¶
Capacity returns the admitter's current snapshot. Returns the zero value when no admitter is configured (e.g. in tests).
func (*Service) ClaimIdempotentRequest ¶ added in v0.1.7
func (*Service) Cluster ¶ added in v0.2.1
Cluster returns the attached cluster.Client. Always non-nil.
func (*Service) ClusterTopologyError ¶ added in v0.2.2
ClusterTopologyError returns a production-topology violation for the current live member set. It is intentionally a runtime check so rolling membership, old nodes that still gossip empty roles, and explicit hybrid roles are all evaluated from the same source of truth the scheduler uses.
func (*Service) CompleteIdempotentRequest ¶ added in v0.1.7
func (*Service) CreateSandbox ¶
func (s *Service) CreateSandbox(ctx context.Context, req models.CreateSandboxRequest) (*models.CreateSandboxResponse, error)
func (*Service) CreateSandboxWithID ¶ added in v0.2.1
func (s *Service) CreateSandboxWithID(ctx context.Context, req models.CreateSandboxRequest, id string) (*models.CreateSandboxResponse, error)
CreateSandboxWithID is the failover-recreate entry point: it behaves like CreateSandbox but uses the supplied ID instead of generating a fresh one. Idempotent at the cluster boundary — if a sandbox with this ID already exists locally we return the existing record without touching docker. Used by the cluster owner watcher to re-materialize a sandbox after its previous owner died.
func (*Service) CreateSnapshot ¶ added in v0.1.7
func (s *Service) CreateSnapshot(ctx context.Context, sandboxID string, req models.CreateSandboxSnapshotRequest) (*models.SandboxSnapshot, error)
CreateSnapshot commits the sandbox container into a reusable local image. Idempotency is by snapshot name: repeated requests for the same sandbox + name return the stored snapshot metadata, while a different sandbox trying to claim the same name is rejected with a conflict.
func (*Service) CreateSnapshotWithOwnership ¶ added in v0.1.7
func (s *Service) CreateSnapshotWithOwnership(ctx context.Context, sandboxID string, req models.CreateSandboxSnapshotRequest) (*models.SandboxSnapshot, bool, error)
CreateSnapshotWithOwnership commits a sandbox image and reports whether this call created the native snapshot row. Callers that add companion metadata can use the flag to avoid rolling back a snapshot that already existed.
func (*Service) DeleteClusterSecrets ¶ added in v0.2.1
func (*Service) DeleteIdempotentRequest ¶ added in v0.1.7
func (*Service) DeleteSnapshot ¶ added in v0.1.7
func (*Service) DeleteSnapshotAlias ¶ added in v0.1.7
func (*Service) DestroySandbox ¶
func (*Service) EnsureClusterReady ¶ added in v0.2.1
EnsureClusterReady blocks until the cluster has elected a leader, mirroring the EnsureLayer4Ready single-flight latch shape. Single-node mode latches immediately. The API wrapper calls this before any RecordPlacement so a just-booted node doesn't 503 a CreateSandbox while raft is still catching up.
func (*Service) EnsureLayer4Ready ¶ added in v0.1.4
EnsureLayer4Ready bootstraps the caddy-l4 app under a single-flight mutex and latches success. Safe to call from boot AND from each L4 exposure path: the atomic fast-path turns it into a single load on the steady state, and a failed boot is recovered by the very next TCP/TLS expose call instead of surfacing as a confusing "layer4 app missing" error from caddy.
func (*Service) EnsureNetstatsReady ¶ added in v0.1.7
EnsureNetstatsReady boots the per-sandbox network byte-counter poller under a single-flight latch. Same lazy-bootstrap shape as EnsureLayer4Ready: caller pays the bootstrap cost only once across the daemon's lifetime, and the poller goroutine survives until ctx — typically the daemon's signal context — is cancelled.
Failure here is non-fatal: callers log and continue. Without the poller, network counters stay at zero and quotas never trigger; both Create and reads of /network/usage still work.
func (*Service) EnsureSandboxAwakeForHTTP ¶ added in v0.2.4
func (s *Service) EnsureSandboxAwakeForHTTP(ctx context.Context, id string) (*models.Sandbox, error)
EnsureSandboxAwakeForHTTP is the single chokepoint the control-plane HTTP proxies (v1 toolbox, daytona toolbox, e2b runtime) call before resolving the toolbox endpoint. It returns the current sandbox row so the caller can route to ContainerIP without re-fetching.
Behavior:
- sandbox started → fast-path return (no lock, no metric latency observation; only a requests_total tick).
- sandbox stopped + (gate off OR Lifecycle.Serverless=false OR WakeArmed=false) → ErrSandboxManuallyStopped. The proxy maps this to 409 so the SDK can explicitly call StartSandbox.
- sandbox stopped + eligible to wake → take the per-sandbox wake lock, recheck under the lock (a concurrent flight may have finished), enforce the circuit breaker, then call StartSandbox with a 15s budget. On admission failure or start error the breaker counter advances; 5 failures inside 60s trip it for 60s. Success resets the counter.
Same single-flight pattern as EnsureLayer4Ready, but per-sandbox: the wave of HTTP requests waking the same sandbox collapses to one StartSandbox; waves to different sandboxes proceed concurrently.
func (*Service) ExposePort ¶
func (s *Service) ExposePort(ctx context.Context, id string, port int, protocol string) (models.ExposePortResponse, error)
ExposePort publishes a sandbox container port through one of three caddy surfaces, selected by protocol:
- "" / "http": existing Caddy HTTP reverse-proxy route, returns https://<id>-<port>.<domain> (or the path-mode equivalent).
- "tcp": allocates a parent-host TCP port from the [SB_L4_PORT_RANGE_START, SB_L4_PORT_RANGE_END] pool, points caddy-l4 at it, and returns tcp://<public-host>:<host-port>. This is what unblocks native Postgres / Redis / MySQL DSNs in the spawn-postgres docs.
- "tls": adds a TLS-SNI route to the shared layer4 server. Requires --domain (so the SNI hostname has a place to resolve) and a non-empty SB_L4_TLS_LISTEN. Returns tls://<id>-<port>.<domain>:<l4-port>.
func (*Service) ForceReconcileHTTPWakeShape ¶ added in v0.2.4
ForceReconcileHTTPWakeShape republishes the HTTP route for every serverless sandbox's HTTP exposures using the current HTTPWakeDirectBypassEnabled value. Called exactly once at daemon startup when the rollback marker shows the previous run had bypass=true and the current run has bypass=false: without this, every serverless sandbox that was last written under bypass=true keeps its direct route pointing straight at the container IP, bypassing the wake-aware ingress proxy entirely — i.e. the rollback knob would not actually roll back. Conversely, on a fresh enable (false→true) routes will be rewritten lazily on the next Start/Stop, so no force-pass is needed in that direction (D5 of plans/warm-direct-route-bypass.md).
Idempotent: chooseRouteShape is a pure function of (cfg, sandbox), so repeated calls converge to the same route set. Best-effort per sandbox — one failed Caddy write logs and the loop continues so a single broken row does not block the rest from being rolled back.
func (*Service) GetCompatState ¶ added in v0.1.7
func (*Service) GetIdempotentRequest ¶ added in v0.1.7
func (*Service) GetNetworkUsage ¶ added in v0.1.7
GetNetworkUsage returns the current cumulative byte counters and configured limits for a sandbox. Callers handle ErrNotFound translation.
Best-effort lazy bootstrap of the netstats poller: if boot's EnsureNetstatsReady failed (cold-start race against the docker daemon, etc.) this call retries it under the same single-flight latch. Failure is logged and swallowed — the caller still gets back whatever counters the store has, and the next call will retry again.
func (*Service) GetSandbox ¶
func (*Service) GetSnapshot ¶ added in v0.1.7
func (*Service) GetSnapshotAlias ¶ added in v0.1.7
func (*Service) IsSandboxStarted ¶ added in v0.2.4
IsSandboxStarted is the preflight check the wake-aware ingress proxy uses to skip request-body buffering for warm sandboxes. It returns (true, nil) only when the sandbox row exists and its status is Started; any other status (Stopped, Destroyed, in-flight create) returns (false, nil). store.ErrNotFound is returned unwrapped so the proxy can map it to 404 without a second store hit.
Hot path: a short-TTL cache (warmCache) lets repeated warm requests to the same sandbox skip the SQLite read. The cache holds positives only ("seen Started in the last warmCacheTTL"); a miss always falls through to s.store.Get so cold/stopped/destroyed sandboxes never short-circuit. Lifecycle transitions (StartSandbox success, every stop path, every destroy path) explicitly invalidate the entry.
func (*Service) ListCompatState ¶ added in v0.1.7
func (*Service) ListMounts ¶
func (s *Service) ListMounts(ctx context.Context, sandboxID string) ([]models.MountSpecRedacted, error)
ListMounts returns the redacted mount config for a sandbox. Credentials are never included in the response — they are write-only via CreateSandbox.
func (*Service) ListSandboxes ¶
func (s *Service) ListSandboxes(ctx context.Context, tagFilter map[string]string) ([]*models.Sandbox, error)
ListSandboxes returns sandboxes whose Tags match every entry in tagFilter. A nil or empty filter returns every sandbox on this node. Filtering happens in-memory after the store read because Tags is JSON-encoded; pushing the filter into SQL via json_extract is a follow-up once row counts make the extra hop worth it. The filter exists so an external control plane can ask "give me the sandboxes belonging to user X" without round-tripping every sandbox in the cluster (see plans/multi-tenancy-via-control-plane.md).
func (*Service) ListSnapshotAliases ¶ added in v0.1.7
func (*Service) ListSnapshots ¶ added in v0.1.7
func (*Service) MarkImported ¶ added in v0.2.3
MarkImported is the spec write-back invoked by the AutoImportReconciler after a successful import. It flips the replicated spec's distribution mode to `aocr_imported` and points ImageRegistryRef at the new cluster- side ref so any future failover/recreate pulls from the cluster namespace (cluster PAT) instead of the original upstream.
Best-effort by contract: a missing spec, a non-leader UpsertSpec rejection, or a Raft commit error all log+return without surfacing the failure. The bytes are already in AOCR; the original ref still resolves through the mirror; the next mutation will refresh the FSM. This is the AutoImportSpecMutator implementation referenced by the reconciler.
func (*Service) NormalizeCreateImageDistribution ¶ added in v0.2.1
func (s *Service) NormalizeCreateImageDistribution(ctx context.Context, req *models.CreateSandboxRequest) error
NormalizeCreateImageDistribution resolves snapshot aliases and fills the create request's image-distribution metadata before placement decisions are made. The method is intentionally side-effect free: it does not pull, push, or verify image contents.
func (*Service) OpenClusterSecrets ¶ added in v0.2.1
func (s *Service) OpenClusterSecrets(ctx context.Context, redacted models.CreateSandboxRequest, secrets cluster.PlacementSecrets) (models.CreateSandboxRequest, error)
func (*Service) OpenClusterSecretsForNode ¶ added in v0.2.1
func (s *Service) OpenClusterSecretsForNode(ctx context.Context, redacted models.CreateSandboxRequest, secrets cluster.PlacementSecrets, nodeID string) (out models.CreateSandboxRequest, err error)
OpenClusterSecretsForNode resolves a replicated secret handle and merges the decrypted credentials back into a redacted spec. LegacySealed is still honored for placements written before the ref model.
func (*Service) PutClusterSecretsForRecipient ¶ added in v0.2.1
func (s *Service) PutClusterSecretsForRecipient(ctx context.Context, sandboxID string, req models.CreateSandboxRequest, recipient string) (cluster.PlacementSecrets, error)
PutClusterSecretsForRecipient stores the credential-bearing parts of req behind a provider ref and returns the handle safe to replicate through Raft. The local provider stores an encrypted recipient-bound envelope in SQLite; external KMS/secret-store providers can replace this boundary without changing cluster placement state.
func (*Service) ReconcileClusterIngress ¶ added in v0.2.1
func (*Service) ReconstructWakeArmedIfNeeded ¶ added in v0.2.4
ReconstructWakeArmedIfNeeded is the D1 reconstruction rule from plans/serverless-sandbox-http-wake.md, extended to also self-heal Caddy state. A serverless sandbox can land in `stopped + wake_armed=false` on a node that does not own the arming history — typically after a cluster owner change (the old owner's SQLite wake_armed bit doesn't migrate) or after a daemon restart that lost the in-memory expected-stop bookkeeping and saw a crash event arrive before the row got an armed value. In those cases the conservative default is the same as the involuntary-stop policy: arm wake so the next HTTP request resurrects the sandbox.
It ALSO re-upserts wake routes when wake_armed=true is already set. Caddy's admin-API routes are dynamic and do not survive a Caddy restart (only certs persist via S3 storage), so the daemon must treat wake_armed=true as a goal state to reassert against Caddy on every reconcile pass, not as a witness that routes still exist. Without this, a Caddy restart between auto-stop and the next wake request leaves the wildcard 404 fallback serving requests forever even though the DB looks correct.
Reconstruction is a no-op unless ALL of:
- cfg.EnableServerless is on
- sandbox is Lifecycle.Serverless
- sandbox is in Stopped status
- the sandbox has at least one exposed port that can receive an inbound wake request or connection
Caddy-first per D5: install wake routes for every exposure before flipping the store bit. If every wake route install fails we leave wake_armed unchanged (false stays false, true stays true) so route state and the bit remain in sync — the next reconcile pass will retry. The store write is skipped when wake_armed is already true so steady-state reconcile passes do not churn the row.
func (*Service) RecreateSandbox ¶ added in v0.2.1
func (s *Service) RecreateSandbox(ctx context.Context, id string, spec models.CreateSandboxRequest, secrets cluster.PlacementSecrets, exposedPorts map[int]cluster.ExposedPortRoute) error
RecreateSandbox satisfies cluster.SandboxRecreator. The cluster owner watcher invokes this for any FSM placement that points to self. If the sandbox already exists locally, we still replay the replicated port intents: a previous recreate attempt may have created the container and then failed while restoring Caddy/L4 ingress.
secrets is the provider handle that can rehydrate the redacted spec; we resolve and re-merge it here via OpenClusterSecretsForNode so the recreated container can pull from the same private registry / mount the same external storage. A decrypt failure is fatal to this attempt but non-fatal globally — the watcher's retry loop (now with reassign-after-K-failures) will eventually move the placement to a node whose key matches.
Port replay tries every port but returns an error when any replay failed so the owner watcher keeps retrying and can eventually reassign the placement. ExposePort is idempotent, so a partial replay is safe to resume.
Only placements whose create spec opted into failover.policy=recreate reach this path; default sandboxes remain non-HA and are orphaned on owner death.
func (*Service) RegisterSnapshot ¶ added in v0.1.7
func (s *Service) RegisterSnapshot(ctx context.Context, snapshot *models.SandboxSnapshot) (*models.SandboxSnapshot, error)
RegisterSnapshot persists a snapshot row whose Image was resolved out-of-band — either a pre-existing registry image the caller supplied by name, or a freshly built local tag produced by the image builder (e.g. the daytona facade's buildInfo path). It does NOT call docker.CreateSnapshot; the image is assumed to already be runnable. Idempotency is by snapshot name; a re-register with matching image is treated as a no-op so SDK retries don't fail. A different image under the same name is a conflict.
func (*Service) ReplayClusterOwnership ¶ added in v0.2.4
ReplayClusterOwnership pushes local sandbox truth into the cluster placement FSM. It is the local-store -> global-index repair path used at boot and by the periodic reconciler after Docker confirms a container still exists.
func (*Service) ReplayReservations ¶
ReplayReservations re-populates the admitter from persistent state. Without this, after a daemon restart the admitter sees zero reservations and the host can be overcommitted on the first wave of new sandboxes. Destroyed AND stopped sandboxes are skipped — neither holds host CPU/RAM (the stop path releases the slot, and StartSandbox re-Admits on the way back up), so counting them here would re-introduce the overcommit-budget bug we fixed when stop began releasing capacity. Best-effort: a store error is logged, not returned, since admission control degrading to "unaware" is preferable to refusing to boot.
func (*Service) ResizeSandbox ¶
func (*Service) ResolveSandboxIDByName ¶ added in v0.1.7
ResolveSandboxIDByName looks up the sandbox owning the given unique name. Empty name returns ErrNotFound (handled inside the store).
func (*Service) SealClusterSecrets ¶ added in v0.2.1
func (s *Service) SealClusterSecrets(req models.CreateSandboxRequest) ([]byte, error)
SealClusterSecrets extracts the secret-bearing portions of req, marshals them as JSON, and encrypts the result with the service cipher. The output is opaque bytes safe to put in the raft log. Returns nil/nil when there are no secrets to seal so the FSM column stays empty for sandboxes that don't need it.
The legacy method emits a wildcard-recipient v2 envelope for compatibility. New cluster placement paths should prefer SealClusterSecretsForRecipient so the encrypted payload is authenticated to the specific owner node ID.
func (*Service) SealClusterSecretsForRecipient ¶ added in v0.2.1
func (*Service) SetNetworkLimits ¶ added in v0.1.7
func (s *Service) SetNetworkLimits(ctx context.Context, id string, bytesInLimit, bytesOutLimit int64) (*models.NetworkUsage, error)
SetNetworkLimits writes new caps (0 = unlimited) and re-evaluates the quota state. If the new limit moves the sandbox back under-quota, the matching ingress/egress block is cleared. If it leaves the sandbox over, the matching block is (re-)applied.
Egress clears are conditional: NetworkBlockAll uses the same DOCKER-USER row as the quota egress block, so we must not lift it here when the operator's blanket egress block is still on. See pkg/docker/netrules commentary on the shared rule.
func (*Service) SnapshotPushReconciler ¶ added in v0.2.4
func (s *Service) SnapshotPushReconciler() *SnapshotPushReconciler
SnapshotPushReconciler exposes the reconciler so main.go can drive it from a ticker. Returns nil when snapshot push is disabled — the ticker wrapper should no-op cleanly in that case.
func (*Service) StartBuiltImageGC ¶ added in v0.1.7
StartBuiltImageGC launches the periodic janitor that removes locally-built images (BuiltImageNamespace, i.e. "aerolvm-build/*") that are no longer referenced by any active sandbox AND were created more than the configured TTL ago. Without this, two failure modes leak images forever:
- POST /v1/images/build called standalone (no follow-up CreateSandbox).
- Build succeeded, CreateSandbox failed AND the daytona facade's inline rollback couldn't reach the daemon (e.g. server-side panic, dropped connection between build success and rollback call).
The TTL keeps the janitor from racing the dominant build+create flow: an image built moments ago must clear ImageBuildGCTTL before it's eligible, so a transient network hiccup between build and create can't have the janitor yanking an image the client is about to consume.
No-op if ImageBuildGCEnabled is false or ImageBuildGCInterval <= 0.
func (*Service) StartCaddyCoalescer ¶ added in v0.2.4
StartCaddyCoalescer starts the periodic-drain goroutine. Called once from cmd/sandboxd/main.go after svc.New on worker nodes. Flush-only callers (installHTTPPortRoute today) don't strictly need this — Flush drives its own drain — but the ticker is the safety net for any future Enqueue (fire-and-forget) callsites and for ops that get stranded when a Flush caller's ctx cancels mid-drain.
func (*Service) StartClusterIngressReconcile ¶ added in v0.2.1
func (*Service) StartEventMonitor ¶
StartEventMonitor launches the Docker event consumer goroutine. It is the realtime counterpart to Reconcile() — when a container dies, OOM-kills, or is destroyed out-of-band, this loop updates the DB and tears down routes within ~1s instead of waiting for the next reconcile tick.
func (*Service) StartL4WakeProxy ¶ added in v0.2.4
StartL4WakeProxy starts the loopback TCP wake listener used by raw-TCP serverless routes. TLS-SNI wake uses per-exposure Unix sockets installed by ensureTLSWakeListener, but those listeners are also closed from this context.
func (*Service) StartLifecycleSweep ¶
StartLifecycleSweep launches the per-sandbox lifecycle ticker. Every minute it evaluates each sandbox's Lifecycle timers (StopIfIdleFor / DestroyIfIdleFor / StopAtAge / DestroyAtAge) plus the legacy global SB_IDLE_TIMEOUT_MIN fallback for sandboxes that don't declare any per-sandbox timers. Without either configured, the sweep still runs but is a no-op — kept on so a later UpdateLifecycle call doesn't need to start a goroutine.
func (*Service) StartReconcileLoop ¶
func (*Service) StartSandbox ¶
func (*Service) StopCaddyCoalescer ¶ added in v0.2.4
func (s *Service) StopCaddyCoalescer()
StopCaddyCoalescer drains any pending op and exits the Run goroutine. No-op when StartCaddyCoalescer was never called — caddyCoalescer.Stop blocks on c.done which is only closed by Run, so calling it without a prior Start would hang shutdown.
func (*Service) StopSandbox ¶
StopSandbox is the operator-initiated stop (API surface, manual). It is a thin wrapper over stopSandboxInternal that pins the stop mode to manual: wake_armed is always cleared, so a serverless sandbox stopped via this path stays down until the operator explicitly starts it again. See serverless.go for the full wake-arming policy.
func (*Service) ToolboxTarget ¶
func (*Service) TouchSandbox ¶
TouchSandbox bumps last_active_at for id with per-sandbox debounce. The wake-aware ingress proxy and the toolbox/session/runtime proxies call this on every request, so it must be cheap under high per-sandbox RPS — see touch_coalescer.go for the rationale and the debounce window. Callers that need a guaranteed immediate flush (lifecycle transitions, tests) should drop to s.store.Touch directly.
func (*Service) UnexposePort ¶
func (*Service) UnsealClusterSecrets ¶ added in v0.2.1
func (s *Service) UnsealClusterSecrets(redacted models.CreateSandboxRequest, sealed []byte) (models.CreateSandboxRequest, error)
UnsealClusterSecrets opens a sealed bag and merges the credentials back into the previously-redacted spec. Returns the merged spec; the input is not mutated. An empty sealed payload returns redacted unchanged so callers don't have to short-circuit themselves.
The merge prefers the redacted spec for non-secret fields (Registry Server/Username, mount Source/Target/Options) and overlays only the secret bits — so a future credential rotation that re-seals doesn't stomp on whatever the latest replicated metadata is.
func (*Service) UnsealClusterSecretsForNode ¶ added in v0.2.1
func (s *Service) UnsealClusterSecretsForNode(redacted models.CreateSandboxRequest, sealed []byte, nodeID string) (models.CreateSandboxRequest, error)
func (*Service) UnsealRegistry ¶ added in v0.2.1
func (s *Service) UnsealRegistry(sealed []byte) (*models.RegistryAuth, error)
UnsealRegistry decrypts a previously sealed RegistryAuth. Returns nil/nil when the input is empty (no credentials persisted). Exported for the boot-time backfill in cmd/sandboxd that rebuilds CreateSandboxRequest from the persisted Sandbox row.
func (*Service) UpdateLifecycle ¶
func (s *Service) UpdateLifecycle(ctx context.Context, id string, l models.Lifecycle) (*models.Sandbox, error)
UpdateLifecycle replaces the lifecycle timers on an existing sandbox. Full-replacement semantics: pass zero in any field to clear that timer. The sweep picks up the new values on its next tick (within ~1 minute), so a tightened deadline can fire as soon as the next sweep runs.
func (*Service) UpdateTags ¶ added in v0.1.7
UpdateTags replaces sandboxes.tags_json for the given sandbox. Tags are the native key/value bag — facades use it for label-style metadata (Daytona labels, E2B metadata).
func (*Service) UpsertCompatState ¶ added in v0.1.7
func (*Service) UpsertSnapshotAlias ¶ added in v0.1.7
func (*Service) WakeAwareL4PortTarget ¶ added in v0.2.4
WakeAwareL4PortTarget resolves a raw TCP upstream, ensuring the sandbox is awake first. It is shared by raw TCP and TLS-SNI wake proxying.
func (*Service) WakeAwarePortTarget ¶ added in v0.2.4
func (s *Service) WakeAwarePortTarget(ctx context.Context, id string, port int) (PortEndpoint, error)
WakeAwarePortTarget resolves a sandbox's exposed-port upstream URL, ensuring the sandbox is awake first. Used by the loopback ingress proxy when Caddy forwards a wake-aware HTTP port route. The same sentinels EnsureSandboxAwakeForHTTP returns flow through here unchanged.
func (*Service) WakeAwareToolboxTarget ¶ added in v0.2.4
WakeAwareToolboxTarget is the entry point every control-plane HTTP proxy (v1 toolbox/sessions, daytona toolbox, e2b runtime) calls in place of ToolboxTarget. It funnels every request through the wake helper so a stopped serverless sandbox with wake_armed=true is resurrected before the proxy attempts to dial the toolbox.
Behavior is exactly ToolboxTarget for: running sandboxes, stopped non-serverless sandboxes, and any case where cfg.EnableServerless is off (rollout-gate contract). The two new sentinels (ErrSandboxManuallyStopped, ErrWakeCircuitOpen) are surfaced upstream so apihttp can map them to 409 and 503+Retry-After:60 respectively.
type SnapshotPushConfig ¶ added in v0.2.4
type SnapshotPushConfig struct {
// Enabled gates the whole feature. When false the snapshot-create
// path stays byte-identical to today (no push, no state transitions).
Enabled bool
// Host is the AOCR registry vhost the daemon pushes to
// (e.g. `aocr.aerol.ai`). Falls back to ImageDistributionAOCRHost
// at the wiring layer when MirrorPushHost is unset.
Host string
// ClusterID is the AOCR-scoped tenant the snapshot ends up under.
// Used both as the registry-auth username (by convention; AOCR
// validates the PAT, not the username) and as the path namespace
// `cluster/<id>/snapshots/<name>`.
ClusterID string
// PATPath is the file path to the bearer token presented as the
// Docker registry password. File-sourced so rotation is just a
// file write; never logged.
PATPath string
}
SnapshotPushConfig wires the optional background AOCR push to a specific registry endpoint. A zero value disables push; the snapshot path runs exactly as before. main() validates that PATPath, ClusterID, and Host are non-empty when Enabled=true.
Auth reuses the auto-import cluster PAT (same secret, same rotation surface). The PAT file is re-read on every call so an operator can rotate it without restarting sandboxd — there's no in-memory cache.
func (SnapshotPushConfig) Validate ¶ added in v0.2.4
func (c SnapshotPushConfig) Validate() error
Validate enforces the consistency rules the snapshot pusher relies on. Returns nil for a disabled config so callers can validate unconditionally.
type SnapshotPushDocker ¶ added in v0.2.4
type SnapshotPushDocker interface {
PushImage(ctx context.Context, req docker.PushImageRequest) (string, error)
}
SnapshotPushDocker is the narrow Docker surface the pusher needs. Pulled out so tests can inject a fake without instantiating a real Docker client, and so the dependency surface of the pusher stays readable at a glance.
type SnapshotPushReconcileStats ¶ added in v0.2.4
SnapshotPushReconcileStats is what RunOnce reports back. Used by tests and by the metrics-collection wrapper in main.go.
type SnapshotPushReconciler ¶ added in v0.2.4
type SnapshotPushReconciler struct {
// contains filtered or unexported fields
}
SnapshotPushReconciler walks rows the snapshot-create path flagged `push_state IN ('pending', 'error')` and re-tries them. One ticker per node; fan-out inside a single tick is capped to keep a recovery storm from saturating the local Docker daemon or AOCR.
Like AutoImportReconciler, the type does not own a goroutine — RunOnce is the unit of work and the ticker lives in main.go. That keeps the type fully testable (table-driven, no time.Tick) and lets an operator trigger an immediate retry from a debug endpoint if they want.
func NewSnapshotPushReconciler ¶ added in v0.2.4
func NewSnapshotPushReconciler(pusher *SnapshotPusher, store SnapshotPushStore, logger *slog.Logger, maxInFlight int) *SnapshotPushReconciler
NewSnapshotPushReconciler builds the reconciler. Returns nil if the pusher is nil (snapshot push disabled) so callers can `if r == nil { return }` cheaply. maxInFlight clamps fan-out per tick; pass 0 for the safe default (2 — pushes are I/O-heavy).
func (*SnapshotPushReconciler) RunOnce ¶ added in v0.2.4
func (r *SnapshotPushReconciler) RunOnce(ctx context.Context) (SnapshotPushReconcileStats, error)
RunOnce processes one sweep of the push-pending queue. Blocks until all per-snapshot attempts have returned (success, failure, or context cancel). The bounded semaphore caps in-flight Docker push streams; the AOCR side is usually fine with the load but the local daemon's network/IO is not infinitely parallel.
type SnapshotPushResult ¶ added in v0.2.4
type SnapshotPushResult struct {
RegistryRef string
// Digest is the manifest digest (`sha256:...`) the registry assigned
// to the pushed tag, as reported by Docker's push stream `aux`
// payload. Empty when the daemon did not surface one — older daemons
// or registries that don't return a manifest digest will leave this
// blank and the reconciler stores "" rather than fabricating a value.
Digest string
}
SnapshotPushResult is what PushOnce returns on success. The destination ref is what the reconciler writes onto the snapshot row's ImageRegistryRef field — and what cluster placement on other nodes will see when they decide where to start a sandbox from this snapshot.
type SnapshotPushStore ¶ added in v0.2.4
type SnapshotPushStore interface {
ListSnapshotsPendingPush(ctx context.Context) ([]*models.SandboxSnapshot, error)
SetSnapshotPushState(ctx context.Context, name, state, errMsg string) error
UpdateSnapshotImageDistribution(ctx context.Context, name, mode, registryRef, digest string) error
}
SnapshotPushStore is the narrow store surface the reconciler needs. Pulled out as an interface so the reconciler is testable without a real SQLite store and so future store backends can satisfy it without growing the reconciler's dependency surface.
type SnapshotPusher ¶ added in v0.2.4
type SnapshotPusher struct {
// contains filtered or unexported fields
}
SnapshotPusher pushes a locally-committed snapshot image to AOCR. It holds no per-snapshot state; instantiate once at boot and reuse. Safe for concurrent use.
func NewSnapshotPusher ¶ added in v0.2.4
func NewSnapshotPusher(cfg SnapshotPushConfig, docker SnapshotPushDocker, logger *slog.Logger) (*SnapshotPusher, error)
NewSnapshotPusher builds a pusher. Returns nil if cfg.Enabled is false so callers can do `if pusher == nil { return }` on the snapshot-create path without an extra config check.
func (*SnapshotPusher) DestRefFor ¶ added in v0.2.4
func (p *SnapshotPusher) DestRefFor(snapshotName string) string
DestRefFor returns the AOCR destination ref this pusher would write for the given snapshot name. Used by callers that need the canonical ref without actually pushing — e.g. cross-node snapshot lookup, where a peer node took the snapshot and we want to pull from AOCR without asking the FSM whether the snapshot exists. Returns "" when the pusher is nil so callers can no-op cleanly when push is disabled.
func (*SnapshotPusher) PushOnce ¶ added in v0.2.4
func (p *SnapshotPusher) PushOnce(ctx context.Context, snapshot *models.SandboxSnapshot) (SnapshotPushResult, error)
PushOnce tags the snapshot's local image as the AOCR destination ref and pushes via the Docker daemon. Returns the destination ref on success — used by the reconciler to fill ImageRegistryRef and flip ImageDistributionMode to "aocr".
Idempotent: re-pushing the same `cluster/<id>/snapshots/<name>:latest` tag is safe; the registry treats the layer set as already-present and the daemon's push completes without error.
The PAT is read fresh from disk on every call so an operator can rotate it without restarting sandboxd. The PAT never appears in the returned error message — only the registry's response body does.
type ToolboxEndpoint ¶
Source Files
¶
- auto_import.go
- auto_import_retry.go
- auto_import_writeback.go
- caddy_coalescer.go
- cluster_ownership.go
- cluster_secrets.go
- events.go
- facade_state.go
- image_distribution.go
- ingress_delta.go
- ingress_metrics.go
- l4wake.go
- metrics.go
- netstats.go
- route_shape.go
- serverless.go
- service.go
- snapshot_push.go
- snapshot_push_retry.go
- touch_coalescer.go