alien_bindings/traits.rs
1use crate::error::Result;
2use crate::presigned::PresignedRequest;
3use alien_core::{BuildConfig, BuildExecution};
4use async_trait::async_trait;
5use object_store::path::Path;
6use object_store::ObjectStore;
7use serde::{Deserialize, Serialize};
8use std::collections::BTreeMap;
9use std::sync::Arc;
10use std::time::Duration;
11use url::Url;
12
13#[cfg(feature = "openapi")]
14use utoipa::ToSchema;
15
16/// Marker trait for all binding types.
17pub trait Binding: Send + Sync + std::fmt::Debug {}
18
19/// A storage binding that provides object store capabilities.
20#[async_trait]
21pub trait Storage: Binding + ObjectStore {
22 /// Gets the base directory path configured for this storage binding.
23 fn get_base_dir(&self) -> Path;
24 /// Gets the underlying URL configured for this storage binding.
25 fn get_url(&self) -> Url;
26
27 /// Creates a presigned request for uploading data to the specified path.
28 /// The request can be serialized, stored, and executed later.
29 async fn presigned_put(&self, path: &Path, expires_in: Duration) -> Result<PresignedRequest>;
30
31 /// Creates a presigned request for downloading data from the specified path.
32 /// The request can be serialized, stored, and executed later.
33 async fn presigned_get(&self, path: &Path, expires_in: Duration) -> Result<PresignedRequest>;
34
35 /// Creates a presigned request for deleting the object at the specified path.
36 /// The request can be serialized, stored, and executed later.
37 async fn presigned_delete(&self, path: &Path, expires_in: Duration)
38 -> Result<PresignedRequest>;
39}
40
41/// A provider-backed key for small wrap and unwrap operations.
42#[async_trait]
43pub trait Key: Binding {
44 async fn encrypt(
45 &self,
46 plaintext: &[u8],
47 context: Option<&BTreeMap<String, String>>,
48 ) -> Result<Vec<u8>>;
49
50 async fn decrypt(
51 &self,
52 ciphertext: &[u8],
53 context: Option<&BTreeMap<String, String>>,
54 ) -> Result<Vec<u8>>;
55}
56
57/// A build binding that provides build execution capabilities.
58#[async_trait]
59pub trait Build: Binding {
60 /// Starts a new build with the given configuration.
61 /// Returns the build execution information.
62 async fn start_build(&self, config: BuildConfig) -> Result<BuildExecution>;
63
64 /// Gets the status of a specific build execution.
65 async fn get_build_status(&self, build_id: &str) -> Result<BuildExecution>;
66
67 /// Stops or cancels a running build.
68 async fn stop_build(&self, build_id: &str) -> Result<()>;
69}
70
71/// AWS IAM Role service account information
72#[derive(Debug, Clone, Serialize, Deserialize)]
73#[serde(rename_all = "camelCase")]
74#[cfg_attr(feature = "openapi", derive(ToSchema))]
75pub struct AwsServiceAccountInfo {
76 /// The IAM role name
77 pub role_name: String,
78 /// The IAM role ARN (for AssumeRole)
79 pub role_arn: String,
80}
81
82/// GCP Service Account information
83#[derive(Debug, Clone, Serialize, Deserialize)]
84#[serde(rename_all = "camelCase")]
85#[cfg_attr(feature = "openapi", derive(ToSchema))]
86pub struct GcpServiceAccountInfo {
87 /// The service account email (for impersonation)
88 pub email: String,
89 /// The service account unique ID
90 pub unique_id: String,
91}
92
93/// Azure User-Assigned Managed Identity information
94#[derive(Debug, Clone, Serialize, Deserialize)]
95#[serde(rename_all = "camelCase")]
96#[cfg_attr(feature = "openapi", derive(ToSchema))]
97pub struct AzureServiceAccountInfo {
98 /// The managed identity client ID (for authentication)
99 pub client_id: String,
100 /// The managed identity resource ID (ARM ID)
101 pub resource_id: String,
102 /// The managed identity principal ID
103 pub principal_id: String,
104}
105
106/// Platform-specific service account information
107#[derive(Debug, Clone, Serialize, Deserialize)]
108#[serde(tag = "platform", rename_all = "camelCase")]
109#[cfg_attr(feature = "openapi", derive(ToSchema))]
110pub enum ServiceAccountInfo {
111 /// AWS IAM Role
112 Aws(AwsServiceAccountInfo),
113 /// GCP Service Account
114 Gcp(GcpServiceAccountInfo),
115 /// Azure User-Assigned Managed Identity
116 Azure(AzureServiceAccountInfo),
117}
118
119/// Configuration for impersonation
120#[derive(Debug, Clone)]
121pub struct ImpersonationRequest {
122 /// Optional session name (AWS only)
123 pub session_name: Option<String>,
124 /// Optional session duration in seconds
125 pub duration_seconds: Option<i32>,
126 /// Optional scopes (GCP only)
127 pub scopes: Option<Vec<String>>,
128}
129
130impl Default for ImpersonationRequest {
131 fn default() -> Self {
132 Self {
133 session_name: None,
134 duration_seconds: Some(3600), // 1 hour default
135 scopes: None,
136 }
137 }
138}
139
140/// A service account binding that provides identity and impersonation capabilities.
141#[async_trait]
142pub trait ServiceAccount: Binding {
143 /// Gets information about the service account
144 async fn get_info(&self) -> Result<ServiceAccountInfo>;
145
146 /// Impersonates the service account and returns credentials as a ClientConfig.
147 ///
148 /// This performs the cloud-specific impersonation:
149 /// - AWS: STS AssumeRole to get temporary credentials
150 /// - GCP: IAM Credentials API generateAccessToken
151 /// - Azure: Uses the attached managed identity (no API call needed)
152 async fn impersonate(&self, request: ImpersonationRequest) -> Result<alien_core::ClientConfig>;
153
154 /// Helper for downcasting trait object
155 fn as_any(&self) -> &dyn std::any::Any;
156}
157
158/// Response from repository operations.
159#[derive(Debug, Clone, Serialize, Deserialize)]
160#[serde(rename_all = "camelCase")]
161#[cfg_attr(feature = "openapi", derive(ToSchema))]
162pub struct RepositoryResponse {
163 /// The **routable name** of the repository — the full, platform-specific
164 /// path used for subsequent calls (`get_repository`, `delete_repository`,
165 /// `generate_credentials`, `*_cross_account_access`).
166 ///
167 /// Per-platform format (matches
168 /// `alien.dev/content/docs/infrastructure/artifact-registry/behavior.mdx`):
169 ///
170 /// | Platform | Format |
171 /// |---|---|
172 /// | AWS (ECR) | `{registry_prefix}-{logical}` (e.g. `alien-artifacts-my-app`) |
173 /// | GCP (GAR) | `{project_id}/{gar_repo}/{logical}` |
174 /// | Azure (ACR) | `{logical}` (used directly) |
175 /// | Local | `{binding_name}/{logical}` |
176 ///
177 /// **Round-trip invariant:** callers MUST be able to pass this value back
178 /// to any other method on the trait without further transformation.
179 /// Implementations MUST NOT re-apply prefixing in receivers — assume
180 /// `repo_id` arguments are already routable.
181 pub name: String,
182 /// Repository URI for pushing/pulling images. None if repository is not ready yet.
183 pub uri: Option<String>,
184 /// Optional creation timestamp in ISO8601 format.
185 pub created_at: Option<String>,
186}
187
188/// Permissions level for artifact registry access.
189#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
190#[serde(rename_all = "kebab-case")]
191#[cfg_attr(feature = "openapi", derive(ToSchema))]
192pub enum ArtifactRegistryPermissions {
193 /// Pull-only access (download artifacts).
194 Pull,
195 /// Push and pull access (upload and download artifacts).
196 PushPull,
197}
198
199/// How the registry expects credentials to be presented.
200#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
201#[serde(rename_all = "lowercase")]
202#[cfg_attr(feature = "openapi", derive(ToSchema))]
203pub enum RegistryAuthMethod {
204 /// HTTP Basic auth (username:password). Used by ECR, GAR, Local.
205 Basic,
206 /// HTTP Bearer token. Used by ACR (Azure).
207 Bearer,
208}
209
210/// Credentials for accessing a repository.
211#[derive(Debug, Clone, Serialize, Deserialize)]
212#[serde(rename_all = "camelCase")]
213#[cfg_attr(feature = "openapi", derive(ToSchema))]
214pub struct ArtifactRegistryCredentials {
215 /// How to present these credentials to the registry.
216 pub auth_method: RegistryAuthMethod,
217 /// Username for authentication (empty for Bearer auth).
218 pub username: String,
219 /// Password or token for authentication.
220 pub password: String,
221 /// Optional expiration time in ISO8601 format.
222 pub expires_at: Option<String>,
223}
224
225/// Types of compute services that can access artifact registries.
226#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
227#[serde(rename_all = "kebab-case")]
228#[cfg_attr(feature = "openapi", derive(ToSchema))]
229pub enum ComputeServiceType {
230 /// Serverless functions
231 Worker,
232 /// Sandbox sessions.
233 ///
234 /// Separate from `Worker` because a cloud may pull each as a different principal (GCP's
235 /// per-service agents); a provider is free to grant nothing further if one principal covers both.
236 Sandbox,
237 // In the future, we could add Container, VirtualMachine, Kubernetes, etc.
238}
239
240impl ComputeServiceType {
241 /// Every variant, for the callers that have to name one principal per compute service with no
242 /// value to match on: a revoke with no record to go on, and the table reading a member back.
243 pub const ALL: &'static [Self] = &[Self::Worker, Self::Sandbox];
244
245 /// Adding a variant fails to compile here, beside the list it also has to join — a compute
246 /// service missing from `ALL` is granted a pull that nothing revokes.
247 #[allow(dead_code)]
248 fn all_is_exhaustive(self) {
249 match self {
250 Self::Worker | Self::Sandbox => {}
251 }
252 }
253}
254
255/// Cross-account access configuration for AWS artifact registries.
256#[derive(Debug, Clone, Serialize, Deserialize)]
257#[serde(rename_all = "camelCase")]
258#[cfg_attr(feature = "openapi", derive(ToSchema))]
259pub struct AwsCrossAccountAccess {
260 /// AWS account IDs that should have cross-account access.
261 pub account_ids: Vec<String>,
262 /// AWS regions where the target Lambda functions run.
263 /// Used to construct `aws:sourceArn` patterns for the Lambda service principal condition.
264 pub regions: Vec<String>,
265 /// Types of compute services that should have access.
266 pub allowed_service_types: Vec<ComputeServiceType>,
267 /// Specific IAM role ARNs to grant access to.
268 /// These are typically deployment/management roles or service-specific roles.
269 pub role_arns: Vec<String>,
270}
271
272/// Cross-account access configuration for GCP artifact registries.
273#[derive(Debug, Clone, Serialize, Deserialize)]
274#[serde(rename_all = "camelCase")]
275#[cfg_attr(feature = "openapi", derive(ToSchema))]
276pub struct GcpCrossAccountAccess {
277 /// GCP project numbers that should have access.
278 pub project_numbers: Vec<String>,
279 /// Types of compute services that should have access.
280 pub allowed_service_types: Vec<ComputeServiceType>,
281 /// Additional service account emails to grant access to.
282 /// These are typically deployment/management service accounts.
283 pub service_account_emails: Vec<String>,
284}
285
286/// Platform-specific cross-account access configuration.
287#[derive(Debug, Clone, Serialize, Deserialize)]
288#[serde(tag = "platform", rename_all = "lowercase")]
289#[cfg_attr(feature = "openapi", derive(ToSchema))]
290pub enum CrossAccountAccess {
291 /// AWS-specific cross-account access configuration.
292 Aws(AwsCrossAccountAccess),
293 /// GCP-specific cross-account access configuration.
294 Gcp(GcpCrossAccountAccess),
295}
296
297/// Current cross-account access permissions for a repository.
298#[derive(Debug, Clone, Serialize, Deserialize)]
299#[serde(rename_all = "camelCase")]
300#[cfg_attr(feature = "openapi", derive(ToSchema))]
301pub struct CrossAccountPermissions {
302 /// Platform-specific access configuration currently applied.
303 pub access: CrossAccountAccess,
304 /// Timestamp when permissions were last updated.
305 pub last_updated: Option<String>,
306}
307
308/// A trait for artifact registry bindings that provide container image repository management.
309#[async_trait]
310pub trait ArtifactRegistry: Binding {
311 /// Returns the raw registry endpoint URL (e.g., "https://123456.dkr.ecr.us-east-1.amazonaws.com"
312 /// or "http://localhost:5000"). Used by the push proxy to forward requests transparently.
313 ///
314 /// Default returns empty string — cloud provider implementations should override.
315 fn registry_endpoint(&self) -> String {
316 String::new()
317 }
318
319 /// Returns the OCI repository path prefix used for upstream operations.
320 ///
321 /// This identifier serves two related roles, both pointing at the same
322 /// upstream location:
323 ///
324 /// 1. **Proxy routing.** When the push proxy forwards push/pull requests
325 /// to the upstream registry, this prefix is prepended to the image
326 /// name portion of the OCI path.
327 /// 2. **Shared deployment-image repository name.** `alien release` pushes
328 /// every function image as `{prefix}:{logical}-{hash}` into one shared
329 /// repository whose routable name is exactly this prefix. Pass it as
330 /// `repo_id` when calling `add_cross_account_access` /
331 /// `remove_cross_account_access` for the deployment cross-account
332 /// flow.
333 ///
334 /// Examples:
335 /// - ECR: `"alien-e2e"` — flat repo prefix; also the routable repo name
336 /// for the shared deployment-image repository
337 /// - GAR: `"my-project/alien-e2e"` — project/repo structure
338 /// - ACR: `"alien-e2e"` when configured, or `""` for registry-root
339 /// repositories; principal pull access is granted on the parent registry
340 /// - Local: `"artifacts"` or similar — cross-account not supported
341 ///
342 /// An empty return value indicates the platform has no shared
343 /// deployment-image repo at the binding level.
344 ///
345 /// Default returns empty string.
346 fn upstream_repository_prefix(&self) -> String {
347 String::new()
348 }
349
350 /// Creates a repository within the artifact registry.
351 ///
352 /// `repo_name` is the **logical** identifier the caller chose (e.g.
353 /// `"my-app"`). The implementation transforms it to the routable
354 /// platform-specific form before calling any backend API; what's
355 /// returned in [`RepositoryResponse::name`] is the routable form.
356 ///
357 /// On platforms where image paths are implicit (GAR, ACR, Local),
358 /// this may not call any backend API — but it still returns a valid
359 /// routable name.
360 async fn create_repository(&self, repo_name: &str) -> Result<RepositoryResponse>;
361
362 /// Gets repository details. `repo_id` is the routable name returned by
363 /// [`Self::create_repository`]; implementations MUST NOT re-apply
364 /// prefixing.
365 async fn get_repository(&self, repo_id: &str) -> Result<RepositoryResponse>;
366
367 /// Adds cross-account access permissions for a repository.
368 /// This adds the specified permissions to any existing cross-account permissions.
369 ///
370 /// `repo_id` is the routable name from [`Self::create_repository`].
371 ///
372 /// For AWS: grants access to specified account IDs with configurable principals and compute service types (ECR repository policy).
373 /// For GCP: grants access to serverless robots and service accounts on the parent GAR registry (image-path-level IAM is not supported).
374 /// For Azure: not supported — returns `OperationNotSupported`.
375 async fn add_cross_account_access(
376 &self,
377 repo_id: &str,
378 access: CrossAccountAccess,
379 ) -> Result<()>;
380
381 /// Removes cross-account access permissions for a repository.
382 ///
383 /// `repo_id` is the routable name from [`Self::create_repository`].
384 ///
385 /// For AWS: removes access from the ECR repository policy.
386 /// For GCP: removes IAM bindings on the parent GAR registry.
387 /// For Azure: not supported — returns `OperationNotSupported`.
388 async fn remove_cross_account_access(
389 &self,
390 repo_id: &str,
391 access: CrossAccountAccess,
392 ) -> Result<()>;
393
394 /// Gets the current cross-account access permissions for a repository.
395 ///
396 /// `repo_id` is the routable name from [`Self::create_repository`].
397 /// For Azure: not supported — returns `OperationNotSupported`.
398 async fn get_cross_account_access(&self, repo_id: &str) -> Result<CrossAccountPermissions>;
399
400 /// Generates credentials for accessing a repository with the specified
401 /// permissions.
402 ///
403 /// `repo_id` is the routable name from [`Self::create_repository`].
404 ///
405 /// Most platforms produce registry-scoped (not repo-scoped) credentials,
406 /// so `repo_id` typically only affects logging — not the credentials
407 /// themselves.
408 async fn generate_credentials(
409 &self,
410 repo_id: &str,
411 permissions: ArtifactRegistryPermissions,
412 ttl_seconds: Option<u32>,
413 ) -> Result<ArtifactRegistryCredentials>;
414
415 /// Deletes a repository and all contained images.
416 ///
417 /// `repo_id` is the routable name from [`Self::create_repository`].
418 /// Implementations MUST NOT delete the *parent* registry (which is owned
419 /// by `alien-infra`); on platforms with implicit image paths (GAR, ACR,
420 /// Local) this is a no-op.
421 async fn delete_repository(&self, repo_id: &str) -> Result<()>;
422}
423
424/// A trait for vault bindings that provide secure secret management.
425#[async_trait]
426pub trait Vault: Binding {
427 /// Gets a secret value by name.
428 async fn get_secret(&self, secret_name: &str) -> Result<String>;
429
430 /// Sets a secret value, creating it if it doesn't exist or updating it if it does.
431 async fn set_secret(&self, secret_name: &str, value: &str) -> Result<()>;
432
433 /// Deletes a secret by name.
434 async fn delete_secret(&self, secret_name: &str) -> Result<()>;
435
436 /// Lists the names of all secrets stored in this vault.
437 ///
438 /// Returned names are in the vault's own namespace (any provider-specific
439 /// prefix is stripped), so each name can be passed straight back to
440 /// [`Vault::get_secret`].
441 ///
442 /// Implementations backed by a flat namespace (a single string prefix
443 /// with no reserved separator, e.g. `"{vault_prefix}-{secret_name}"`)
444 /// cannot implement this safely with a prefix/`BeginsWith` scan: a vault
445 /// named `"app"` would also match a sibling vault named `"app-prod"`,
446 /// aliasing across vaults. Such providers should return
447 /// `OperationNotSupported` rather than list under this hazard.
448 async fn list_secrets(&self) -> Result<Vec<String>>;
449}
450
451/// TLS policy used when building a Postgres connection string.
452#[derive(Debug, Clone, Copy, PartialEq, Eq)]
453pub enum SslMode {
454 /// Plain TCP, no TLS (Local or an explicit BYO opt-out).
455 Disable,
456 /// Require TLS and verify the server certificate against a trusted CA.
457 VerifyCa,
458 /// Require TLS and verify both the trusted CA chain and the dialed hostname
459 /// (BYO / External default, Aurora, and Flexible Server).
460 VerifyFull,
461}
462
463impl SslMode {
464 /// The `sslmode` query-parameter value, and the wire form embedders (the napi addon)
465 /// hand to their own callers.
466 pub fn as_str(self) -> &'static str {
467 match self {
468 SslMode::Disable => "disable",
469 SslMode::VerifyCa => "verify-ca",
470 SslMode::VerifyFull => "verify-full",
471 }
472 }
473}
474
475/// Invalid PEM roots supplied to a verified Postgres TLS policy.
476#[derive(Debug, Clone, Copy, PartialEq, Eq)]
477pub struct InvalidPostgresCaCertificates;
478
479impl std::fmt::Display for InvalidPostgresCaCertificates {
480 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
481 f.write_str("expected one or more non-empty PEM certificates")
482 }
483}
484
485impl std::error::Error for InvalidPostgresCaCertificates {}
486
487/// A complete Postgres TLS policy.
488///
489/// The fields are private so callers cannot pair plaintext with CA roots or construct
490/// `verify-ca` without a root. Cloning the policy is cheap: certificate bundles are
491/// reference-counted and copied only when crossing an FFI boundary.
492#[derive(Clone, PartialEq, Eq)]
493pub struct PostgresTlsPolicy {
494 sslmode: SslMode,
495 ca_certificates: Arc<[String]>,
496}
497
498impl std::fmt::Debug for PostgresTlsPolicy {
499 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
500 f.debug_struct("PostgresTlsPolicy")
501 .field("sslmode", &self.sslmode)
502 .field("ca_certificate_count", &self.ca_certificates.len())
503 .finish()
504 }
505}
506
507impl PostgresTlsPolicy {
508 /// Plain TCP with no certificate roots.
509 pub fn disabled() -> Self {
510 Self {
511 sslmode: SslMode::Disable,
512 ca_certificates: Arc::default(),
513 }
514 }
515
516 /// TLS with CA verification but no hostname verification.
517 ///
518 /// This mode requires at least one valid PEM root.
519 pub fn verify_ca(
520 ca_certificates: Vec<String>,
521 ) -> std::result::Result<Self, InvalidPostgresCaCertificates> {
522 Self::verified(SslMode::VerifyCa, ca_certificates, true)
523 }
524
525 /// TLS with CA and hostname verification.
526 ///
527 /// An empty root set intentionally selects the runtime's system trust store.
528 pub fn verify_full(
529 ca_certificates: Vec<String>,
530 ) -> std::result::Result<Self, InvalidPostgresCaCertificates> {
531 Self::verified(SslMode::VerifyFull, ca_certificates, false)
532 }
533
534 /// TLS with CA and hostname verification using the runtime's system trust store.
535 pub fn verify_full_with_system_roots() -> Self {
536 Self {
537 sslmode: SslMode::VerifyFull,
538 ca_certificates: Arc::default(),
539 }
540 }
541
542 fn verified(
543 sslmode: SslMode,
544 ca_certificates: Vec<String>,
545 require_roots: bool,
546 ) -> std::result::Result<Self, InvalidPostgresCaCertificates> {
547 if (require_roots && ca_certificates.is_empty())
548 || ca_certificates.iter().any(|certificate| {
549 let certificate = certificate.trim();
550 !certificate.starts_with("-----BEGIN CERTIFICATE-----")
551 || !certificate.ends_with("-----END CERTIFICATE-----")
552 })
553 {
554 return Err(InvalidPostgresCaCertificates);
555 }
556
557 Ok(Self {
558 sslmode,
559 ca_certificates: ca_certificates.into(),
560 })
561 }
562
563 /// The libpq-compatible `sslmode` represented by this complete policy.
564 pub fn sslmode(&self) -> SslMode {
565 self.sslmode
566 }
567
568 /// PEM-encoded roots, or an empty slice when this policy uses no roots or the
569 /// runtime's system trust store.
570 pub fn ca_certificates(&self) -> &[String] {
571 &self.ca_certificates
572 }
573}
574
575/// Resolved connection details for a Postgres database.
576#[derive(Clone)]
577pub struct PostgresConnectionParams {
578 pub host: String,
579 pub port: u16,
580 pub database: String,
581 pub username: String,
582 pub password: String,
583 pub tls: PostgresTlsPolicy,
584}
585
586// Hand-written Debug so the resolved password never reaches logs, error chains, or panic output
587// (`Binding`/`BindingsProviderApi` require Debug). Mirrors the KV providers' redacting Debug.
588impl std::fmt::Debug for PostgresConnectionParams {
589 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
590 f.debug_struct("PostgresConnectionParams")
591 .field("host", &self.host)
592 .field("port", &self.port)
593 .field("database", &self.database)
594 .field("username", &self.username)
595 .field("password", &"<redacted>")
596 .field("tls", &self.tls)
597 .finish()
598 }
599}
600
601impl PostgresConnectionParams {
602 /// Creates resolved connection details from a complete, internally consistent TLS
603 /// policy.
604 pub fn new(
605 host: String,
606 port: u16,
607 database: String,
608 username: String,
609 password: String,
610 tls: PostgresTlsPolicy,
611 ) -> Self {
612 Self {
613 host,
614 port,
615 database,
616 username,
617 password,
618 tls,
619 }
620 }
621
622 /// The libpq-compatible TLS mode used by this connection.
623 pub fn sslmode(&self) -> SslMode {
624 self.tls.sslmode()
625 }
626
627 /// PEM-encoded root CA certificates.
628 pub fn ca_certificates(&self) -> &[String] {
629 self.tls.ca_certificates()
630 }
631
632 /// Builds a `postgres://` URL. Username and password are percent-encoded so a
633 /// generated password containing URL-special characters can never corrupt it.
634 pub fn connection_string(&self) -> String {
635 format!(
636 "postgres://{}:{}@{}:{}/{}?sslmode={}",
637 encode_userinfo(&self.username),
638 encode_userinfo(&self.password),
639 self.host,
640 self.port,
641 // Encode the database path segment so this URL stays byte-identical to the TS
642 // resolver's `encodeUserinfo` (the same RFC 3986 unreserved-set encoding).
643 encode_userinfo(&self.database),
644 self.sslmode().as_str(),
645 )
646 }
647}
648
649/// Percent-encodes a URL userinfo component; the RFC 3986 unreserved set passes through.
650fn encode_userinfo(value: &str) -> String {
651 let mut out = String::with_capacity(value.len());
652 for byte in value.bytes() {
653 match byte {
654 b'A'..=b'Z' | b'a'..=b'z' | b'0'..=b'9' | b'-' | b'.' | b'_' | b'~' => {
655 out.push(byte as char)
656 }
657 _ => out.push_str(&format!("%{:02X}", byte)),
658 }
659 }
660 out
661}
662
663/// Connection-only Postgres binding. Unlike every other resource, Postgres ships no
664/// gRPC service and wraps no operations (by design): every backend
665/// speaks the same wire protocol, so the binding returns connection details and the
666/// application uses its own driver.
667pub trait Postgres: Binding {
668 fn connection_params(&self) -> &PostgresConnectionParams;
669
670 /// `postgres://` connection string, derived from `connection_params` and never stored.
671 fn connection_string(&self) -> String {
672 self.connection_params().connection_string()
673 }
674}
675
676/// A precondition for a KV put operation.
677#[derive(Debug, Clone, Default, PartialEq, Eq)]
678pub enum PutCondition {
679 /// Replace any current value, or create the key when it is absent.
680 #[default]
681 None,
682 /// Create the key only when it is absent or logically expired.
683 Absent,
684 /// Replace the value only when the key still has this opaque version.
685 Version(String),
686}
687
688/// Options for KV put operations.
689#[derive(Debug, Clone, Default)]
690pub struct PutOptions {
691 /// Optional TTL for automatic expiration (soft hint - items MAY be deleted after expiry)
692 pub ttl: Option<Duration>,
693 /// Optional atomic write precondition.
694 pub condition: PutCondition,
695}
696
697/// A KV entry together with its opaque version.
698#[derive(Debug, Clone, PartialEq, Eq)]
699pub struct KvEntry {
700 /// Stored key.
701 pub key: String,
702 /// Stored value bytes.
703 pub value: Vec<u8>,
704 /// Provider-neutral version for conditional writes. Callers must treat it as opaque.
705 pub version: String,
706}
707
708/// Represents the result of a scan operation.
709#[derive(Debug)]
710pub struct ScanResult {
711 /// Entries found (may be ≤ limit, no guarantee to fill).
712 pub items: Vec<KvEntry>,
713 /// Opaque, prefix-bound cursor for pagination. None if the traversal is complete.
714 /// A non-empty cursor may lead to an empty page when records expire between pages.
715 pub next_cursor: Option<String>,
716}
717
718/// A trait for key-value store bindings that provide minimal, platform-agnostic KV operations.
719/// This API is designed to work consistently across DynamoDB, Firestore, Azure Table Storage,
720/// and the local provider.
721#[async_trait]
722pub trait Kv: Binding {
723 /// Get an entry by key. Returns `None` if the key doesn't exist or has expired.
724 ///
725 /// **TTL Behavior**: TTL is a soft hint for automatic cleanup. If `now >= expires_at`,
726 /// implementations SHOULD behave as if the key is absent, even if the item still exists
727 /// physically in the backend. Physical deletion is eventual and not guaranteed.
728 ///
729 /// **Validation**: Keys are validated against MAX_KEY_BYTES and portable charset.
730 /// Invalid keys return `KvError::InvalidKey` immediately.
731 async fn get(&self, key: &str) -> Result<Option<KvEntry>>;
732
733 /// Put a value with optional options. Conditional writes return `false` when their precondition
734 /// does not match. Unconditional writes return `true`.
735 ///
736 /// **Size Limits**:
737 /// - Keys: ≤ MAX_KEY_BYTES (512 bytes) with portable ASCII charset
738 /// - Values: ≤ MAX_VALUE_BYTES (24,576 bytes = 24 KiB)
739 ///
740 /// **Validation**: Size and charset constraints are enforced before backend calls.
741 /// Invalid inputs return `KvError::InvalidKey` or `KvError::InvalidValue` immediately.
742 ///
743 /// **TTL Behavior**: TTL is a soft hint for automatic cleanup. If TTL is specified,
744 /// item expires at `put_time + ttl`. Expired items SHOULD appear absent on subsequent
745 /// reads, but physical deletion is eventual and not guaranteed.
746 ///
747 /// **Conditional Logic**: [`PutCondition::Absent`] uses each backend's atomic create
748 /// primitive, with a version-guarded takeover when an expired row still exists physically.
749 /// [`PutCondition::Version`] uses DynamoDB conditions, Firestore update-time preconditions,
750 /// Azure entity tags, or the local database's conditional update.
751 async fn put(&self, key: &str, value: Vec<u8>, options: Option<PutOptions>) -> Result<bool>;
752
753 /// Delete a key. No error if key doesn't exist.
754 ///
755 /// **Validation**: Keys are validated against MAX_KEY_BYTES and portable charset.
756 /// Invalid keys return `KvError::InvalidKey` immediately.
757 /// When `if_version` is supplied, delete only if the key still has that opaque version.
758 /// Returns `false` when a conditional delete finds the key absent, expired, or changed.
759 /// Unconditional deletes return `true`, including when the key is already absent.
760 async fn delete(&self, key: &str, if_version: Option<&str>) -> Result<bool>;
761
762 /// Check if a key exists without retrieving the value.
763 ///
764 /// **TTL Behavior**: TTL is a soft hint for automatic cleanup. If `now >= expires_at`,
765 /// SHOULD return false even if physically present. Physical deletion is eventual and not guaranteed.
766 ///
767 /// **Validation**: Keys are validated against MAX_KEY_BYTES and portable charset.
768 /// Invalid keys return `KvError::InvalidKey` immediately.
769 async fn exists(&self, key: &str) -> Result<bool>;
770
771 /// Scan keys with a prefix, with pagination support.
772 ///
773 /// **Scan Contract**:
774 /// - Returns an **arbitrary, unordered subset** in backend-natural order
775 /// - **No ordering guarantees** across backends (for example, Azure partition fan-out)
776 /// - **May return ≤ limit items** (not guaranteed to fill even if more data exists)
777 /// - **Clients MUST de-duplicate** keys across pages (backends may return duplicates)
778 /// - **No completeness guarantee** under concurrent writes (may miss or duplicate)
779 /// - Without concurrent changes, following every returned cursor until `None`
780 /// visits every matching, unexpired key visible to the backend
781 ///
782 /// **Cursor Behavior**:
783 /// - Opaque string, implementation-specific format
784 /// - Bound to the prefix that created it; using it with another prefix is invalid
785 /// - May become invalid after backend state changes
786 /// - **No TTL guarantees** - can expire without notice
787 /// - Passing invalid cursor should return error, not partial results
788 ///
789 /// **TTL Behavior**: TTL is a soft hint for automatic cleanup. Expired items SHOULD
790 /// be filtered out from results, but physical deletion is eventual and not guaranteed.
791 ///
792 /// **Validation**: Prefix follows same key validation rules.
793 /// Invalid prefix returns `KvError::InvalidKey` immediately.
794 async fn scan_prefix(
795 &self,
796 prefix: &str,
797 limit: Option<usize>,
798 cursor: Option<String>,
799 ) -> Result<ScanResult>;
800}
801
802/// JSON/Text message payload for Queue
803#[derive(Debug, Clone, Serialize, Deserialize)]
804#[serde(tag = "type", rename_all = "lowercase")]
805#[cfg_attr(feature = "openapi", derive(ToSchema))]
806pub enum MessagePayload {
807 /// JSON-serializable value
808 Json(serde_json::Value),
809 /// UTF-8 text payload
810 Text(String),
811}
812
813/// A queue message with payload and receipt handle for acknowledgment
814#[derive(Debug, Clone, Serialize, Deserialize)]
815#[serde(rename_all = "camelCase")]
816#[cfg_attr(feature = "openapi", derive(ToSchema))]
817pub struct QueueMessage {
818 /// JSON-first message payload
819 pub payload: MessagePayload,
820 /// Opaque receipt handle for acknowledgment (backend-specific, short-lived)
821 pub receipt_handle: String,
822 /// Delivery attempt for this message, 1-based (1 = first delivery).
823 ///
824 /// Providers that do not report redelivery counts always set 1; the local
825 /// provider reports the real per-message count so handlers can enforce
826 /// retry limits.
827 #[serde(default = "first_attempt")]
828 pub attempt: u32,
829}
830
831/// Serde default for [`QueueMessage::attempt`]: treat missing counts as the
832/// first delivery.
833fn first_attempt() -> u32 {
834 1
835}
836
837/// Maximum message size in bytes (64 KiB = 65,536 bytes)
838///
839/// This limit ensures compatibility across all queue backends:
840/// - **AWS SQS**: 256KB message limit (much higher, not constraining)
841/// - **Azure Service Bus**: 1MB message limit (much higher, not constraining)
842/// - **GCP Pub/Sub**: 10MB message limit (much higher, not constraining)
843///
844/// The 64KB limit provides:
845/// - Reasonable message sizes for most use cases
846/// - Fast network transfer and low latency
847/// - Consistent behavior across all cloud providers
848/// - Efficient memory usage during batch processing
849pub const MAX_MESSAGE_BYTES: usize = 65_536; // 64 KiB
850
851/// Maximum number of messages per receive call
852///
853/// This limit balances throughput with processing simplicity:
854/// - **AWS SQS**: Supports up to 10 messages per ReceiveMessage call
855/// - **Azure Service Bus**: Can receive multiple messages via prefetch/batching
856/// - **GCP Pub/Sub**: Supports configurable max_messages per Pull request
857///
858/// The 10-message limit ensures:
859/// - Portable batch sizes across all backends
860/// - Manageable memory usage
861/// - Reasonable processing latency per batch
862pub const MAX_BATCH_SIZE: usize = 10;
863
864/// Fixed lease duration in seconds
865///
866/// Messages are leased for exactly 30 seconds after delivery:
867/// - Long enough for most processing tasks
868/// - Short enough to enable fast retry on failures
869/// - Eliminates complexity of dynamic lease management
870/// - Consistent across all platforms
871pub const LEASE_SECONDS: u64 = 30;
872
873/// A trait for queue bindings providing minimal, portable queue operations.
874#[async_trait]
875pub trait Queue: Binding {
876 /// Send a message to the specified queue
877 async fn send(&self, queue: &str, message: MessagePayload) -> Result<()>;
878
879 /// Receive up to `max_messages` (1..=10) from the specified queue
880 async fn receive(&self, queue: &str, max_messages: usize) -> Result<Vec<QueueMessage>>;
881
882 /// Acknowledge a message using its receipt handle (idempotent)
883 async fn ack(&self, queue: &str, receipt_handle: &str) -> Result<()>;
884
885 /// Negative-acknowledge a message: release its lease so it becomes
886 /// immediately available for redelivery, without waiting out the
887 /// visibility timeout. Receipt-handle rules mirror [`Queue::ack`].
888 async fn nack(&self, queue: &str, receipt_handle: &str) -> Result<()>;
889
890 /// Delete every message in the queue, whether visible or in flight.
891 async fn purge(&self, queue: &str) -> Result<()>;
892}
893
894/// Request for invoking a function directly
895#[derive(Debug, Clone, Serialize, Deserialize)]
896#[serde(rename_all = "camelCase")]
897#[cfg_attr(feature = "openapi", derive(ToSchema))]
898pub struct WorkerInvokeRequest {
899 /// Worker identifier (name, ARN, URL, etc.)
900 pub target_worker: String,
901 /// HTTP method
902 pub method: String,
903 /// Request path
904 pub path: String,
905 /// HTTP headers
906 pub headers: BTreeMap<String, String>,
907 /// Request body bytes
908 pub body: Vec<u8>,
909 /// Optional timeout for the invocation
910 pub timeout: Option<Duration>,
911}
912
913/// Response from worker invocation.
914#[derive(Debug, Clone, Serialize, Deserialize)]
915#[serde(rename_all = "camelCase")]
916#[cfg_attr(feature = "openapi", derive(ToSchema))]
917pub struct WorkerInvokeResponse {
918 /// HTTP status code
919 pub status: u16,
920 /// HTTP response headers
921 pub headers: BTreeMap<String, String>,
922 /// Response body bytes
923 pub body: Vec<u8>,
924}
925
926/// A trait for worker bindings that enable direct worker-to-worker calls.
927#[async_trait]
928pub trait Worker: Binding {
929 /// Invoke a worker with HTTP request data.
930 ///
931 /// This enables direct, low-latency worker-to-worker communication within
932 /// the same cloud environment, bypassing Commands for internal calls.
933 ///
934 /// Platform implementations:
935 /// - AWS: Uses InvokeWorker API directly
936 /// - GCP: Calls private service URL directly
937 /// - Azure: Calls private container app URL directly
938 /// - Kubernetes: HTTP call to internal service
939 async fn invoke(&self, request: WorkerInvokeRequest) -> Result<WorkerInvokeResponse>;
940
941 /// Get the public URL of the worker, if available.
942 ///
943 /// Returns the worker's public URL if it exists and is accessible.
944 /// This is useful for exposing public endpoints or getting URLs for
945 /// external integration.
946 ///
947 /// Platform implementations:
948 /// - AWS: Uses GetWorkerUrlConfig API or returns URL from binding
949 /// - GCP: Returns Cloud Run service URL or calls get_service API
950 /// - Azure: Returns Container App URL or calls get_container_app API
951 async fn get_worker_url(&self) -> Result<Option<String>>;
952
953 /// Get a reference to this object as `Any` for dynamic casting
954 fn as_any(&self) -> &dyn std::any::Any;
955}
956
957/// A trait for container bindings that enable container-to-container communication
958#[async_trait]
959pub trait Container: Binding {
960 /// Get the internal URL for container-to-container communication.
961 ///
962 /// This returns the internal service discovery URL that other containers
963 /// in the same network can use to communicate with this container.
964 ///
965 /// Platform implementations:
966 /// - Managed cloud (AWS/GCP/Azure): Returns internal DNS URL (e.g., "http://api.svc:8080")
967 /// - Local (Docker): Returns Docker network DNS URL (e.g., "http://api.svc:3000")
968 fn get_internal_url(&self) -> &str;
969
970 /// Get the public URL of the container, if available.
971 ///
972 /// Returns the container's public URL if it exists and is accessible
973 /// from outside the cluster/network.
974 ///
975 /// Platform implementations:
976 /// - Managed cloud: Returns load balancer URL if exposed publicly
977 /// - Local: Returns localhost URL with mapped port (e.g., "http://localhost:62844")
978 fn get_public_url(&self) -> Option<&str>;
979
980 /// Get the container name/ID.
981 fn get_container_name(&self) -> &str;
982
983 /// Get a reference to this object as `Any` for dynamic casting
984 fn as_any(&self) -> &dyn std::any::Any;
985}
986
987/// A request to create a sandbox.
988#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
989#[cfg_attr(feature = "openapi", derive(ToSchema))]
990#[serde(rename_all = "camelCase", deny_unknown_fields)]
991pub struct CreateSandboxRequest {
992 /// Sandbox id to reconnect to, for the verbs that take one. AWS, Azure and GCP always
993 /// allocate their own on `create` and ignore this; only Local and Kubernetes honor it as the
994 /// new sandbox's id. Read the id back from the response rather than assume the one sent.
995 #[serde(skip_serializing_if = "Option::is_none")]
996 pub sandbox_id: Option<String>,
997 /// Opaque tenant key. Never sent to a provider verbatim — the binding derives a
998 /// fixed-length identifier from it with a deployment-scoped HMAC.
999 #[serde(skip_serializing_if = "Option::is_none")]
1000 pub tenant_key: Option<String>,
1001 /// Environment variables to place in the sandbox.
1002 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1003 pub env: BTreeMap<String, String>,
1004 /// Wall-clock lifetime, after which the platform terminates the sandbox.
1005 ///
1006 /// Requires `sandboxLifetime`; a backend without it refuses rather than accepting a ceiling
1007 /// it would never apply. It only ever shortens: a declared ceiling still bounds the sandbox,
1008 /// so a caller cannot buy itself more life than the deployment allows.
1009 #[serde(default, skip_serializing_if = "Option::is_none")]
1010 pub timeout_ms: Option<u64>,
1011}
1012
1013/// A live sandbox.
1014#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1015#[cfg_attr(feature = "openapi", derive(ToSchema))]
1016#[serde(rename_all = "camelCase")]
1017pub struct SandboxInstance {
1018 /// Provider-scoped sandbox identifier
1019 pub sandbox_id: String,
1020 /// Current lifecycle state
1021 pub state: SandboxState,
1022 /// Lifecycle generation. A capability from another generation is rejected, which is how
1023 /// terminate revokes without distributing a revocation list.
1024 ///
1025 /// Differs by backend: AWS/Azure allocate a fresh id per sandbox, so a constant carries the
1026 /// whole meaning; GCP has none, so this is the guest's boot id — which a snapshot restore
1027 /// reports unchanged, so restore and replacement need the resource name too to tell apart.
1028 pub generation: u64,
1029}
1030
1031/// A sandbox from `get_or_create`, and which of the two things happened.
1032///
1033/// A reconnect and a create are separate paths in every provider; `created` carries that
1034/// distinction out to the caller, which would otherwise make its own first-run work idempotent.
1035#[derive(Debug, Clone, PartialEq, Eq)]
1036pub struct ResolvedSandbox {
1037 /// The sandbox, whether it was made here or found.
1038 pub sandbox: SandboxInstance,
1039 /// Whether this request is what created it.
1040 pub created: bool,
1041}
1042
1043impl ResolvedSandbox {
1044 /// This request created the sandbox.
1045 pub fn created(sandbox: SandboxInstance) -> Self {
1046 Self {
1047 sandbox,
1048 created: true,
1049 }
1050 }
1051
1052 /// An existing sandbox, whatever it took to reach it — a wait, a wake — but not made here.
1053 pub fn found(sandbox: SandboxInstance) -> Self {
1054 Self {
1055 sandbox,
1056 created: false,
1057 }
1058 }
1059}
1060
1061/// Lifecycle state of a sandbox.
1062#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1063#[cfg_attr(feature = "openapi", derive(ToSchema))]
1064#[serde(rename_all = "camelCase")]
1065pub enum SandboxState {
1066 /// Created but not yet able to run a command.
1067 ///
1068 /// A real state, not a placeholder: a MicroVM takes seconds to reach `RUNNING`, and calling
1069 /// that Running would tell a caller to send commands to something that cannot answer.
1070 Starting,
1071 /// Executing, consuming CPU and memory
1072 Running,
1073 /// Paused with state preserved
1074 Paused,
1075 /// Terminated; the id will not run again
1076 Terminated,
1077}
1078
1079/// A command to run inside a sandbox.
1080#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1081#[cfg_attr(feature = "openapi", derive(ToSchema))]
1082#[serde(rename_all = "camelCase")]
1083pub struct RunCommandRequest {
1084 /// The program to run
1085 pub command: String,
1086 /// Arguments handed to the program, each as one argument — nothing re-parses their text
1087 #[serde(default, skip_serializing_if = "Vec::is_empty")]
1088 pub args: Vec<String>,
1089 /// Working directory inside the sandbox
1090 #[serde(skip_serializing_if = "Option::is_none")]
1091 pub cwd: Option<String>,
1092 /// Environment overlaid on the sandbox's own
1093 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
1094 pub env: BTreeMap<String, String>,
1095 /// How long this command may run. Required — a defaulted timeout is a hang waiting for a slow
1096 /// day.
1097 ///
1098 /// `CreateSandboxRequest::timeout_ms` is an outer bound this timeout cannot see. Neither
1099 /// shortens the other, but on AWS and GCP the platform reaps the sandbox at its lifetime
1100 /// whatever is running inside: a command still going is cut off mid-flight and reports
1101 /// `SANDBOX_OUTCOME_UNKNOWN`, never `timeoutExceeded`, because nothing survived to say what it
1102 /// did. A lifetime has to leave room for the longest command it must cover.
1103 ///
1104 /// It bounds the command, not the call, and the call lands just after it. Where the agent
1105 /// supervises the process it kills the process group; where the data plane has no timeout of
1106 /// its own the command runs under `timeout` inside the sandbox. Either way the sandbox stays
1107 /// usable. Only a sandbox that cannot run `timeout` is ended instead, and that call returns
1108 /// once the sandbox is gone.
1109 ///
1110 /// On every backend `timeoutExceeded` is reported only once the command has verifiably
1111 /// stopped — the agent waits for its kill, and where there is no agent the sandbox kills the
1112 /// command itself and says so. It is never reported on a stop that was merely requested: a
1113 /// timeout that leaves untrusted code running is not a timeout.
1114 ///
1115 /// What stops is the command and its process group. A descendant that detaches itself into a
1116 /// session of its own is beyond any signal sent from inside, on every backend; it is bounded
1117 /// by the sandbox, which ends on terminate or at its own lifetime ceiling.
1118 pub timeout: Duration,
1119}
1120
1121impl RunCommandRequest {
1122 /// The program followed by its arguments, which is the shape every backend's exec takes.
1123 pub fn argv(&self) -> Vec<String> {
1124 std::iter::once(self.command.clone())
1125 .chain(self.args.iter().cloned())
1126 .collect()
1127 }
1128}
1129
1130/// One frame of a running command's output.
1131#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1132#[cfg_attr(feature = "openapi", derive(ToSchema))]
1133#[serde(rename_all = "camelCase", tag = "type")]
1134pub enum CommandOutput {
1135 /// Bytes written to stdout
1136 #[serde(rename_all = "camelCase")]
1137 Stdout {
1138 /// Monotonic across both streams, so a caller can interleave them in production order
1139 seq: u64,
1140 /// Raw bytes; command output is not necessarily UTF-8
1141 data: Vec<u8>,
1142 },
1143 /// Bytes written to stderr
1144 #[serde(rename_all = "camelCase")]
1145 Stderr {
1146 /// Monotonic across both streams
1147 seq: u64,
1148 /// Raw bytes
1149 data: Vec<u8>,
1150 },
1151 /// The command finished. Exactly one terminal frame is emitted, always last.
1152 #[serde(rename_all = "camelCase")]
1153 Exit {
1154 /// Process exit code
1155 code: i32,
1156 /// Set when output was cut short by a bound rather than by the command finishing
1157 #[serde(default)]
1158 truncated: bool,
1159 },
1160}
1161
1162/// The id a started job answers to.
1163#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1164#[cfg_attr(feature = "openapi", derive(ToSchema))]
1165#[serde(rename_all = "camelCase")]
1166pub struct JobStart {
1167 /// Identifier every later poll and cancel addresses
1168 pub job_id: String,
1169}
1170
1171/// A job's output so far, and how it ended once it has.
1172#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1173#[cfg_attr(feature = "openapi", derive(ToSchema))]
1174#[serde(rename_all = "camelCase")]
1175pub struct JobPoll {
1176 /// Whether the command is still running
1177 pub running: bool,
1178 /// Output produced after the polled sequence. The ending is `exit` or `error`, never a frame.
1179 pub frames: Vec<CommandOutput>,
1180 /// How the command exited, once it has
1181 #[serde(skip_serializing_if = "Option::is_none")]
1182 pub exit: Option<JobExit>,
1183 /// Why the command ended without exiting
1184 #[serde(skip_serializing_if = "Option::is_none")]
1185 pub error: Option<JobError>,
1186}
1187
1188/// How a job's command exited.
1189#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
1190#[cfg_attr(feature = "openapi", derive(ToSchema))]
1191#[serde(rename_all = "camelCase")]
1192pub struct JobExit {
1193 /// Process exit code
1194 pub code: i32,
1195 /// Set when output was cut short by a bound rather than by the command finishing
1196 #[serde(default)]
1197 pub truncated: bool,
1198}
1199
1200/// Why a job ended without its command exiting.
1201#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1202#[cfg_attr(feature = "openapi", derive(ToSchema))]
1203#[serde(rename_all = "camelCase")]
1204pub struct JobError {
1205 /// Machine-readable cause, e.g. `timeoutExceeded`
1206 pub code: String,
1207 /// Human-readable detail
1208 pub message: String,
1209}
1210
1211/// An authenticated, port-scoped capability to reach a service inside a sandbox.
1212///
1213/// Not a URL string: AWS needs a JWE and a port header, Azure an Entra token, and a bare
1214/// string cannot carry either. Returning one would push callers into building the request
1215/// themselves and getting the auth wrong.
1216#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
1217#[cfg_attr(feature = "openapi", derive(ToSchema))]
1218#[serde(rename_all = "camelCase")]
1219pub struct PreviewCapability {
1220 /// Endpoint the request must be sent to
1221 pub endpoint: String,
1222 /// Headers that must accompany every request
1223 pub headers: BTreeMap<String, String>,
1224 /// Ports this capability admits. A request to any other port is refused upstream.
1225 pub allowed_ports: Vec<u16>,
1226 /// Seconds until the capability expires
1227 pub expires_in_seconds: u64,
1228}
1229
1230/// A sandbox binding: create sandboxes, run untrusted code in them, and tear them down.
1231///
1232/// Capabilities differ per platform. Call `capabilities()` and branch, or call and handle the
1233/// typed error — an unsupported capability is never a silent no-op.
1234#[async_trait]
1235pub trait Sandbox: Binding {
1236 /// What this platform's backend supports.
1237 fn capabilities(&self) -> alien_core::SandboxCapabilities;
1238
1239 /// Creates a sandbox that can already take work.
1240 ///
1241 /// Returning before the sandbox can serve pushes a readiness poll into every caller for a
1242 /// condition only the backend can observe, and the resulting race fails a fraction of the
1243 /// time rather than every time. A backend whose start API returns early waits here.
1244 async fn create(&self, request: CreateSandboxRequest) -> Result<SandboxInstance>;
1245
1246 /// Fetches a sandbox by id, or `None` if it does not exist.
1247 ///
1248 /// Requires `reconnect`. `None` means the sandbox does not exist; a backend that cannot
1249 /// answer the question returns the typed error instead, so an absent sandbox and an
1250 /// unreachable one are never the same result.
1251 async fn get(&self, sandbox_id: &str) -> Result<Option<SandboxInstance>>;
1252
1253 /// Fetches a sandbox, creating it if absent, and reports which it did.
1254 ///
1255 /// Answer `created` from the path taken, never from a timestamp: the clocks are the provider's,
1256 /// not ours. No backend offers an atomic create-if-absent, so two callers naming one id can
1257 /// both be told they created it; a caller that cannot tolerate its setup running twice needs
1258 /// its own lock.
1259 ///
1260 /// `timeoutMs` bounds a sandbox this call creates. No backend can move a running sandbox's
1261 /// deadline, so one that is found keeps the lifetime it was created with, and `created` is
1262 /// how a caller tells the two apart.
1263 async fn get_or_create(&self, request: CreateSandboxRequest) -> Result<ResolvedSandbox>;
1264
1265 /// Lists the sandboxes belonging to this binding's parent.
1266 ///
1267 /// Offered only where the backend has a verb for it and the grant covers it — GCP lists
1268 /// under its engine. AWS and Azure raise `OperationNotSupported`: enumerating there costs an
1269 /// account-wide grant the sandbox role deliberately withholds. Reaching a known id is `get`.
1270 async fn list(&self) -> Result<Vec<SandboxInstance>>;
1271
1272 /// Runs a command, streaming output frames until exactly one terminal frame.
1273 ///
1274 /// The stream carries backpressure: a consumer that stops reading stops the sandbox's
1275 /// writer, rather than buffering without bound.
1276 async fn run_command(
1277 &self,
1278 sandbox_id: &str,
1279 request: RunCommandRequest,
1280 ) -> Result<futures::stream::BoxStream<'static, Result<CommandOutput>>>;
1281
1282 /// Starts a command as a job, which outlives the call that started it. Requires `jobs`.
1283 ///
1284 /// A start that goes unanswered is `SANDBOX_OUTCOME_UNKNOWN` and not retryable: the sandbox
1285 /// may have taken the command, and repeating it would run it twice.
1286 async fn start_job(&self, sandbox_id: &str, request: RunCommandRequest) -> Result<JobStart>;
1287
1288 /// Reads a job's output after `since_seq`, and its ending once it has one. Requires `jobs`.
1289 ///
1290 /// `None` reads from the first frame. Repeating a poll costs nothing and changes nothing, so
1291 /// a sandbox that cannot be reached is `SANDBOX_UNREACHABLE` and retryable.
1292 async fn poll_job(
1293 &self,
1294 sandbox_id: &str,
1295 job_id: &str,
1296 since_seq: Option<u64>,
1297 ) -> Result<JobPoll>;
1298
1299 /// Cancels a job, stopping its command. Requires `jobs`.
1300 ///
1301 /// Classified like `poll_job`: a cancel that is repeated stops nothing a second time.
1302 async fn cancel_job(&self, sandbox_id: &str, job_id: &str) -> Result<()>;
1303
1304 /// Reads a file out of the sandbox. Requires `files`. Paths are normalised and may not
1305 /// escape the root.
1306 async fn read_file(&self, sandbox_id: &str, path: &str) -> Result<Vec<u8>>;
1307
1308 /// Writes files into the sandbox. Requires `files`. Parent directories are created as needed.
1309 async fn write_files(&self, sandbox_id: &str, files: BTreeMap<String, Vec<u8>>) -> Result<()>;
1310
1311 /// Mints a capability to reach a declared port. Requires `preview`.
1312 async fn preview(&self, sandbox_id: &str, port: u16) -> Result<PreviewCapability>;
1313
1314 /// Pauses a sandbox, preserving state. Requires `pauseResume`.
1315 ///
1316 /// A command already running is frozen with the guest rather than drained or stopped, and its
1317 /// own deadline freezes with it, so nothing inside the sandbox will end it. A caller streaming
1318 /// that command's output is failed within a bound rather than held for the length of the pause.
1319 /// On Azure the in-guest deadline is what stops an overrunning command, so one frozen by the
1320 /// pause is ended by terminating the sandbox instead and a later `resume` finds nothing.
1321 async fn pause(&self, sandbox_id: &str) -> Result<()>;
1322
1323 /// Resumes a paused sandbox. Requires `pauseResume`.
1324 async fn resume(&self, sandbox_id: &str) -> Result<()>;
1325
1326 /// Captures full sandbox state and returns its identifier. Requires `snapshot`.
1327 async fn snapshot(&self, sandbox_id: &str) -> Result<String>;
1328
1329 /// Terminates a sandbox. Idempotent: terminating an absent sandbox succeeds.
1330 async fn terminate(&self, sandbox_id: &str) -> Result<()>;
1331
1332 /// Get a reference to this object as `Any` for dynamic casting
1333 fn as_any(&self) -> &dyn std::any::Any;
1334}
1335
1336/// A provider must implement methods to load the various types of bindings
1337/// based on environment variables or other configuration sources.
1338#[async_trait]
1339pub trait BindingsProviderApi: Send + Sync + std::fmt::Debug {
1340 /// Given a binding identifier, builds a Storage implementation.
1341 async fn load_storage(&self, binding_name: &str) -> Result<Arc<dyn Storage>>;
1342
1343 /// Given a binding identifier, builds a Key implementation.
1344 async fn load_key(&self, binding_name: &str) -> Result<Arc<dyn Key>> {
1345 Err(alien_error::AlienError::new(
1346 crate::error::ErrorData::OperationNotSupported {
1347 operation: "load_key".to_string(),
1348 reason: format!("Key resource '{binding_name}' is not supported by this provider"),
1349 },
1350 ))
1351 }
1352
1353 /// Given a binding identifier, builds a Build implementation.
1354 async fn load_build(&self, binding_name: &str) -> Result<Arc<dyn Build>>;
1355
1356 /// Given a binding identifier, builds an ArtifactRegistry implementation.
1357 async fn load_artifact_registry(&self, binding_name: &str)
1358 -> Result<Arc<dyn ArtifactRegistry>>;
1359
1360 /// Given a binding identifier, builds a Vault implementation.
1361 async fn load_vault(&self, binding_name: &str) -> Result<Arc<dyn Vault>>;
1362
1363 /// Given a binding identifier, builds a KV implementation.
1364 async fn load_kv(&self, binding_name: &str) -> Result<Arc<dyn Kv>>;
1365
1366 /// Given a binding identifier, builds a Postgres implementation.
1367 ///
1368 /// Every backend is resolved here. Local and External carry their password inline;
1369 /// Aurora, Cloud SQL, and Azure Flexible Server carry only a locator for it and read
1370 /// the value from that cloud's secret store during this call, using the workload's own
1371 /// identity. Resolution happens once per load, so the returned handle is synchronous.
1372 async fn load_postgres(&self, binding_name: &str) -> Result<Arc<dyn Postgres>>;
1373
1374 /// Given a binding identifier, builds a Queue implementation.
1375 async fn load_queue(&self, binding_name: &str) -> Result<Arc<dyn Queue>>;
1376
1377 /// Given a binding identifier, builds a Worker implementation.
1378 async fn load_worker(&self, binding_name: &str) -> Result<Arc<dyn Worker>>;
1379
1380 /// Given a binding identifier, builds a Container implementation.
1381 async fn load_container(&self, binding_name: &str) -> Result<Arc<dyn Container>>;
1382
1383 /// Given a binding identifier, builds a ServiceAccount implementation.
1384 async fn load_service_account(&self, binding_name: &str) -> Result<Arc<dyn ServiceAccount>>;
1385
1386 /// Given a binding identifier, builds a Sandbox implementation.
1387 async fn load_sandbox(&self, binding_name: &str) -> Result<Arc<dyn Sandbox>>;
1388
1389 /// Runtime-only binding env vars (a local Postgres connection with its password, a local
1390 /// BYO-key AI binding) for the given resource — re-resolved on every (re)start so the secret
1391 /// reaches the worker process but is never written to persisted worker metadata. The resource
1392 /// type routes resolution to the right local source. Default `None`: cloud providers carry a
1393 /// secret locator (not a raw secret) and use the normal persisted path.
1394 async fn resolve_runtime_only_binding_env(
1395 &self,
1396 _binding_name: &str,
1397 _resource_type: &str,
1398 ) -> Result<Option<std::collections::HashMap<String, String>>> {
1399 Ok(None)
1400 }
1401}
1402
1403#[cfg(test)]
1404mod tests {
1405 use super::CreateSandboxRequest;
1406
1407 /// A field this struct does not know is a caller asking for something it will not get, and
1408 /// the field it would most want is the one that bounds the sandbox. Read as "absent",
1409 /// `timeoutMs` misspelled is an unbounded sandbox reported as a bounded one.
1410 #[test]
1411 fn a_create_request_refuses_a_field_it_does_not_know() {
1412 for body in [
1413 r#"{"sessionId":"s1"}"#,
1414 r#"{"timeoutMillis":60000}"#,
1415 r#"{"timeout_ms":60000}"#,
1416 r#"{"workingDirectory":"/work"}"#,
1417 ] {
1418 let error = serde_json::from_str::<CreateSandboxRequest>(body)
1419 .expect_err("a field the request does not declare must be refused, not dropped");
1420 assert!(
1421 error.to_string().contains("unknown field"),
1422 "{body} was refused for the wrong reason: {error}"
1423 );
1424 }
1425 }
1426
1427 #[test]
1428 fn a_create_request_takes_every_field_it_declares() {
1429 let parsed: CreateSandboxRequest = serde_json::from_str(
1430 r#"{"sandboxId":"s1","tenantKey":"t","env":{"A":"b"},"timeoutMs":60000}"#,
1431 )
1432 .expect("the declared fields parse");
1433
1434 assert_eq!(parsed.sandbox_id.as_deref(), Some("s1"));
1435 assert_eq!(parsed.tenant_key.as_deref(), Some("t"));
1436 assert_eq!(parsed.env.get("A").map(String::as_str), Some("b"));
1437 assert_eq!(parsed.timeout_ms, Some(60_000));
1438 }
1439}