1#[cfg(feature = "handlers")]
10use crate::error::Error;
11use crate::error::Result;
12use boatramp_core::deploy::DeployStore;
13use boatramp_core::envelope::KeyEnvelope;
14use boatramp_core::kv::KvStore;
15use std::path::Path;
16use std::sync::Arc;
17
18#[cfg(feature = "handlers")]
24const DEFAULT_ASYNC_TIMEOUT_MS: u64 = 15 * 60 * 1000;
25
26#[cfg(feature = "handlers")]
30const DEFAULT_ASYNC_CONCURRENCY: usize = 8;
31
32#[cfg(feature = "handlers")]
35const DEFAULT_STREAMING_TIMEOUT_MS: u64 = 15 * 60 * 1000;
36
37#[cfg(feature = "handlers")]
41const DEFAULT_STREAMING_CONCURRENCY: usize = 64;
42
43#[cfg(feature = "handlers")]
47#[allow(clippy::too_many_arguments)]
48pub async fn build_handler_runtime(
49 kv: Arc<dyn KvStore>,
50 storage: Arc<dyn boatramp_core::Storage>,
51 data_dir: &Path,
52 handlers_cfg: Option<&crate::config::HandlersConfig>,
53 messaging_override: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
54 max_blob_bytes: u64,
55 max_component_bytes: u64,
56 allow_guest_private_egress: bool,
58 self_egress_addrs: Vec<std::net::SocketAddr>,
60 allow_env_secret_refs: bool,
64 allow_guest_email: bool,
68 deploy: &DeployStore,
71 secrets_envelope: Option<Arc<dyn KeyEnvelope>>,
72) -> Result<boatramp_server::HandlerRuntime> {
73 #[cfg(not(feature = "email"))]
76 let _ = allow_guest_email;
77 let defaults = boatramp_handlers::Limits::default();
87 let sync_limits = boatramp_handlers::Limits {
88 timeout_ms: handlers_cfg
89 .and_then(|h| h.sync_max_timeout_ms)
90 .unwrap_or(defaults.timeout_ms),
91 ..defaults
92 };
93 let async_limits = boatramp_handlers::Limits {
94 timeout_ms: handlers_cfg
95 .and_then(|h| h.async_max_timeout_ms)
96 .unwrap_or(DEFAULT_ASYNC_TIMEOUT_MS),
97 max_concurrency: handlers_cfg
98 .and_then(|h| h.async_max_concurrency)
99 .unwrap_or(DEFAULT_ASYNC_CONCURRENCY),
100 fuel: handlers_cfg.and_then(|h| h.async_max_fuel),
101 ..sync_limits
102 };
103 let streaming_limits = boatramp_handlers::Limits {
108 timeout_ms: handlers_cfg
109 .and_then(|h| h.streaming_max_timeout_ms)
110 .unwrap_or(DEFAULT_STREAMING_TIMEOUT_MS),
111 max_concurrency: handlers_cfg
112 .and_then(|h| h.streaming_max_concurrency)
113 .unwrap_or(DEFAULT_STREAMING_CONCURRENCY),
114 fuel: handlers_cfg.and_then(|h| h.streaming_max_fuel),
115 ..sync_limits
116 };
117 let outbound_timeout = handlers_cfg
118 .and_then(|h| h.outbound_timeout_ms)
119 .map(std::time::Duration::from_millis);
120 let engine = if handlers_cfg.is_some_and(|h| h.pooling) {
123 boatramp_handlers::HandlerEngine::with_pooling(sync_limits, 64)?
124 } else {
125 boatramp_handlers::HandlerEngine::new(sync_limits, 64)?
126 }
127 .with_async_limits(async_limits)
128 .with_streaming_limits(streaming_limits)
129 .with_outbound_timeout(outbound_timeout)
130 .with_private_egress(allow_guest_private_egress)
131 .with_self_egress(self_egress_addrs);
132 let sql = build_sql_backends(
133 handlers_cfg.and_then(|h| h.bindings.sql.as_ref()),
134 data_dir,
135 deploy,
136 &kv,
137 secrets_envelope.as_ref(),
138 )
139 .await?;
140 let messaging: Arc<dyn boatramp_core::messaging::Messaging> = messaging_override
143 .unwrap_or_else(|| {
144 Arc::new(boatramp_core::messaging::LogMessaging::new(
145 storage.clone(),
146 kv.clone(),
147 ))
148 });
149 let kv_for_secrets = kv.clone();
152 #[cfg(feature = "email")]
156 let kv_for_email = kv.clone();
157 #[cfg(feature = "email")]
158 let messaging_for_email = messaging.clone();
159 #[cfg(feature = "email")]
160 let email_envelope = secrets_envelope.clone();
161 let runtime =
162 boatramp_server::HandlerRuntime::new(engine, kv, storage, Some(sql), Some(messaging));
163 runtime.set_max_blob_bytes(max_blob_bytes);
165 runtime.set_max_component_bytes(max_component_bytes);
166 runtime.set_allow_env_secret_refs(allow_env_secret_refs);
168 if let Some(envelope) = secrets_envelope {
172 runtime.set_secret_store(Arc::new(boatramp_core::secret_store::SecretStore::new(
173 kv_for_secrets,
174 envelope,
175 )));
176 }
177 #[cfg(feature = "email")]
185 if allow_guest_email {
186 if let Some(envelope) = email_envelope {
187 let store = Arc::new(boatramp_core::email_config::EmailProfileStore::new(
188 kv_for_email,
189 envelope,
190 ));
191 let backend = Arc::new(boatramp_handlers::LettreBackend::new(
192 allow_guest_private_egress,
193 ));
194 let spool = boatramp_server::NodeEmailSpool::spawn(
195 backend,
196 Some(messaging_for_email),
197 store.clone(),
198 );
199 runtime.set_email_profile_store(store);
200 runtime.set_email_spool(spool);
201 }
202 }
203 Ok(runtime)
204}
205
206#[cfg(feature = "handlers")]
211async fn build_sql_backends(
212 cfg: Option<&crate::config::SqlBindingConfig>,
213 data_dir: &Path,
214 deploy: &DeployStore,
215 kv: &Arc<dyn KvStore>,
216 secrets_envelope: Option<&Arc<dyn KeyEnvelope>>,
217) -> Result<Arc<dyn boatramp_core::sql::SqlBackends>> {
218 let resolve_env = |var: &Option<String>| -> Result<Option<String>> {
219 match var {
220 Some(var) => Ok(Some(
221 std::env::var(var).map_err(|_| Error::SqlEnvUnset(var.clone()))?,
222 )),
223 None => Ok(None),
224 }
225 };
226
227 let backend = match cfg.and_then(|c| c.url.as_ref()) {
228 Some(url) => {
231 let cfg = cfg.expect("url implies cfg");
232 let admin_url = cfg.admin_url.as_ref().ok_or(Error::SqlAdminUrlRequired)?;
233 let token = resolve_env(&cfg.token_env)?.unwrap_or_default();
234 let admin_token = resolve_env(&cfg.admin_token_env)?;
235 let backends = boatramp_storage::LibsqlSqlBackends::remote(
236 url.clone(),
237 admin_url.clone(),
238 token,
239 admin_token,
240 );
241 match &cfg.replica_url {
243 Some(replica_url) => backends.with_read_replica(replica_url.clone()),
244 None => backends,
245 }
246 }
247 None => {
249 let dir = cfg
250 .and_then(|c| c.dir.clone())
251 .unwrap_or_else(|| data_dir.join("handlers-sql"));
252 boatramp_storage::LibsqlSqlBackends::local(dir)
253 }
254 };
255 let preview_mode = match cfg.and_then(|c| c.preview_mode.as_deref()) {
257 None | Some("empty") => boatramp_core::sql::PreviewSqlMode::Empty,
258 Some("branch") => boatramp_core::sql::PreviewSqlMode::Branch,
259 Some("shared") => boatramp_core::sql::PreviewSqlMode::Shared,
260 Some(other) => return Err(Error::UnknownPreviewMode(other.to_string())),
261 };
262 let preview_init = match cfg.and_then(|c| c.preview_init.as_ref()) {
263 Some(path) => {
264 Some(
265 std::fs::read_to_string(path).map_err(|err| Error::PreviewInitRead {
266 path: path.clone(),
267 source: err,
268 })?,
269 )
270 }
271 None => None,
272 };
273 let default: Arc<dyn boatramp_core::sql::SqlBackends> =
274 Arc::new(backend.with_preview_policy(preview_mode, preview_init));
275
276 let databases = cfg.map(|c| &c.databases);
280 if databases.is_none_or(std::collections::BTreeMap::is_empty) {
281 return Ok(default);
282 }
283 let databases = databases.expect("checked non-empty above");
284
285 #[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
286 {
287 use boatramp_core::sql::SqlBackend;
288 use boatramp_storage::sql_sqlx::{
289 connect, CompositeSqlBackends, ExternalSqlKind, ExternalSqlOptions,
290 };
291 let timeout = |db: &crate::config::ExternalDatabaseConfig| {
292 db.connect_timeout_secs.map(std::time::Duration::from_secs)
293 };
294 let mut composite = CompositeSqlBackends::new(default);
295 for (name, db) in databases {
296 let kind = ExternalSqlKind::parse(&db.kind).ok_or_else(|| Error::SqlExternalKind {
297 name: name.clone(),
298 kind: db.kind.clone(),
299 })?;
300 if db.compute.as_deref().is_some_and(|c| !c.is_empty()) {
301 if let Some(var) = db.password_env.as_deref().filter(|v| !v.is_empty()) {
309 let workload = db.compute.as_deref().expect("compute checked above");
310 let password =
311 std::env::var(var).map_err(|_| Error::SqlEnvUnset(var.into()))?;
312 let resolver = Arc::new(crate::managed_sql::DeployEndpointResolver::new(
313 deploy.clone(),
314 boatramp_core::project::DEFAULT_PROJECT,
315 ));
316 let external: Arc<dyn SqlBackend> = Arc::new(
317 boatramp_storage::sql_compute::ComputeResolvedSqlBackend::new(
318 resolver,
319 workload,
320 kind,
321 db.database.clone().unwrap_or_default(),
322 db.user.clone().unwrap_or_default(),
323 password,
324 db.pool_max,
325 db.read_only,
326 timeout(db),
327 ),
328 );
329 composite = composite.with_external(name.clone(), external, db.allow_preview);
330 continue;
331 }
332 let envelope = secrets_envelope
335 .cloned()
336 .ok_or_else(|| Error::SqlManagedNeedsSecrets(name.clone()))?;
337 let resolver = crate::tenant_sql::NodeTenantSqlResolver::new(
338 deploy.clone(),
339 kv.clone(),
340 envelope,
341 db,
342 )
343 .expect("a compute-backed managed binding builds a per-tenant resolver");
344 let site_scoped = resolver.site_scoped();
345 composite = composite.with_per_tenant(
346 name.clone(),
347 Arc::new(resolver),
348 site_scoped,
349 db.allow_preview,
350 );
351 } else {
352 if db.url_env.trim().is_empty() {
355 return Err(Error::SqlExternalUrlEnvMissing(name.clone()));
356 }
357 let url = std::env::var(&db.url_env)
358 .map_err(|_| Error::SqlEnvUnset(db.url_env.clone()))?;
359 let read_url = match &db.read_url_env {
360 Some(var) => {
361 Some(std::env::var(var).map_err(|_| Error::SqlEnvUnset(var.clone()))?)
362 }
363 None => None,
364 };
365 let opts = ExternalSqlOptions::new(url)
366 .with_read_url(read_url)
367 .with_max_connections(db.pool_max)
368 .read_only(db.read_only)
369 .with_connect_timeout(timeout(db));
370 let external: Arc<dyn SqlBackend> =
371 connect(kind, &opts).map_err(|source| Error::SqlExternalConnect {
372 name: name.clone(),
373 source,
374 })?;
375 composite = composite.with_external(name.clone(), external, db.allow_preview);
376 }
377 }
378 Ok(Arc::new(composite))
379 }
380 #[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
381 {
382 let _ = (deploy, kv, secrets_envelope);
385 let name = databases.keys().next().cloned().unwrap_or_default();
386 Err(Error::SqlExternalUnavailable(name))
387 }
388}
389
390#[cfg(not(feature = "handlers"))]
391#[allow(clippy::too_many_arguments)]
392pub async fn build_handler_runtime(
393 _kv: Arc<dyn KvStore>,
394 _storage: Arc<dyn boatramp_core::Storage>,
395 _data_dir: &Path,
396 _handlers_cfg: Option<&crate::config::HandlersConfig>,
397 _messaging_override: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
398 _max_blob_bytes: u64,
399 _max_component_bytes: u64,
400 _allow_guest_private_egress: bool,
401 _self_egress_addrs: Vec<std::net::SocketAddr>,
402 _allow_env_secret_refs: bool,
406 _allow_guest_email: bool,
407 _deploy: &DeployStore,
408 _secrets_envelope: Option<Arc<dyn KeyEnvelope>>,
409) -> Result<boatramp_server::HandlerRuntime> {
410 Ok(boatramp_server::HandlerRuntime::disabled())
411}
412
413#[cfg(all(test, any(feature = "sql-postgres", feature = "sql-mysql")))]
414mod tests {
415 use super::*;
416 use std::result::Result;
419
420 use async_trait::async_trait;
421 use boatramp_core::envelope::EnvelopeError;
422 use boatramp_core::kv::MemoryKv;
423 use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
424
425 struct TestEnvelope;
427 #[async_trait]
428 impl KeyEnvelope for TestEnvelope {
429 async fn wrap(&self, p: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
430 Ok(p.iter().rev().copied().collect())
431 }
432 async fn unwrap(&self, w: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
433 Ok(w.iter().rev().copied().collect())
434 }
435 }
436
437 struct NullStorage;
440 #[async_trait]
441 impl Storage for NullStorage {
442 async fn get(&self, _: &str) -> Result<GetObject, StorageError> {
443 Err(StorageError::NotFound(String::new()))
444 }
445 async fn get_range(
446 &self,
447 _: &str,
448 _: u64,
449 _: Option<u64>,
450 ) -> Result<GetObject, StorageError> {
451 Err(StorageError::NotFound(String::new()))
452 }
453 async fn put(
454 &self,
455 _: &str,
456 _: ByteStream,
457 _: PutMeta,
458 ) -> Result<ObjectMeta, StorageError> {
459 Err(StorageError::unsupported("null"))
460 }
461 async fn head(&self, _: &str) -> Result<ObjectMeta, StorageError> {
462 Err(StorageError::NotFound(String::new()))
463 }
464 async fn delete(&self, _: &str) -> Result<(), StorageError> {
465 Ok(())
466 }
467 async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
468 Ok(Vec::new())
469 }
470 }
471
472 fn managed_sql_cfg() -> crate::config::SqlBindingConfig {
474 let mut databases = std::collections::BTreeMap::new();
475 databases.insert(
476 "analytics".to_string(),
477 crate::config::ExternalDatabaseConfig {
478 kind: "postgres".into(),
479 compute: Some("pg".into()),
480 database: Some("analytics".into()),
481 user: Some("app".into()),
482 ..Default::default()
483 },
484 );
485 crate::config::SqlBindingConfig {
486 databases,
487 ..Default::default()
488 }
489 }
490
491 #[tokio::test]
492 async fn managed_sql_fails_closed_without_secrets() {
493 let tmp = tempfile::tempdir().unwrap();
494 let deploy = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
495 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
496 let cfg = managed_sql_cfg();
497 match build_sql_backends(Some(&cfg), tmp.path(), &deploy, &kv, None).await {
499 Err(Error::SqlManagedNeedsSecrets(name)) => assert_eq!(name, "analytics"),
500 Ok(_) => panic!("a managed DB without [secrets] must fail closed, got Ok"),
501 Err(other) => panic!("expected SqlManagedNeedsSecrets, got: {other}"),
502 }
503 }
504
505 #[tokio::test]
506 async fn managed_sql_builds_lazily_and_seals_the_credential() {
507 let tmp = tempfile::tempdir().unwrap();
508 let deploy = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
509 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
510 let envelope: Arc<dyn KeyEnvelope> = Arc::new(TestEnvelope);
511 let cfg = managed_sql_cfg();
512 let backends = build_sql_backends(Some(&cfg), tmp.path(), &deploy, &kv, Some(&envelope))
517 .await
518 .expect("managed sql builds without a live DB (lazy connect)");
519 assert!(
520 kv.get("managed-sql-cred/default/pg")
521 .await
522 .unwrap()
523 .is_none(),
524 "nothing sealed at build — per-tenant credentials are minted on first open"
525 );
526
527 let _ = backends
531 .database("default", "blog", "analytics")
532 .await
533 .unwrap();
534 let sealed = kv
535 .get("managed-sql-cred/default/pg")
536 .await
537 .unwrap()
538 .expect("credential sealed on first resolve under the default project");
539 assert_ne!(sealed.len(), 0);
540 }
541}