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 deploy: &DeployStore,
63 secrets_envelope: Option<Arc<dyn KeyEnvelope>>,
64) -> Result<boatramp_server::HandlerRuntime> {
65 let defaults = boatramp_handlers::Limits::default();
75 let sync_limits = boatramp_handlers::Limits {
76 timeout_ms: handlers_cfg
77 .and_then(|h| h.sync_max_timeout_ms)
78 .unwrap_or(defaults.timeout_ms),
79 ..defaults
80 };
81 let async_limits = boatramp_handlers::Limits {
82 timeout_ms: handlers_cfg
83 .and_then(|h| h.async_max_timeout_ms)
84 .unwrap_or(DEFAULT_ASYNC_TIMEOUT_MS),
85 max_concurrency: handlers_cfg
86 .and_then(|h| h.async_max_concurrency)
87 .unwrap_or(DEFAULT_ASYNC_CONCURRENCY),
88 fuel: handlers_cfg.and_then(|h| h.async_max_fuel),
89 ..sync_limits
90 };
91 let streaming_limits = boatramp_handlers::Limits {
96 timeout_ms: handlers_cfg
97 .and_then(|h| h.streaming_max_timeout_ms)
98 .unwrap_or(DEFAULT_STREAMING_TIMEOUT_MS),
99 max_concurrency: handlers_cfg
100 .and_then(|h| h.streaming_max_concurrency)
101 .unwrap_or(DEFAULT_STREAMING_CONCURRENCY),
102 fuel: handlers_cfg.and_then(|h| h.streaming_max_fuel),
103 ..sync_limits
104 };
105 let outbound_timeout = handlers_cfg
106 .and_then(|h| h.outbound_timeout_ms)
107 .map(std::time::Duration::from_millis);
108 let engine = if handlers_cfg.is_some_and(|h| h.pooling) {
111 boatramp_handlers::HandlerEngine::with_pooling(sync_limits, 64)?
112 } else {
113 boatramp_handlers::HandlerEngine::new(sync_limits, 64)?
114 }
115 .with_async_limits(async_limits)
116 .with_streaming_limits(streaming_limits)
117 .with_outbound_timeout(outbound_timeout)
118 .with_private_egress(allow_guest_private_egress)
119 .with_self_egress(self_egress_addrs);
120 let sql = build_sql_backends(
121 handlers_cfg.and_then(|h| h.bindings.sql.as_ref()),
122 data_dir,
123 deploy,
124 &kv,
125 secrets_envelope.as_ref(),
126 )
127 .await?;
128 let messaging: Arc<dyn boatramp_core::messaging::Messaging> = messaging_override
131 .unwrap_or_else(|| {
132 Arc::new(boatramp_core::messaging::LogMessaging::new(
133 storage.clone(),
134 kv.clone(),
135 ))
136 });
137 let runtime =
138 boatramp_server::HandlerRuntime::new(engine, kv, storage, Some(sql), Some(messaging));
139 runtime.set_max_blob_bytes(max_blob_bytes);
141 runtime.set_max_component_bytes(max_component_bytes);
142 Ok(runtime)
143}
144
145#[cfg(feature = "handlers")]
150async fn build_sql_backends(
151 cfg: Option<&crate::config::SqlBindingConfig>,
152 data_dir: &Path,
153 deploy: &DeployStore,
154 kv: &Arc<dyn KvStore>,
155 secrets_envelope: Option<&Arc<dyn KeyEnvelope>>,
156) -> Result<Arc<dyn boatramp_core::sql::SqlBackends>> {
157 let resolve_env = |var: &Option<String>| -> Result<Option<String>> {
158 match var {
159 Some(var) => Ok(Some(
160 std::env::var(var).map_err(|_| Error::SqlEnvUnset(var.clone()))?,
161 )),
162 None => Ok(None),
163 }
164 };
165
166 let backend = match cfg.and_then(|c| c.url.as_ref()) {
167 Some(url) => {
170 let cfg = cfg.expect("url implies cfg");
171 let admin_url = cfg.admin_url.as_ref().ok_or(Error::SqlAdminUrlRequired)?;
172 let token = resolve_env(&cfg.token_env)?.unwrap_or_default();
173 let admin_token = resolve_env(&cfg.admin_token_env)?;
174 let backends = boatramp_storage::LibsqlSqlBackends::remote(
175 url.clone(),
176 admin_url.clone(),
177 token,
178 admin_token,
179 );
180 match &cfg.replica_url {
182 Some(replica_url) => backends.with_read_replica(replica_url.clone()),
183 None => backends,
184 }
185 }
186 None => {
188 let dir = cfg
189 .and_then(|c| c.dir.clone())
190 .unwrap_or_else(|| data_dir.join("handlers-sql"));
191 boatramp_storage::LibsqlSqlBackends::local(dir)
192 }
193 };
194 let preview_mode = match cfg.and_then(|c| c.preview_mode.as_deref()) {
196 None | Some("empty") => boatramp_core::sql::PreviewSqlMode::Empty,
197 Some("branch") => boatramp_core::sql::PreviewSqlMode::Branch,
198 Some("shared") => boatramp_core::sql::PreviewSqlMode::Shared,
199 Some(other) => return Err(Error::UnknownPreviewMode(other.to_string())),
200 };
201 let preview_init = match cfg.and_then(|c| c.preview_init.as_ref()) {
202 Some(path) => {
203 Some(
204 std::fs::read_to_string(path).map_err(|err| Error::PreviewInitRead {
205 path: path.clone(),
206 source: err,
207 })?,
208 )
209 }
210 None => None,
211 };
212 let default: Arc<dyn boatramp_core::sql::SqlBackends> =
213 Arc::new(backend.with_preview_policy(preview_mode, preview_init));
214
215 let databases = cfg.map(|c| &c.databases);
219 if databases.is_none_or(std::collections::BTreeMap::is_empty) {
220 return Ok(default);
221 }
222 let databases = databases.expect("checked non-empty above");
223
224 #[cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
225 {
226 use boatramp_core::project::DEFAULT_PROJECT;
227 use boatramp_core::sql::SqlBackend;
228 use boatramp_storage::sql_compute::ComputeResolvedSqlBackend;
229 use boatramp_storage::sql_sqlx::{
230 connect, CompositeSqlBackends, ExternalSqlKind, ExternalSqlOptions,
231 };
232 let timeout = |db: &crate::config::ExternalDatabaseConfig| {
233 db.connect_timeout_secs.map(std::time::Duration::from_secs)
234 };
235 let mut composite = CompositeSqlBackends::new(default);
236 for (name, db) in databases {
237 let kind = ExternalSqlKind::parse(&db.kind).ok_or_else(|| Error::SqlExternalKind {
238 name: name.clone(),
239 kind: db.kind.clone(),
240 })?;
241 let external: Arc<dyn SqlBackend> = if let Some(workload) =
242 db.compute.as_deref().filter(|c| !c.is_empty())
243 {
244 let password = match db.password_env.as_deref().filter(|v| !v.is_empty()) {
248 Some(var) => std::env::var(var).map_err(|_| Error::SqlEnvUnset(var.into()))?,
249 None => {
250 let envelope = secrets_envelope
253 .cloned()
254 .ok_or_else(|| Error::SqlManagedNeedsSecrets(name.clone()))?;
255 crate::managed_sql::ManagedSqlCredentials::new(kv.clone(), envelope)
256 .password(DEFAULT_PROJECT, workload)
257 .await
258 .map_err(|reason| Error::SqlManagedCredential {
259 name: name.clone(),
260 reason,
261 })?
262 }
263 };
264 let resolver = Arc::new(crate::managed_sql::DeployEndpointResolver::new(
265 deploy.clone(),
266 DEFAULT_PROJECT,
267 ));
268 Arc::new(ComputeResolvedSqlBackend::new(
269 resolver,
270 workload,
271 kind,
272 db.database.clone().unwrap_or_default(),
273 db.user.clone().unwrap_or_default(),
274 password,
275 db.pool_max,
276 db.read_only,
277 timeout(db),
278 ))
279 } else {
280 if db.url_env.trim().is_empty() {
283 return Err(Error::SqlExternalUrlEnvMissing(name.clone()));
284 }
285 let url = std::env::var(&db.url_env)
286 .map_err(|_| Error::SqlEnvUnset(db.url_env.clone()))?;
287 let read_url = match &db.read_url_env {
288 Some(var) => {
289 Some(std::env::var(var).map_err(|_| Error::SqlEnvUnset(var.clone()))?)
290 }
291 None => None,
292 };
293 let opts = ExternalSqlOptions::new(url)
294 .with_read_url(read_url)
295 .with_max_connections(db.pool_max)
296 .read_only(db.read_only)
297 .with_connect_timeout(timeout(db));
298 connect(kind, &opts).map_err(|source| Error::SqlExternalConnect {
299 name: name.clone(),
300 source,
301 })?
302 };
303 composite = composite.with_external(name.clone(), external, db.allow_preview);
304 }
305 Ok(Arc::new(composite))
306 }
307 #[cfg(not(any(feature = "sql-postgres", feature = "sql-mysql")))]
308 {
309 let _ = (deploy, kv, secrets_envelope);
312 let name = databases.keys().next().cloned().unwrap_or_default();
313 Err(Error::SqlExternalUnavailable(name))
314 }
315}
316
317#[cfg(not(feature = "handlers"))]
318#[allow(clippy::too_many_arguments)]
319pub async fn build_handler_runtime(
320 _kv: Arc<dyn KvStore>,
321 _storage: Arc<dyn boatramp_core::Storage>,
322 _data_dir: &Path,
323 _handlers_cfg: Option<&crate::config::HandlersConfig>,
324 _messaging_override: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
325 _max_blob_bytes: u64,
326 _max_component_bytes: u64,
327 _allow_guest_private_egress: bool,
328 _self_egress_addrs: Vec<std::net::SocketAddr>,
329 _deploy: &DeployStore,
330 _secrets_envelope: Option<Arc<dyn KeyEnvelope>>,
331) -> Result<boatramp_server::HandlerRuntime> {
332 Ok(boatramp_server::HandlerRuntime::disabled())
333}
334
335#[cfg(all(test, any(feature = "sql-postgres", feature = "sql-mysql")))]
336mod tests {
337 use super::*;
338 use std::result::Result;
341
342 use async_trait::async_trait;
343 use boatramp_core::envelope::EnvelopeError;
344 use boatramp_core::kv::MemoryKv;
345 use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
346
347 struct TestEnvelope;
349 #[async_trait]
350 impl KeyEnvelope for TestEnvelope {
351 async fn wrap(&self, p: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
352 Ok(p.iter().rev().copied().collect())
353 }
354 async fn unwrap(&self, w: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
355 Ok(w.iter().rev().copied().collect())
356 }
357 }
358
359 struct NullStorage;
362 #[async_trait]
363 impl Storage for NullStorage {
364 async fn get(&self, _: &str) -> Result<GetObject, StorageError> {
365 Err(StorageError::NotFound(String::new()))
366 }
367 async fn get_range(
368 &self,
369 _: &str,
370 _: u64,
371 _: Option<u64>,
372 ) -> Result<GetObject, StorageError> {
373 Err(StorageError::NotFound(String::new()))
374 }
375 async fn put(
376 &self,
377 _: &str,
378 _: ByteStream,
379 _: PutMeta,
380 ) -> Result<ObjectMeta, StorageError> {
381 Err(StorageError::unsupported("null"))
382 }
383 async fn head(&self, _: &str) -> Result<ObjectMeta, StorageError> {
384 Err(StorageError::NotFound(String::new()))
385 }
386 async fn delete(&self, _: &str) -> Result<(), StorageError> {
387 Ok(())
388 }
389 async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
390 Ok(Vec::new())
391 }
392 }
393
394 fn managed_sql_cfg() -> crate::config::SqlBindingConfig {
396 let mut databases = std::collections::BTreeMap::new();
397 databases.insert(
398 "analytics".to_string(),
399 crate::config::ExternalDatabaseConfig {
400 kind: "postgres".into(),
401 compute: Some("pg".into()),
402 database: Some("analytics".into()),
403 user: Some("app".into()),
404 ..Default::default()
405 },
406 );
407 crate::config::SqlBindingConfig {
408 databases,
409 ..Default::default()
410 }
411 }
412
413 #[tokio::test]
414 async fn managed_sql_fails_closed_without_secrets() {
415 let tmp = tempfile::tempdir().unwrap();
416 let deploy = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
417 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
418 let cfg = managed_sql_cfg();
419 match build_sql_backends(Some(&cfg), tmp.path(), &deploy, &kv, None).await {
421 Err(Error::SqlManagedNeedsSecrets(name)) => assert_eq!(name, "analytics"),
422 Ok(_) => panic!("a managed DB without [secrets] must fail closed, got Ok"),
423 Err(other) => panic!("expected SqlManagedNeedsSecrets, got: {other}"),
424 }
425 }
426
427 #[tokio::test]
428 async fn managed_sql_builds_lazily_and_seals_the_credential() {
429 let tmp = tempfile::tempdir().unwrap();
430 let deploy = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
431 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
432 let envelope: Arc<dyn KeyEnvelope> = Arc::new(TestEnvelope);
433 let cfg = managed_sql_cfg();
434 let backends = build_sql_backends(Some(&cfg), tmp.path(), &deploy, &kv, Some(&envelope))
437 .await
438 .expect("managed sql builds without a live DB (lazy connect)");
439 let sealed = kv
440 .get("managed-sql-cred/default/pg")
441 .await
442 .unwrap()
443 .expect("managed credential sealed at build under the default project");
444 assert_ne!(sealed.len(), 0);
445 let _: Arc<dyn boatramp_core::sql::SqlBackends> = backends;
447 }
448}