Skip to main content

alien_bindings/
bindings.rs

1//! App-facing convenience API for accessing bindings.
2//!
3//! [`Bindings`] wraps a [`crate::provider::LazyEnvBindingsProvider`], giving application
4//! code a small, stable surface — `storage`, `kv`, `queue`, `vault`, `container`,
5//! `postgres`, `sandbox` — instead of the full [`crate::traits::BindingsProviderApi`] used internally
6//! by the manager and controllers.
7
8use crate::error::Result;
9use crate::provider::{BindingsProvider, LazyEnvBindingsProvider};
10use crate::refreshing::{
11    RefreshingKey, RefreshingKv, RefreshingQueue, RefreshingStorage, RefreshingVault,
12};
13use crate::traits::{
14    BindingsProviderApi, Container, Key, Kv, MessagePayload, Postgres, Queue, QueueMessage,
15    Sandbox, Storage, Vault,
16};
17use std::collections::HashMap;
18use std::sync::Arc;
19
20/// App-facing entry point for environment-backed bindings.
21///
22/// Construction is synchronous and only validates each configured binding's JSON
23/// shape (see [`BindingsProvider::from_env_deferred`]); the deployment platform,
24/// cloud client configuration, and each binding's backing client are resolved
25/// lazily, on first use. A first operation against a binding that is not
26/// configured reports `BINDING_NOT_CONFIGURED` before any platform resolution, so
27/// a zero-environment process still constructs and fails cleanly.
28///
29/// # Examples
30///
31/// This is the canonical usage for a Container/Daemon-shaped app (a long-running
32/// resident process that only needs bindings, with no Worker event handlers):
33///
34/// ```no_run
35/// use alien_bindings::Bindings;
36/// use object_store::{path::Path, PutPayload};
37///
38/// # async fn run() -> Result<(), Box<dyn std::error::Error>> {
39/// let bindings = Bindings::from_env()?;
40///
41/// let storage = bindings.storage("files").await?;
42/// storage
43///     .put(&Path::from("greeting.txt"), PutPayload::from_static(b"hello"))
44///     .await?;
45/// # Ok(())
46/// # }
47/// ```
48#[derive(Debug)]
49pub struct Bindings {
50    provider: Arc<LazyEnvBindingsProvider>,
51}
52
53/// A queue binding scoped to its configured queue name.
54#[derive(Clone)]
55pub struct BoundQueue {
56    inner: Arc<dyn Queue>,
57    name: Arc<str>,
58}
59
60impl std::fmt::Debug for BoundQueue {
61    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
62        formatter
63            .debug_struct("Queue")
64            .field("name", &self.name)
65            .finish_non_exhaustive()
66    }
67}
68
69impl BoundQueue {
70    fn new(inner: Arc<dyn Queue>, name: impl Into<Arc<str>>) -> Self {
71        Self {
72            inner,
73            name: name.into(),
74        }
75    }
76
77    /// Send a message to this queue.
78    pub async fn send(&self, message: MessagePayload) -> Result<()> {
79        self.inner.send(&self.name, message).await
80    }
81
82    /// Receive up to `max_messages` messages from this queue.
83    pub async fn receive(&self, max_messages: usize) -> Result<Vec<QueueMessage>> {
84        self.inner.receive(&self.name, max_messages).await
85    }
86
87    /// Acknowledge a received message.
88    pub async fn ack(&self, receipt_handle: &str) -> Result<()> {
89        self.inner.ack(&self.name, receipt_handle).await
90    }
91
92    /// Release a received message for redelivery.
93    pub async fn nack(&self, receipt_handle: &str) -> Result<()> {
94        self.inner.nack(&self.name, receipt_handle).await
95    }
96
97    /// Delete every message in this queue.
98    pub async fn purge(&self) -> Result<()> {
99        self.inner.purge(&self.name).await
100    }
101}
102
103impl Bindings {
104    /// Sync-constructs `Bindings` from the current process environment.
105    pub fn from_env() -> Result<Self> {
106        Self::from_env_map(std::env::vars().collect())
107    }
108
109    /// Sync-constructs `Bindings` from an explicit environment map instead of the process
110    /// environment.
111    ///
112    /// This is public for embedders that resolve bindings from a caller-supplied map rather
113    /// than `std::env` — notably the napi addon, which merges `std::env::vars()` with
114    /// per-call overrides before constructing `Bindings`. It is also what `from_env`
115    /// delegates to and what this module's tests use to inject `ALIEN_*_BINDING` variables
116    /// (avoiding process-global state that's unsafe to share across parallel tests).
117    pub fn from_env_map(env: HashMap<String, String>) -> Result<Self> {
118        Ok(Self {
119            provider: Arc::new(BindingsProvider::from_env_deferred(env)?),
120        })
121    }
122
123    /// Loads the object storage binding named `binding_name`.
124    ///
125    /// The returned handle checks credential freshness before each operation.
126    /// Native credentials and fresh short-lived credentials remain cached; a
127    /// provider inside its refresh window is refreshed once under its shared
128    /// resolver's single-flight guard.
129    pub async fn storage(&self, binding_name: &str) -> Result<Arc<dyn Storage>> {
130        let initial = self.provider.load_storage(binding_name).await?;
131        Ok(Arc::new(RefreshingStorage::new(
132            self.provider.clone(),
133            binding_name.to_string(),
134            initial,
135        )))
136    }
137
138    /// Loads a provider-backed key binding that refreshes minted credentials before use.
139    pub async fn key(&self, binding_name: &str) -> Result<Arc<dyn Key>> {
140        self.provider.load_key(binding_name).await?;
141        Ok(Arc::new(RefreshingKey::new(
142            self.provider.clone(),
143            binding_name.to_string(),
144        )))
145    }
146
147    /// Loads an environment-backed key-value binding that refreshes minted
148    /// credentials before use.
149    pub async fn kv(&self, binding_name: &str) -> Result<Arc<dyn Kv>> {
150        self.provider.load_kv(binding_name).await?;
151        Ok(Arc::new(RefreshingKv::new(
152            self.provider.clone(),
153            binding_name.to_string(),
154        )))
155    }
156
157    /// Loads a queue binding that refreshes minted credentials before use.
158    pub async fn queue(&self, binding_name: &str) -> Result<BoundQueue> {
159        self.provider.load_queue(binding_name).await?;
160        let queue: Arc<dyn Queue> = Arc::new(RefreshingQueue::new(
161            self.provider.clone(),
162            binding_name.to_string(),
163        ));
164        Ok(BoundQueue::new(queue, binding_name))
165    }
166
167    /// Loads an environment-backed vault binding that refreshes minted
168    /// credentials before use.
169    pub async fn vault(&self, binding_name: &str) -> Result<Arc<dyn Vault>> {
170        self.provider.load_vault(binding_name).await?;
171        Ok(Arc::new(RefreshingVault::new(
172            self.provider.clone(),
173            binding_name.to_string(),
174        )))
175    }
176
177    /// Loads a linked container for read-only service discovery.
178    pub async fn container(&self, binding_name: &str) -> Result<Arc<dyn Container>> {
179        self.provider.load_container(binding_name).await
180    }
181
182    /// Loads the connection details for a linked Postgres database.
183    ///
184    /// Unlike the other kinds this returns no operations — Postgres has no gRPC service
185    /// and every backend speaks the same wire protocol, so the handle carries connection
186    /// details and the application connects with its own driver.
187    ///
188    /// The Local and External backends carry their password inline in the binding
189    /// environment variable, so their handle is resolved once and then cached.
190    ///
191    /// The three cloud backends carry only a locator for their password and read it from
192    /// the cloud secret store on **every** call — their handle is never cached. There is
193    /// no refreshing wrapper, so a handle keeps the password that was current when it was
194    /// created and calling this again is what picks up a rotated one. Each call is
195    /// therefore one secret-store read: hold the returned handle for the lifetime of a
196    /// connection pool rather than calling this per query.
197    pub async fn postgres(&self, binding_name: &str) -> Result<Arc<dyn Postgres>> {
198        self.provider.load_postgres(binding_name).await
199    }
200
201    /// Loads a linked sandbox for running untrusted code.
202    ///
203    /// No refreshing wrapper: a sandbox handle addresses a control plane rather than holding
204    /// data-plane credentials, and each sandbox capability is minted per call.
205    pub async fn sandbox(&self, binding_name: &str) -> Result<Arc<dyn Sandbox>> {
206        self.provider.load_sandbox(binding_name).await
207    }
208}
209
210#[cfg(test)]
211mod tests {
212    use super::*;
213    use std::net::SocketAddr;
214    use std::sync::atomic::{AtomicUsize, Ordering};
215
216    use crate::error::binding_env_var;
217    use crate::traits::MessagePayload;
218    use alien_core::{
219        Platform, ENV_ALIEN_DEPLOYMENT_ID, ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT,
220        ENV_ALIEN_DEPLOYMENT_TOKEN, ENV_ALIEN_DEPLOYMENT_TYPE, ENV_ALIEN_MANAGER_URL,
221        ENV_ALIEN_RESOURCE_ID,
222    };
223    use axum::{extract::State, routing::post, Json, Router};
224    use object_store::{path::Path as ObjectPath, PutPayload};
225    use std::collections::HashMap;
226    use tempfile::TempDir;
227
228    /// Minimal valid environment (no bindings configured yet).
229    fn base_env() -> HashMap<String, String> {
230        HashMap::from([(
231            ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
232            Platform::Local.as_str().to_string(),
233        )])
234    }
235
236    fn with_binding(
237        mut env: HashMap<String, String>,
238        binding_name: &str,
239        json: &str,
240    ) -> HashMap<String, String> {
241        env.insert(binding_env_var(binding_name), json.to_string());
242        env
243    }
244
245    #[derive(Clone)]
246    struct MintServerState {
247        calls: Arc<AtomicUsize>,
248        state_directory: String,
249    }
250
251    async fn mint_handler(State(state): State<MintServerState>) -> Json<serde_json::Value> {
252        let call = state.calls.fetch_add(1, Ordering::SeqCst) + 1;
253        let lifetime_seconds = if call == 1 { 120 } else { 3600 };
254        let expires_at =
255            (chrono::Utc::now() + chrono::Duration::seconds(lifetime_seconds)).to_rfc3339();
256        Json(serde_json::json!({
257            "clientConfig": {
258                "platform": "local",
259                "state_directory": state.state_directory,
260            },
261            "expiresAt": expires_at,
262            "principal": "local:refreshing-binding-test",
263        }))
264    }
265
266    async fn spawn_mint_server(state_directory: &str) -> (String, Arc<AtomicUsize>) {
267        let calls = Arc::new(AtomicUsize::new(0));
268        let app = Router::new()
269            .route("/v1/credentials/mint", post(mint_handler))
270            .with_state(MintServerState {
271                calls: calls.clone(),
272                state_directory: state_directory.to_string(),
273            });
274        let listener = tokio::net::TcpListener::bind(SocketAddr::from(([127, 0, 0, 1], 0)))
275            .await
276            .expect("bind fake mint server");
277        let address = listener.local_addr().expect("read fake server address");
278        tokio::spawn(async move {
279            axum::serve(listener, app)
280                .await
281                .expect("serve fake mint endpoint");
282        });
283        (format!("http://{address}"), calls)
284    }
285
286    fn mint_env(manager_url: &str) -> HashMap<String, String> {
287        HashMap::from([
288            (
289                ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
290                Platform::Aws.as_str().to_string(),
291            ),
292            ("AWS_EC2_METADATA_DISABLED".to_string(), "true".to_string()),
293            (
294                "AWS_PROFILE".to_string(),
295                "__alien_missing_refresh_test_profile__".to_string(),
296            ),
297            (ENV_ALIEN_MANAGER_URL.to_string(), manager_url.to_string()),
298            (
299                ENV_ALIEN_DEPLOYMENT_TOKEN.to_string(),
300                "refresh-test-token".to_string(),
301            ),
302            (
303                ENV_ALIEN_DEPLOYMENT_ID.to_string(),
304                "refresh-test-deployment".to_string(),
305            ),
306            (
307                ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT.to_string(),
308                "refresh-test-service-account".to_string(),
309            ),
310            (
311                ENV_ALIEN_RESOURCE_ID.to_string(),
312                "refresh-test-resource".to_string(),
313            ),
314        ])
315    }
316
317    #[test]
318    fn from_env_map_constructs_synchronously_from_injected_env() {
319        // No `.await` here at all: proves construction is a plain sync function,
320        // not something that merely returns a Future.
321        let bindings =
322            Bindings::from_env_map(base_env()).expect("valid env should construct Bindings");
323        drop(bindings);
324    }
325
326    #[tokio::test]
327    async fn storage_delegates_to_local_provider_and_performs_real_io() {
328        let temp_dir = TempDir::new().expect("tempdir");
329        let json = format!(
330            r#"{{"service":"local-storage","storagePath":"{}"}}"#,
331            temp_dir.path().display()
332        );
333        let env = with_binding(base_env(), "files", &json);
334        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
335
336        let storage = bindings
337            .storage("files")
338            .await
339            .expect("storage binding should load");
340
341        let path = ObjectPath::from("greeting.txt");
342        storage
343            .put(&path, PutPayload::from(bytes::Bytes::from_static(b"hello")))
344            .await
345            .expect("put should succeed");
346        let fetched = storage
347            .get(&path)
348            .await
349            .expect("get should succeed")
350            .bytes()
351            .await
352            .expect("reading bytes should succeed");
353        assert_eq!(fetched.as_ref(), b"hello");
354    }
355
356    #[tokio::test]
357    async fn kv_delegates_to_local_provider_and_performs_real_io() {
358        let temp_dir = TempDir::new().expect("tempdir");
359        let json = format!(
360            r#"{{"service":"local-kv","dataDir":"{}"}}"#,
361            temp_dir.path().display()
362        );
363        let env = with_binding(base_env(), "cache", &json);
364        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
365
366        let kv = bindings.kv("cache").await.expect("kv binding should load");
367
368        kv.put("greeting", b"hi".to_vec(), None)
369            .await
370            .expect("put should succeed");
371        let value = kv
372            .get("greeting")
373            .await
374            .expect("get should succeed")
375            .expect("value should exist");
376        assert_eq!(value.value, b"hi");
377    }
378
379    #[tokio::test]
380    async fn long_lived_kv_handle_refreshes_minted_provider_before_expiry() {
381        let temp_dir = TempDir::new().expect("tempdir");
382        let (manager_url, calls) = spawn_mint_server(
383            temp_dir
384                .path()
385                .to_str()
386                .expect("tempdir path must be valid UTF-8"),
387        )
388        .await;
389        let json = format!(
390            r#"{{"service":"local-kv","dataDir":"{}"}}"#,
391            temp_dir.path().display()
392        );
393        let env = with_binding(mint_env(&manager_url), "cache", &json);
394        let bindings = Bindings::from_env_map(env).expect("minting env should construct Bindings");
395
396        let kv = bindings
397            .kv("cache")
398            .await
399            .expect("first binding resolution should mint credentials");
400        assert_eq!(calls.load(Ordering::SeqCst), 1);
401
402        kv.put("greeting", b"hi".to_vec(), None)
403            .await
404            .expect("the long-lived handle should refresh and write");
405        assert_eq!(
406            calls.load(Ordering::SeqCst),
407            2,
408            "the first mint is still unexpired but inside the refresh window"
409        );
410
411        let value = kv
412            .get("greeting")
413            .await
414            .expect("the same long-lived handle should read")
415            .expect("value should exist");
416        assert_eq!(value.value, b"hi");
417        assert_eq!(
418            calls.load(Ordering::SeqCst),
419            2,
420            "the refreshed provider should stay cached while fresh"
421        );
422    }
423
424    #[tokio::test]
425    async fn queue_delegates_to_local_provider_and_performs_real_io() {
426        let temp_dir = TempDir::new().expect("tempdir");
427        let json = format!(
428            r#"{{"service":"local-queue","queuePath":"{}"}}"#,
429            temp_dir.path().join("queue.db").display()
430        );
431        let env = with_binding(base_env(), "jobs", &json);
432        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
433
434        let queue = bindings
435            .queue("jobs")
436            .await
437            .expect("queue binding should load");
438
439        queue
440            .send(MessagePayload::Text("hello".to_string()))
441            .await
442            .expect("send should succeed");
443        let messages = queue.receive(1).await.expect("receive should succeed");
444        assert_eq!(messages.len(), 1);
445    }
446
447    #[tokio::test]
448    async fn bound_queue_uses_its_configured_name_for_every_operation() {
449        let temp_dir = TempDir::new().expect("tempdir");
450        let json = format!(
451            r#"{{"service":"local-queue","queuePath":"{}"}}"#,
452            temp_dir.path().join("queue.db").display()
453        );
454        let env = with_binding(base_env(), "jobs", &json);
455        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
456
457        let queue = bindings
458            .queue("jobs")
459            .await
460            .expect("queue binding should load");
461
462        // nack: an in-flight message under the default lease is hidden, but a
463        // nack makes it immediately redeliverable.
464        queue
465            .send(MessagePayload::Text("retry".to_string()))
466            .await
467            .expect("send should succeed");
468        let first = queue.receive(1).await.expect("receive should succeed");
469        assert_eq!(first.len(), 1);
470        assert!(
471            queue
472                .receive(1)
473                .await
474                .expect("receive should succeed")
475                .is_empty(),
476            "in-flight message must be hidden before nack"
477        );
478        queue
479            .nack(&first[0].receipt_handle)
480            .await
481            .expect("nack should succeed");
482        let redelivered = queue.receive(1).await.expect("receive should succeed");
483        assert_eq!(redelivered.len(), 1, "nacked message must be redelivered");
484
485        // purge: clears everything, in flight or visible.
486        queue.purge().await.expect("purge should succeed");
487        assert!(
488            queue
489                .receive(1)
490                .await
491                .expect("receive should succeed")
492                .is_empty(),
493            "purge must empty the queue"
494        );
495    }
496
497    #[tokio::test]
498    async fn container_exposes_internal_and_optional_public_urls() {
499        let env = with_binding(
500            base_env(),
501            "database",
502            r#"{"service":"local","containerName":"database","internalUrl":"http://database.internal:5432","publicUrl":"http://localhost:15432"}"#,
503        );
504        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
505
506        let container = bindings
507            .container("database")
508            .await
509            .expect("container binding should load");
510
511        assert_eq!(
512            container.get_internal_url(),
513            "http://database.internal:5432"
514        );
515        assert_eq!(container.get_public_url(), Some("http://localhost:15432"));
516    }
517
518    #[tokio::test]
519    async fn vault_delegates_to_local_provider_and_performs_real_io() {
520        let temp_dir = TempDir::new().expect("tempdir");
521        let json = format!(
522            r#"{{"service":"local-vault","vaultName":"secrets","dataDir":"{}"}}"#,
523            temp_dir.path().display()
524        );
525        let env = with_binding(base_env(), "secrets", &json);
526        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
527
528        let vault = bindings
529            .vault("secrets")
530            .await
531            .expect("vault binding should load");
532
533        vault
534            .set_secret("api-key", "sekrit")
535            .await
536            .expect("set_secret should succeed");
537        let value = vault
538            .get_secret("api-key")
539            .await
540            .expect("get_secret should succeed");
541        assert_eq!(value, "sekrit");
542
543        // list_secrets must be reachable through the `Arc<dyn Vault>` surface
544        // and return the stored names.
545        vault
546            .set_secret("db-url", "postgres://…")
547            .await
548            .expect("set_secret should succeed");
549        let mut names = vault
550            .list_secrets()
551            .await
552            .expect("list_secrets should succeed");
553        names.sort();
554        assert_eq!(names, vec!["api-key".to_string(), "db-url".to_string()]);
555    }
556
557    #[tokio::test]
558    async fn missing_storage_binding_returns_binding_not_configured() {
559        let bindings = Bindings::from_env_map(base_env())
560            .expect("construction should succeed with no bindings configured");
561
562        let error = bindings
563            .storage("files")
564            .await
565            .expect_err("missing binding should error");
566
567        assert_eq!(error.code, "BINDING_NOT_CONFIGURED");
568        assert!(
569            error.to_string().contains("ALIEN_FILES_BINDING"),
570            "message should name the env var, got: {error}"
571        );
572    }
573
574    #[tokio::test]
575    async fn zero_env_construct_then_missing_binding_is_binding_not_configured() {
576        // The app-facing contract: with NO deployment type and NO credentials,
577        // construction must succeed and the FIRST op on a missing binding must
578        // report BINDING_NOT_CONFIGURED (naming ALIEN_<NAME>_BINDING) BEFORE any
579        // platform / client-config resolution. There is deliberately no
580        // ALIEN_DEPLOYMENT_TYPE in this environment. Table test over all four
581        // app-facing kinds so a future kind added to `Bindings` without wiring
582        // `ensure_binding_present` into its `load_*` method fails this test
583        // instead of silently regressing to ENVIRONMENT_VARIABLE_MISSING.
584        for kind in ["storage", "kv", "queue", "vault"] {
585            let bindings = Bindings::from_env_map(HashMap::new())
586                .expect("zero-env construction must succeed (platform resolution deferred)");
587
588            let error = match kind {
589                "storage" => bindings.storage("x").await.unwrap_err(),
590                "kv" => bindings.kv("x").await.unwrap_err(),
591                "queue" => bindings.queue("x").await.unwrap_err(),
592                "vault" => bindings.vault("x").await.unwrap_err(),
593                other => unreachable!("unhandled kind in table test: {other}"),
594            };
595
596            assert_eq!(
597                error.code, "BINDING_NOT_CONFIGURED",
598                "{kind}: expected the missing-binding error, not a platform/deployment error: {error}"
599            );
600            assert!(
601                error.to_string().contains("ALIEN_X_BINDING"),
602                "{kind}: message should name the env var, got: {error}"
603            );
604        }
605    }
606
607    #[test]
608    fn malformed_binding_json_returns_binding_config_invalid_naming_env_var() {
609        let env = with_binding(base_env(), "files", "not-json");
610
611        let error =
612            Bindings::from_env_map(env).expect_err("malformed binding JSON should fail to load");
613
614        assert_eq!(error.code, "BINDING_CONFIG_INVALID");
615        assert!(
616            error.to_string().contains("ALIEN_FILES_BINDING"),
617            "message should name the env var, got: {error}"
618        );
619    }
620
621    #[tokio::test]
622    async fn redis_kv_binding_returns_unsupported_binding_provider() {
623        let json = r#"{"service":"redis","connectionUrl":"redis://localhost:6379"}"#;
624        let env = with_binding(base_env(), "cache", json);
625        let bindings = Bindings::from_env_map(env).expect("valid JSON should construct");
626
627        let error = bindings
628            .kv("cache")
629            .await
630            .expect_err("redis is not a supported kv provider in this build");
631
632        assert_eq!(error.code, "UNSUPPORTED_BINDING_PROVIDER");
633        assert!(
634            error.to_string().contains("ALIEN_CACHE_BINDING"),
635            "message should name the env var, got: {error}"
636        );
637    }
638}