Skip to main content

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}