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 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
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 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 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 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 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 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 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}