1use 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#[derive(Debug)]
49pub struct Bindings {
50 provider: Arc<LazyEnvBindingsProvider>,
51}
52
53#[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 pub async fn send(&self, message: MessagePayload) -> Result<()> {
79 self.inner.send(&self.name, message).await
80 }
81
82 pub async fn receive(&self, max_messages: usize) -> Result<Vec<QueueMessage>> {
84 self.inner.receive(&self.name, max_messages).await
85 }
86
87 pub async fn ack(&self, receipt_handle: &str) -> Result<()> {
89 self.inner.ack(&self.name, receipt_handle).await
90 }
91
92 pub async fn nack(&self, receipt_handle: &str) -> Result<()> {
94 self.inner.nack(&self.name, receipt_handle).await
95 }
96
97 pub async fn purge(&self) -> Result<()> {
99 self.inner.purge(&self.name).await
100 }
101}
102
103impl Bindings {
104 pub fn from_env() -> Result<Self> {
106 Self::from_env_map(std::env::vars().collect())
107 }
108
109 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 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 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 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 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 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 pub async fn container(&self, binding_name: &str) -> Result<Arc<dyn Container>> {
179 self.provider.load_container(binding_name).await
180 }
181
182 pub async fn postgres(&self, binding_name: &str) -> Result<Arc<dyn Postgres>> {
198 self.provider.load_postgres(binding_name).await
199 }
200
201 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 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 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 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 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 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 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}