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