service

package
v0.2.4 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: May 24, 2026 License: MIT Imports: 42 Imported by: 0

Documentation

Index

Constants

This section is empty.

Variables

View Source
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.

View Source
var ErrPreferredHostPortUnavailable = errors.New("preferred host port unavailable on this node; exposure parked")

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.

View Source
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.

View Source
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

func GenerateSandboxID() (string, error)

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

type AutoImportReconcileStats struct {
	Scanned   int
	Succeeded int
	Failed    int
	Skipped   int
}

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

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

type AutoImportResult struct {
	RegistryRef    string
	AlreadyPresent bool
}

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

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.

const (
	RouteKindHTTP RouteKind = iota
	RouteKindL4
)

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 New

func New(cfg config.Config, logger *slog.Logger, db *store.Store, runtimeDriver runtime.Runtime, eventsClient *docker.Client, caddyClient *caddy.Client, cipher *secrets.Cipher, mountManager *mounts.Manager, admitter *capacity.Admitter) *Service

func (*Service) AttachCluster added in v0.2.1

func (s *Service) AttachCluster(c cluster.Client)

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

func (s *Service) Capacity() capacity.Snapshot

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 (s *Service) ClaimIdempotentRequest(ctx context.Context, scope, fingerprint string, now time.Time, pendingTTL time.Duration) (*models.IdempotentRequestRecord, bool, error)

func (*Service) Cluster added in v0.2.1

func (s *Service) Cluster() cluster.Client

Cluster returns the attached cluster.Client. Always non-nil.

func (*Service) ClusterTopologyError added in v0.2.2

func (s *Service) ClusterTopologyError() error

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 (s *Service) CompleteIdempotentRequest(ctx context.Context, scope, fingerprint, targetID string, now time.Time, replayTTL time.Duration) error

func (*Service) CreateSandbox

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 (s *Service) DeleteClusterSecrets(ctx context.Context, sandboxID string) error

func (*Service) DeleteIdempotentRequest added in v0.1.7

func (s *Service) DeleteIdempotentRequest(ctx context.Context, scope, fingerprint string) error

func (*Service) DeleteSnapshot added in v0.1.7

func (s *Service) DeleteSnapshot(ctx context.Context, idOrName string) error

func (*Service) DeleteSnapshotAlias added in v0.1.7

func (s *Service) DeleteSnapshotAlias(ctx context.Context, alias string) error

func (*Service) DestroySandbox

func (s *Service) DestroySandbox(ctx context.Context, id string) error

func (*Service) EnsureClusterReady added in v0.2.1

func (s *Service) EnsureClusterReady(ctx context.Context) error

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

func (s *Service) EnsureLayer4Ready(ctx context.Context) error

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

func (s *Service) EnsureNetstatsReady(ctx context.Context) error

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

func (s *Service) ForceReconcileHTTPWakeShape(ctx context.Context) error

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 (s *Service) GetCompatState(ctx context.Context, sandboxID, facade string) (*models.SandboxCompatState, error)

func (*Service) GetIdempotentRequest added in v0.1.7

func (s *Service) GetIdempotentRequest(ctx context.Context, scope, fingerprint string) (*models.IdempotentRequestRecord, error)

func (*Service) GetNetworkUsage added in v0.1.7

func (s *Service) GetNetworkUsage(ctx context.Context, id string) (*models.NetworkUsage, error)

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 (s *Service) GetSandbox(ctx context.Context, id string) (*models.Sandbox, error)

func (*Service) GetSnapshot added in v0.1.7

func (s *Service) GetSnapshot(ctx context.Context, idOrName string) (*models.SandboxSnapshot, error)

func (*Service) GetSnapshotAlias added in v0.1.7

func (s *Service) GetSnapshotAlias(ctx context.Context, alias string) (*models.SnapshotAlias, error)

func (*Service) Health

func (s *Service) Health(ctx context.Context) (models.HealthStatus, error)

func (*Service) IsSandboxStarted added in v0.2.4

func (s *Service) IsSandboxStarted(ctx context.Context, id string) (bool, error)

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 (s *Service) ListCompatState(ctx context.Context, facade string) (map[string]models.SandboxCompatState, error)

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 (s *Service) ListSnapshotAliases(ctx context.Context, facade string) (map[string]models.SnapshotAlias, error)

func (*Service) ListSnapshots added in v0.1.7

func (s *Service) ListSnapshots(ctx context.Context) ([]*models.SandboxSnapshot, error)

func (*Service) MarkImported added in v0.2.3

func (s *Service) MarkImported(ctx context.Context, sandboxID string, registryRef string)

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 (*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) Reconcile

func (s *Service) Reconcile(ctx context.Context) error

func (*Service) ReconcileClusterIngress added in v0.2.1

func (s *Service) ReconcileClusterIngress(ctx context.Context) error

func (*Service) ReconstructWakeArmedIfNeeded added in v0.2.4

func (s *Service) ReconstructWakeArmedIfNeeded(ctx context.Context, sandbox *models.Sandbox)

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

func (s *Service) ReplayClusterOwnership(ctx context.Context) (int, error)

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

func (s *Service) ReplayReservations(ctx context.Context)

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 (s *Service) ResizeSandbox(ctx context.Context, id string, req models.ResizeSandboxRequest) (*models.Sandbox, error)

func (*Service) ResolveSandboxIDByName added in v0.1.7

func (s *Service) ResolveSandboxIDByName(ctx context.Context, name string) (string, error)

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 (s *Service) SealClusterSecretsForRecipient(req models.CreateSandboxRequest, recipient string) ([]byte, error)

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

func (s *Service) StartBuiltImageGC(ctx context.Context)

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

func (s *Service) StartCaddyCoalescer(ctx context.Context)

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 (s *Service) StartClusterIngressReconcile(ctx context.Context)

func (*Service) StartEventMonitor

func (s *Service) StartEventMonitor(ctx context.Context)

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

func (s *Service) StartL4WakeProxy(ctx context.Context) error

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

func (s *Service) StartLifecycleSweep(ctx context.Context)

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 (s *Service) StartReconcileLoop(ctx context.Context)

func (*Service) StartSandbox

func (s *Service) StartSandbox(ctx context.Context, id string) (*models.Sandbox, error)

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

func (s *Service) StopSandbox(ctx context.Context, id string) (*models.Sandbox, error)

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 (s *Service) ToolboxTarget(ctx context.Context, id string) (ToolboxEndpoint, error)

func (*Service) TouchSandbox

func (s *Service) TouchSandbox(ctx context.Context, id string) error

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 (s *Service) UnexposePort(ctx context.Context, id string, port int) error

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

func (s *Service) UpdateTags(ctx context.Context, sandboxID string, tags map[string]string) error

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 (s *Service) UpsertCompatState(ctx context.Context, sandboxID, facade, stateJSON string) error

func (*Service) UpsertSnapshotAlias added in v0.1.7

func (s *Service) UpsertSnapshotAlias(ctx context.Context, alias models.SnapshotAlias) error

func (*Service) WakeAwareL4PortTarget added in v0.2.4

func (s *Service) WakeAwareL4PortTarget(ctx context.Context, id string, port int) (string, error)

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

func (s *Service) WakeAwareToolboxTarget(ctx context.Context, id string) (ToolboxEndpoint, error)

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

type SnapshotPushReconcileStats struct {
	Scanned   int
	Succeeded int
	Failed    int
	Skipped   int
}

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

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

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

type ToolboxEndpoint struct {
	URL   string
	Token string
}

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL