1use crate::error::Result;
10use crate::provider::{BindingsProvider, LazyEnvBindingsProvider};
11use crate::refreshing::{
12 RefreshingKey, RefreshingKv, RefreshingQueue, RefreshingStorage, RefreshingVault,
13 RefreshingWorker,
14};
15use crate::traits::{
16 BindingsProviderApi, Container, Key, Kv, MessagePayload, Postgres, Queue, QueueMessage,
17 Sandbox, Storage, Vault, Worker,
18};
19use std::collections::HashMap;
20use std::sync::Arc;
21
22#[derive(Debug)]
51pub struct Bindings {
52 provider: Arc<LazyEnvBindingsProvider>,
53}
54
55#[derive(Clone)]
57pub struct BoundQueue {
58 inner: Arc<dyn Queue>,
59 name: Arc<str>,
60}
61
62impl std::fmt::Debug for BoundQueue {
63 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
64 formatter
65 .debug_struct("Queue")
66 .field("name", &self.name)
67 .finish_non_exhaustive()
68 }
69}
70
71impl BoundQueue {
72 fn new(inner: Arc<dyn Queue>, name: impl Into<Arc<str>>) -> Self {
73 Self {
74 inner,
75 name: name.into(),
76 }
77 }
78
79 pub async fn send(&self, message: MessagePayload) -> Result<()> {
81 self.inner.send(&self.name, message).await
82 }
83
84 pub async fn receive(&self, max_messages: usize) -> Result<Vec<QueueMessage>> {
86 self.inner.receive(&self.name, max_messages).await
87 }
88
89 pub async fn ack(&self, receipt_handle: &str) -> Result<()> {
91 self.inner.ack(&self.name, receipt_handle).await
92 }
93
94 pub async fn nack(&self, receipt_handle: &str) -> Result<()> {
96 self.inner.nack(&self.name, receipt_handle).await
97 }
98
99 pub async fn purge(&self) -> Result<()> {
101 self.inner.purge(&self.name).await
102 }
103}
104
105impl Bindings {
106 pub fn from_env() -> Result<Self> {
108 Self::from_env_map(std::env::vars().collect())
109 }
110
111 pub fn from_env_map(env: HashMap<String, String>) -> Result<Self> {
120 Ok(Self {
121 provider: Arc::new(BindingsProvider::from_env_deferred(env)?),
122 })
123 }
124
125 pub async fn storage(&self, binding_name: &str) -> Result<Arc<dyn Storage>> {
132 let initial = self.provider.load_storage(binding_name).await?;
133 Ok(Arc::new(RefreshingStorage::new(
134 self.provider.clone(),
135 binding_name.to_string(),
136 initial,
137 )))
138 }
139
140 pub async fn key(&self, binding_name: &str) -> Result<Arc<dyn Key>> {
142 self.provider.load_key(binding_name).await?;
143 Ok(Arc::new(RefreshingKey::new(
144 self.provider.clone(),
145 binding_name.to_string(),
146 )))
147 }
148
149 pub async fn kv(&self, binding_name: &str) -> Result<Arc<dyn Kv>> {
152 self.provider.load_kv(binding_name).await?;
153 Ok(Arc::new(RefreshingKv::new(
154 self.provider.clone(),
155 binding_name.to_string(),
156 )))
157 }
158
159 pub async fn queue(&self, binding_name: &str) -> Result<BoundQueue> {
161 self.provider.load_queue(binding_name).await?;
162 let queue: Arc<dyn Queue> = Arc::new(RefreshingQueue::new(
163 self.provider.clone(),
164 binding_name.to_string(),
165 ));
166 Ok(BoundQueue::new(queue, binding_name))
167 }
168
169 pub async fn vault(&self, binding_name: &str) -> Result<Arc<dyn Vault>> {
172 self.provider.load_vault(binding_name).await?;
173 Ok(Arc::new(RefreshingVault::new(
174 self.provider.clone(),
175 binding_name.to_string(),
176 )))
177 }
178
179 pub async fn container(&self, binding_name: &str) -> Result<Arc<dyn Container>> {
181 self.provider.load_container(binding_name).await
182 }
183
184 pub async fn worker(&self, binding_name: &str) -> Result<Arc<dyn Worker>> {
186 self.provider.load_worker(binding_name).await?;
187 Ok(Arc::new(RefreshingWorker::new(
188 self.provider.clone(),
189 binding_name.to_string(),
190 )))
191 }
192
193 pub async fn postgres(&self, binding_name: &str) -> Result<Arc<dyn Postgres>> {
209 self.provider.load_postgres(binding_name).await
210 }
211
212 pub async fn sandbox(&self, binding_name: &str) -> Result<Arc<dyn Sandbox>> {
217 self.provider.load_sandbox(binding_name).await
218 }
219}
220
221#[cfg(test)]
222mod tests {
223 use super::*;
224 use std::net::SocketAddr;
225 use std::sync::atomic::{AtomicUsize, Ordering};
226
227 use crate::error::binding_env_var;
228 use crate::traits::{MessagePayload, WorkerInvokeRequest};
229 use alien_core::{
230 Platform, ENV_ALIEN_DEPLOYMENT_ID, ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT,
231 ENV_ALIEN_DEPLOYMENT_TOKEN, ENV_ALIEN_DEPLOYMENT_TYPE, ENV_ALIEN_MANAGER_URL,
232 ENV_ALIEN_RESOURCE_ID,
233 };
234 use axum::{extract::State, routing::post, Json, Router};
235 use object_store::{path::Path as ObjectPath, PutPayload};
236 use std::collections::HashMap;
237 use tempfile::TempDir;
238
239 fn base_env() -> HashMap<String, String> {
241 HashMap::from([(
242 ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
243 Platform::Local.as_str().to_string(),
244 )])
245 }
246
247 fn with_binding(
248 mut env: HashMap<String, String>,
249 binding_name: &str,
250 json: &str,
251 ) -> HashMap<String, String> {
252 env.insert(binding_env_var(binding_name), json.to_string());
253 env
254 }
255
256 #[derive(Clone)]
257 struct MintServerState {
258 calls: Arc<AtomicUsize>,
259 state_directory: String,
260 }
261
262 async fn mint_handler(State(state): State<MintServerState>) -> Json<serde_json::Value> {
263 let call = state.calls.fetch_add(1, Ordering::SeqCst) + 1;
264 let lifetime_seconds = if call == 1 { 120 } else { 3600 };
265 let expires_at =
266 (chrono::Utc::now() + chrono::Duration::seconds(lifetime_seconds)).to_rfc3339();
267 Json(serde_json::json!({
268 "clientConfig": {
269 "platform": "local",
270 "state_directory": state.state_directory,
271 },
272 "expiresAt": expires_at,
273 "principal": "local:refreshing-binding-test",
274 }))
275 }
276
277 async fn spawn_mint_server(state_directory: &str) -> (String, Arc<AtomicUsize>) {
278 let calls = Arc::new(AtomicUsize::new(0));
279 let app = Router::new()
280 .route("/v1/credentials/mint", post(mint_handler))
281 .with_state(MintServerState {
282 calls: calls.clone(),
283 state_directory: state_directory.to_string(),
284 });
285 let listener = tokio::net::TcpListener::bind(SocketAddr::from(([127, 0, 0, 1], 0)))
286 .await
287 .expect("bind fake mint server");
288 let address = listener.local_addr().expect("read fake server address");
289 tokio::spawn(async move {
290 axum::serve(listener, app)
291 .await
292 .expect("serve fake mint endpoint");
293 });
294 (format!("http://{address}"), calls)
295 }
296
297 fn mint_env(manager_url: &str) -> HashMap<String, String> {
298 HashMap::from([
299 (
300 ENV_ALIEN_DEPLOYMENT_TYPE.to_string(),
301 Platform::Aws.as_str().to_string(),
302 ),
303 ("AWS_EC2_METADATA_DISABLED".to_string(), "true".to_string()),
304 (
305 "AWS_PROFILE".to_string(),
306 "__alien_missing_refresh_test_profile__".to_string(),
307 ),
308 (ENV_ALIEN_MANAGER_URL.to_string(), manager_url.to_string()),
309 (
310 ENV_ALIEN_DEPLOYMENT_TOKEN.to_string(),
311 "refresh-test-token".to_string(),
312 ),
313 (
314 ENV_ALIEN_DEPLOYMENT_ID.to_string(),
315 "refresh-test-deployment".to_string(),
316 ),
317 (
318 ENV_ALIEN_DEPLOYMENT_SERVICE_ACCOUNT.to_string(),
319 "refresh-test-service-account".to_string(),
320 ),
321 (
322 ENV_ALIEN_RESOURCE_ID.to_string(),
323 "refresh-test-resource".to_string(),
324 ),
325 ])
326 }
327
328 #[test]
329 fn from_env_map_constructs_synchronously_from_injected_env() {
330 let bindings =
333 Bindings::from_env_map(base_env()).expect("valid env should construct Bindings");
334 drop(bindings);
335 }
336
337 #[tokio::test]
338 async fn storage_delegates_to_local_provider_and_performs_real_io() {
339 let temp_dir = TempDir::new().expect("tempdir");
340 let json = format!(
341 r#"{{"service":"local-storage","storagePath":"{}"}}"#,
342 temp_dir.path().display()
343 );
344 let env = with_binding(base_env(), "files", &json);
345 let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
346
347 let storage = bindings
348 .storage("files")
349 .await
350 .expect("storage binding should load");
351
352 let path = ObjectPath::from("greeting.txt");
353 storage
354 .put(&path, PutPayload::from(bytes::Bytes::from_static(b"hello")))
355 .await
356 .expect("put should succeed");
357 let fetched = storage
358 .get(&path)
359 .await
360 .expect("get should succeed")
361 .bytes()
362 .await
363 .expect("reading bytes should succeed");
364 assert_eq!(fetched.as_ref(), b"hello");
365 }
366
367 #[tokio::test]
368 async fn kv_delegates_to_local_provider_and_performs_real_io() {
369 let temp_dir = TempDir::new().expect("tempdir");
370 let json = format!(
371 r#"{{"service":"local-kv","dataDir":"{}"}}"#,
372 temp_dir.path().display()
373 );
374 let env = with_binding(base_env(), "cache", &json);
375 let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
376
377 let kv = bindings.kv("cache").await.expect("kv binding should load");
378
379 kv.put("greeting", b"hi".to_vec(), None)
380 .await
381 .expect("put should succeed");
382 let value = kv
383 .get("greeting")
384 .await
385 .expect("get should succeed")
386 .expect("value should exist");
387 assert_eq!(value.value, b"hi");
388 }
389
390 #[tokio::test]
391 async fn long_lived_kv_handle_refreshes_minted_provider_before_expiry() {
392 let temp_dir = TempDir::new().expect("tempdir");
393 let (manager_url, calls) = spawn_mint_server(
394 temp_dir
395 .path()
396 .to_str()
397 .expect("tempdir path must be valid UTF-8"),
398 )
399 .await;
400 let json = format!(
401 r#"{{"service":"local-kv","dataDir":"{}"}}"#,
402 temp_dir.path().display()
403 );
404 let env = with_binding(mint_env(&manager_url), "cache", &json);
405 let bindings = Bindings::from_env_map(env).expect("minting env should construct Bindings");
406
407 let kv = bindings
408 .kv("cache")
409 .await
410 .expect("first binding resolution should mint credentials");
411 assert_eq!(calls.load(Ordering::SeqCst), 1);
412
413 kv.put("greeting", b"hi".to_vec(), None)
414 .await
415 .expect("the long-lived handle should refresh and write");
416 assert_eq!(
417 calls.load(Ordering::SeqCst),
418 2,
419 "the first mint is still unexpired but inside the refresh window"
420 );
421
422 let value = kv
423 .get("greeting")
424 .await
425 .expect("the same long-lived handle should read")
426 .expect("value should exist");
427 assert_eq!(value.value, b"hi");
428 assert_eq!(
429 calls.load(Ordering::SeqCst),
430 2,
431 "the refreshed provider should stay cached while fresh"
432 );
433 }
434
435 #[tokio::test]
436 async fn long_lived_worker_handle_refreshes_minted_provider_before_invocation() {
437 let temp_dir = TempDir::new().expect("tempdir");
438 let (manager_url, calls) = spawn_mint_server(
439 temp_dir
440 .path()
441 .to_str()
442 .expect("tempdir path must be valid UTF-8"),
443 )
444 .await;
445 let app = Router::new().route("/processor/jobs", post(|| async { "accepted" }));
446 let listener = tokio::net::TcpListener::bind(SocketAddr::from(([127, 0, 0, 1], 0)))
447 .await
448 .expect("bind fake worker server");
449 let address = listener.local_addr().expect("read fake worker address");
450 tokio::spawn(async move {
451 axum::serve(listener, app)
452 .await
453 .expect("serve fake worker endpoint");
454 });
455
456 let worker_url = format!("http://{address}");
457 let json = format!(r#"{{"service":"local","workerUrl":"{worker_url}"}}"#);
458 let env = with_binding(mint_env(&manager_url), "processor", &json);
459 let bindings = Bindings::from_env_map(env).expect("minting env should construct Bindings");
460 let worker = bindings
461 .worker("processor")
462 .await
463 .expect("first binding resolution should mint credentials");
464 assert_eq!(calls.load(Ordering::SeqCst), 1);
465
466 let response = worker
467 .invoke(WorkerInvokeRequest {
468 target_worker: "processor".to_string(),
469 method: "POST".to_string(),
470 path: "/jobs".to_string(),
471 headers: Default::default(),
472 body: Vec::new(),
473 timeout: None,
474 })
475 .await
476 .expect("the long-lived handle should refresh and invoke");
477 assert_eq!(response.status, 200);
478 assert_eq!(response.body, b"accepted");
479 assert_eq!(calls.load(Ordering::SeqCst), 2);
480
481 assert_eq!(
482 worker
483 .get_worker_url()
484 .await
485 .expect("the same long-lived handle should resolve its URL"),
486 Some(worker_url)
487 );
488 assert_eq!(
489 calls.load(Ordering::SeqCst),
490 2,
491 "the refreshed provider should stay cached while fresh"
492 );
493 }
494
495 #[tokio::test]
496 async fn queue_delegates_to_local_provider_and_performs_real_io() {
497 let temp_dir = TempDir::new().expect("tempdir");
498 let json = format!(
499 r#"{{"service":"local-queue","queuePath":"{}"}}"#,
500 temp_dir.path().join("queue.db").display()
501 );
502 let env = with_binding(base_env(), "jobs", &json);
503 let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
504
505 let queue = bindings
506 .queue("jobs")
507 .await
508 .expect("queue binding should load");
509
510 queue
511 .send(MessagePayload::Text("hello".to_string()))
512 .await
513 .expect("send should succeed");
514 let messages = queue.receive(1).await.expect("receive should succeed");
515 assert_eq!(messages.len(), 1);
516 }
517
518 #[tokio::test]
519 async fn bound_queue_uses_its_configured_name_for_every_operation() {
520 let temp_dir = TempDir::new().expect("tempdir");
521 let json = format!(
522 r#"{{"service":"local-queue","queuePath":"{}"}}"#,
523 temp_dir.path().join("queue.db").display()
524 );
525 let env = with_binding(base_env(), "jobs", &json);
526 let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
527
528 let queue = bindings
529 .queue("jobs")
530 .await
531 .expect("queue binding should load");
532
533 queue
536 .send(MessagePayload::Text("retry".to_string()))
537 .await
538 .expect("send should succeed");
539 let first = queue.receive(1).await.expect("receive should succeed");
540 assert_eq!(first.len(), 1);
541 assert!(
542 queue
543 .receive(1)
544 .await
545 .expect("receive should succeed")
546 .is_empty(),
547 "in-flight message must be hidden before nack"
548 );
549 queue
550 .nack(&first[0].receipt_handle)
551 .await
552 .expect("nack should succeed");
553 let redelivered = queue.receive(1).await.expect("receive should succeed");
554 assert_eq!(redelivered.len(), 1, "nacked message must be redelivered");
555
556 queue.purge().await.expect("purge should succeed");
558 assert!(
559 queue
560 .receive(1)
561 .await
562 .expect("receive should succeed")
563 .is_empty(),
564 "purge must empty the queue"
565 );
566 }
567
568 #[tokio::test]
569 async fn container_exposes_internal_and_optional_public_urls() {
570 let env = with_binding(
571 base_env(),
572 "database",
573 r#"{"service":"local","containerName":"database","internalUrl":"http://database.internal:5432","publicUrl":"http://localhost:15432"}"#,
574 );
575 let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
576
577 let container = bindings
578 .container("database")
579 .await
580 .expect("container binding should load");
581
582 assert_eq!(
583 container.get_internal_url(),
584 "http://database.internal:5432"
585 );
586 assert_eq!(container.get_public_url(), Some("http://localhost:15432"));
587 }
588
589 #[tokio::test]
590 async fn vault_delegates_to_local_provider_and_performs_real_io() {
591 let temp_dir = TempDir::new().expect("tempdir");
592 let json = format!(
593 r#"{{"service":"local-vault","vaultName":"secrets","dataDir":"{}"}}"#,
594 temp_dir.path().display()
595 );
596 let env = with_binding(base_env(), "secrets", &json);
597 let bindings = Bindings::from_env_map(env).expect("valid env should construct Bindings");
598
599 let vault = bindings
600 .vault("secrets")
601 .await
602 .expect("vault binding should load");
603
604 vault
605 .set_secret("api-key", "sekrit")
606 .await
607 .expect("set_secret should succeed");
608 let value = vault
609 .get_secret("api-key")
610 .await
611 .expect("get_secret should succeed");
612 assert_eq!(value, "sekrit");
613
614 vault
617 .set_secret("db-url", "postgres://…")
618 .await
619 .expect("set_secret should succeed");
620 let mut names = vault
621 .list_secrets()
622 .await
623 .expect("list_secrets should succeed");
624 names.sort();
625 assert_eq!(names, vec!["api-key".to_string(), "db-url".to_string()]);
626 }
627
628 #[tokio::test]
629 async fn missing_storage_binding_returns_binding_not_configured() {
630 let bindings = Bindings::from_env_map(base_env())
631 .expect("construction should succeed with no bindings configured");
632
633 let error = bindings
634 .storage("files")
635 .await
636 .expect_err("missing binding should error");
637
638 assert_eq!(error.code, "BINDING_NOT_CONFIGURED");
639 assert!(
640 error.to_string().contains("ALIEN_FILES_BINDING"),
641 "message should name the env var, got: {error}"
642 );
643 }
644
645 #[tokio::test]
646 async fn zero_env_construct_then_missing_binding_is_binding_not_configured() {
647 for kind in ["storage", "kv", "queue", "vault"] {
656 let bindings = Bindings::from_env_map(HashMap::new())
657 .expect("zero-env construction must succeed (platform resolution deferred)");
658
659 let error = match kind {
660 "storage" => bindings.storage("x").await.unwrap_err(),
661 "kv" => bindings.kv("x").await.unwrap_err(),
662 "queue" => bindings.queue("x").await.unwrap_err(),
663 "vault" => bindings.vault("x").await.unwrap_err(),
664 other => unreachable!("unhandled kind in table test: {other}"),
665 };
666
667 assert_eq!(
668 error.code, "BINDING_NOT_CONFIGURED",
669 "{kind}: expected the missing-binding error, not a platform/deployment error: {error}"
670 );
671 assert!(
672 error.to_string().contains("ALIEN_X_BINDING"),
673 "{kind}: message should name the env var, got: {error}"
674 );
675 }
676 }
677
678 #[test]
679 fn malformed_binding_json_returns_binding_config_invalid_naming_env_var() {
680 let env = with_binding(base_env(), "files", "not-json");
681
682 let error =
683 Bindings::from_env_map(env).expect_err("malformed binding JSON should fail to load");
684
685 assert_eq!(error.code, "BINDING_CONFIG_INVALID");
686 assert!(
687 error.to_string().contains("ALIEN_FILES_BINDING"),
688 "message should name the env var, got: {error}"
689 );
690 }
691
692 #[tokio::test]
693 async fn redis_kv_binding_returns_unsupported_binding_provider() {
694 let json = r#"{"service":"redis","connectionUrl":"redis://localhost:6379"}"#;
695 let env = with_binding(base_env(), "cache", json);
696 let bindings = Bindings::from_env_map(env).expect("valid JSON should construct");
697
698 let error = bindings
699 .kv("cache")
700 .await
701 .expect_err("redis is not a supported kv provider in this build");
702
703 assert_eq!(error.code, "UNSUPPORTED_BINDING_PROVIDER");
704 assert!(
705 error.to_string().contains("ALIEN_CACHE_BINDING"),
706 "message should name the env var, got: {error}"
707 );
708 }
709}