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` — 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    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
202#[cfg(test)]
203mod tests {
204    use super::*;
205    use std::net::SocketAddr;
206    use std::sync::atomic::{AtomicUsize, Ordering};
207
208    use crate::error::binding_env_var;
209    use crate::traits::MessagePayload;
210    use alien_core::{
211        Platform, ENV_ALIEN_DEPLOYMENT_ID, ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT,
212        ENV_ALIEN_DEPLOYMENT_TOKEN, ENV_ALIEN_DEPLOYMENT_TYPE, ENV_ALIEN_MANAGER_URL,
213        ENV_ALIEN_RESOURCE_ID,
214    };
215    use axum::{extract::State, routing::post, Json, Router};
216    use object_store::{path::Path as ObjectPath, PutPayload};
217    use std::collections::HashMap;
218    use tempfile::TempDir;
219
220    /// Minimal valid environment (no bindings configured yet).
221    fn base_env() -> HashMap<String, String> {
222        HashMap::from([(
223            ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
224            Platform::Local.as_str().to_string(),
225        )])
226    }
227
228    fn with_binding(
229        mut env: HashMap<String, String>,
230        binding_name: &str,
231        json: &str,
232    ) -> HashMap<String, String> {
233        env.insert(binding_env_var(binding_name), json.to_string());
234        env
235    }
236
237    #[derive(Clone)]
238    struct MintServerState {
239        calls: Arc<AtomicUsize>,
240        state_directory: String,
241    }
242
243    async fn mint_handler(State(state): State<MintServerState>) -> Json<serde_json::Value> {
244        let call = state.calls.fetch_add(1, Ordering::SeqCst) + 1;
245        let lifetime_seconds = if call == 1 { 120 } else { 3600 };
246        let expires_at =
247            (chrono::Utc::now() + chrono::Duration::seconds(lifetime_seconds)).to_rfc3339();
248        Json(serde_json::json!({
249            "clientConfig": {
250                "platform": "local",
251                "state_directory": state.state_directory,
252            },
253            "expiresAt": expires_at,
254            "principal": "local:refreshing-binding-test",
255        }))
256    }
257
258    async fn spawn_mint_server(state_directory: &str) -> (String, Arc<AtomicUsize>) {
259        let calls = Arc::new(AtomicUsize::new(0));
260        let app = Router::new()
261            .route("/v1/credentials/mint", post(mint_handler))
262            .with_state(MintServerState {
263                calls: calls.clone(),
264                state_directory: state_directory.to_string(),
265            });
266        let listener = tokio::net::TcpListener::bind(SocketAddr::from(([127, 0, 0, 1], 0)))
267            .await
268            .expect("bind fake mint server");
269        let address = listener.local_addr().expect("read fake server address");
270        tokio::spawn(async move {
271            axum::serve(listener, app)
272                .await
273                .expect("serve fake mint endpoint");
274        });
275        (format!("http://{address}"), calls)
276    }
277
278    fn mint_env(manager_url: &str) -> HashMap<String, String> {
279        HashMap::from([
280            (
281                ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
282                Platform::Aws.as_str().to_string(),
283            ),
284            ("AWS_EC2_METADATA_DISABLED".to_string(), "true".to_string()),
285            (
286                "AWS_PROFILE".to_string(),
287                "__alien_missing_refresh_test_profile__".to_string(),
288            ),
289            (ENV_ALIEN_MANAGER_URL.to_string(), manager_url.to_string()),
290            (
291                ENV_ALIEN_DEPLOYMENT_TOKEN.to_string(),
292                "refresh-test-token".to_string(),
293            ),
294            (
295                ENV_ALIEN_DEPLOYMENT_ID.to_string(),
296                "refresh-test-deployment".to_string(),
297            ),
298            (
299                ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT.to_string(),
300                "refresh-test-service-account".to_string(),
301            ),
302            (
303                ENV_ALIEN_RESOURCE_ID.to_string(),
304                "refresh-test-resource".to_string(),
305            ),
306        ])
307    }
308
309    #[test]
310    fn from_env_map_constructs_synchronously_from_injected_env() {
311        // No `.await` here at all: proves construction is a plain sync function,
312        // not something that merely returns a Future.
313        let bindings =
314            Bindings::from_env_map(base_env()).expect("valid env should construct Bindings");
315        drop(bindings);
316    }
317
318    #[tokio::test]
319    async fn storage_delegates_to_local_provider_and_performs_real_io() {
320        let temp_dir = TempDir::new().expect("tempdir");
321        let json = format!(
322            r#"{{"service":"local-storage","storagePath":"{}"}}"#,
323            temp_dir.path().display()
324        );
325        let env = with_binding(base_env(), "files", &json);
326        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
327
328        let storage = bindings
329            .storage("files")
330            .await
331            .expect("storage binding should load");
332
333        let path = ObjectPath::from("greeting.txt");
334        storage
335            .put(&path, PutPayload::from(bytes::Bytes::from_static(b"hello")))
336            .await
337            .expect("put should succeed");
338        let fetched = storage
339            .get(&path)
340            .await
341            .expect("get should succeed")
342            .bytes()
343            .await
344            .expect("reading bytes should succeed");
345        assert_eq!(fetched.as_ref(), b"hello");
346    }
347
348    #[tokio::test]
349    async fn kv_delegates_to_local_provider_and_performs_real_io() {
350        let temp_dir = TempDir::new().expect("tempdir");
351        let json = format!(
352            r#"{{"service":"local-kv","dataDir":"{}"}}"#,
353            temp_dir.path().display()
354        );
355        let env = with_binding(base_env(), "cache", &json);
356        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
357
358        let kv = bindings.kv("cache").await.expect("kv binding should load");
359
360        kv.put("greeting", b"hi".to_vec(), None)
361            .await
362            .expect("put should succeed");
363        let value = kv
364            .get("greeting")
365            .await
366            .expect("get should succeed")
367            .expect("value should exist");
368        assert_eq!(value.value, b"hi");
369    }
370
371    #[tokio::test]
372    async fn long_lived_kv_handle_refreshes_minted_provider_before_expiry() {
373        let temp_dir = TempDir::new().expect("tempdir");
374        let (manager_url, calls) = spawn_mint_server(
375            temp_dir
376                .path()
377                .to_str()
378                .expect("tempdir path must be valid UTF-8"),
379        )
380        .await;
381        let json = format!(
382            r#"{{"service":"local-kv","dataDir":"{}"}}"#,
383            temp_dir.path().display()
384        );
385        let env = with_binding(mint_env(&manager_url), "cache", &json);
386        let bindings = Bindings::from_env_map(env).expect("minting env should construct Bindings");
387
388        let kv = bindings
389            .kv("cache")
390            .await
391            .expect("first binding resolution should mint credentials");
392        assert_eq!(calls.load(Ordering::SeqCst), 1);
393
394        kv.put("greeting", b"hi".to_vec(), None)
395            .await
396            .expect("the long-lived handle should refresh and write");
397        assert_eq!(
398            calls.load(Ordering::SeqCst),
399            2,
400            "the first mint is still unexpired but inside the refresh window"
401        );
402
403        let value = kv
404            .get("greeting")
405            .await
406            .expect("the same long-lived handle should read")
407            .expect("value should exist");
408        assert_eq!(value.value, b"hi");
409        assert_eq!(
410            calls.load(Ordering::SeqCst),
411            2,
412            "the refreshed provider should stay cached while fresh"
413        );
414    }
415
416    #[tokio::test]
417    async fn queue_delegates_to_local_provider_and_performs_real_io() {
418        let temp_dir = TempDir::new().expect("tempdir");
419        let json = format!(
420            r#"{{"service":"local-queue","queuePath":"{}"}}"#,
421            temp_dir.path().join("queue.db").display()
422        );
423        let env = with_binding(base_env(), "jobs", &json);
424        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
425
426        let queue = bindings
427            .queue("jobs")
428            .await
429            .expect("queue binding should load");
430
431        queue
432            .send(MessagePayload::Text("hello".to_string()))
433            .await
434            .expect("send should succeed");
435        let messages = queue.receive(1).await.expect("receive should succeed");
436        assert_eq!(messages.len(), 1);
437    }
438
439    #[tokio::test]
440    async fn bound_queue_uses_its_configured_name_for_every_operation() {
441        let temp_dir = TempDir::new().expect("tempdir");
442        let json = format!(
443            r#"{{"service":"local-queue","queuePath":"{}"}}"#,
444            temp_dir.path().join("queue.db").display()
445        );
446        let env = with_binding(base_env(), "jobs", &json);
447        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
448
449        let queue = bindings
450            .queue("jobs")
451            .await
452            .expect("queue binding should load");
453
454        // nack: an in-flight message under the default lease is hidden, but a
455        // nack makes it immediately redeliverable.
456        queue
457            .send(MessagePayload::Text("retry".to_string()))
458            .await
459            .expect("send should succeed");
460        let first = queue.receive(1).await.expect("receive should succeed");
461        assert_eq!(first.len(), 1);
462        assert!(
463            queue
464                .receive(1)
465                .await
466                .expect("receive should succeed")
467                .is_empty(),
468            "in-flight message must be hidden before nack"
469        );
470        queue
471            .nack(&first[0].receipt_handle)
472            .await
473            .expect("nack should succeed");
474        let redelivered = queue.receive(1).await.expect("receive should succeed");
475        assert_eq!(redelivered.len(), 1, "nacked message must be redelivered");
476
477        // purge: clears everything, in flight or visible.
478        queue.purge().await.expect("purge should succeed");
479        assert!(
480            queue
481                .receive(1)
482                .await
483                .expect("receive should succeed")
484                .is_empty(),
485            "purge must empty the queue"
486        );
487    }
488
489    #[tokio::test]
490    async fn container_exposes_internal_and_optional_public_urls() {
491        let env = with_binding(
492            base_env(),
493            "database",
494            r#"{"service":"local","containerName":"database","internalUrl":"http://database.internal:5432","publicUrl":"http://localhost:15432"}"#,
495        );
496        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
497
498        let container = bindings
499            .container("database")
500            .await
501            .expect("container binding should load");
502
503        assert_eq!(
504            container.get_internal_url(),
505            "http://database.internal:5432"
506        );
507        assert_eq!(container.get_public_url(), Some("http://localhost:15432"));
508    }
509
510    #[tokio::test]
511    async fn vault_delegates_to_local_provider_and_performs_real_io() {
512        let temp_dir = TempDir::new().expect("tempdir");
513        let json = format!(
514            r#"{{"service":"local-vault","vaultName":"secrets","dataDir":"{}"}}"#,
515            temp_dir.path().display()
516        );
517        let env = with_binding(base_env(), "secrets", &json);
518        let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
519
520        let vault = bindings
521            .vault("secrets")
522            .await
523            .expect("vault binding should load");
524
525        vault
526            .set_secret("api-key", "sekrit")
527            .await
528            .expect("set_secret should succeed");
529        let value = vault
530            .get_secret("api-key")
531            .await
532            .expect("get_secret should succeed");
533        assert_eq!(value, "sekrit");
534
535        // list_secrets must be reachable through the `Arc<dyn Vault>` surface
536        // and return the stored names.
537        vault
538            .set_secret("db-url", "postgres://…")
539            .await
540            .expect("set_secret should succeed");
541        let mut names = vault
542            .list_secrets()
543            .await
544            .expect("list_secrets should succeed");
545        names.sort();
546        assert_eq!(names, vec!["api-key".to_string(), "db-url".to_string()]);
547    }
548
549    #[tokio::test]
550    async fn missing_storage_binding_returns_binding_not_configured() {
551        let bindings = Bindings::from_env_map(base_env())
552            .expect("construction should succeed with no bindings configured");
553
554        let error = bindings
555            .storage("files")
556            .await
557            .expect_err("missing binding should error");
558
559        assert_eq!(error.code, "BINDING_NOT_CONFIGURED");
560        assert!(
561            error.to_string().contains("ALIEN_FILES_BINDING"),
562            "message should name the env var, got: {error}"
563        );
564    }
565
566    #[tokio::test]
567    async fn zero_env_construct_then_missing_binding_is_binding_not_configured() {
568        // The app-facing contract: with NO deployment type and NO credentials,
569        // construction must succeed and the FIRST op on a missing binding must
570        // report BINDING_NOT_CONFIGURED (naming ALIEN_<NAME>_BINDING) BEFORE any
571        // platform / client-config resolution. There is deliberately no
572        // ALIEN_DEPLOYMENT_TYPE in this environment. Table test over all four
573        // app-facing kinds so a future kind added to `Bindings` without wiring
574        // `ensure_binding_present` into its `load_*` method fails this test
575        // instead of silently regressing to ENVIRONMENT_VARIABLE_MISSING.
576        for kind in ["storage", "kv", "queue", "vault"] {
577            let bindings = Bindings::from_env_map(HashMap::new())
578                .expect("zero-env construction must succeed (platform resolution deferred)");
579
580            let error = match kind {
581                "storage" => bindings.storage("x").await.unwrap_err(),
582                "kv" => bindings.kv("x").await.unwrap_err(),
583                "queue" => bindings.queue("x").await.unwrap_err(),
584                "vault" => bindings.vault("x").await.unwrap_err(),
585                other => unreachable!("unhandled kind in table test: {other}"),
586            };
587
588            assert_eq!(
589                error.code, "BINDING_NOT_CONFIGURED",
590                "{kind}: expected the missing-binding error, not a platform/deployment error: {error}"
591            );
592            assert!(
593                error.to_string().contains("ALIEN_X_BINDING"),
594                "{kind}: message should name the env var, got: {error}"
595            );
596        }
597    }
598
599    #[test]
600    fn malformed_binding_json_returns_binding_config_invalid_naming_env_var() {
601        let env = with_binding(base_env(), "files", "not-json");
602
603        let error =
604            Bindings::from_env_map(env).expect_err("malformed binding JSON should fail to load");
605
606        assert_eq!(error.code, "BINDING_CONFIG_INVALID");
607        assert!(
608            error.to_string().contains("ALIEN_FILES_BINDING"),
609            "message should name the env var, got: {error}"
610        );
611    }
612
613    #[tokio::test]
614    async fn redis_kv_binding_returns_unsupported_binding_provider() {
615        let json = r#"{"service":"redis","connectionUrl":"redis://localhost:6379"}"#;
616        let env = with_binding(base_env(), "cache", json);
617        let bindings = Bindings::from_env_map(env).expect("valid JSON should construct");
618
619        let error = bindings
620            .kv("cache")
621            .await
622            .expect_err("redis is not a supported kv provider in this build");
623
624        assert_eq!(error.code, "UNSUPPORTED_BINDING_PROVIDER");
625        assert!(
626            error.to_string().contains("ALIEN_CACHE_BINDING"),
627            "message should name the env var, got: {error}"
628        );
629    }
630}