1use std::future::Future;
15use std::net::{IpAddr, SocketAddr};
16use std::sync::Arc;
17use std::time::Duration;
18
19use axum::body::Body;
20use axum::extract::{ConnectInfo, Path, Query, Request, State};
21use axum::http::{header, HeaderMap, HeaderName, HeaderValue, Method, StatusCode};
22use axum::response::{IntoResponse, Response};
23use axum::routing::{any, get, post, put};
24use axum::{Extension, Json, Router};
25use boatramp_core::access::{AccessConfig, BasicAuth};
26use boatramp_core::authz::{GrantedRole, TokenMeta};
27use boatramp_core::config::{DeployConfig, SiteConfig};
28use boatramp_core::cose::{self, Claims, Signer};
29use boatramp_core::deploy::{
30 DeployMetaInput, DeployStore, FileEntry, GcOptions, GcReport, Manifest,
31};
32use boatramp_core::matcher::Pattern;
33use boatramp_core::route::{self, Outcome};
34use boatramp_core::{DeployError, StorageError};
35use futures::StreamExt;
36use serde::{Deserialize, Serialize};
37
38mod admin_api;
39pub mod sql_shim;
40#[cfg(feature = "oidc")]
41pub(crate) use admin_api::auth_exchange;
42pub(crate) use admin_api::{
43 activate_deployment, cert_status, compute_dns, compute_dns_resolve, compute_exec, compute_ipam,
44 compute_netdiag, compute_reconcile, compute_restart, compute_set_health, compute_status,
45 create_deployment, current_deployment, delete_compute, delete_compute_volume,
46 delete_project_tenancy, delete_site, get_compute, get_daemon_config, get_deployment,
47 get_project_tenancy, get_site_config, invalidate_cache, list_aliases, list_compute,
48 list_compute_volumes, list_deployments, list_sites, prune_delete, prune_report, put_blob,
49 put_compute, put_daemon_config, put_project_tenancy, put_site_config, remove_alias,
50 rollback_daemon_config, scrub_blobs, set_alias, sql_exec, sql_ping, sql_query,
51};
52#[cfg(feature = "handlers")]
53pub(crate) use admin_api::{
54 delete_graphql_safelist, delete_graphql_subgraph, get_graphql_supergraph,
55 list_graphql_safelist, put_graphql_function_subgraph, put_graphql_sql_subgraph,
56 put_graphql_subgraph, register_graphql_safelist,
57};
58#[cfg(feature = "admin")]
60mod admin_controller;
61mod auth;
62#[cfg(feature = "console")]
63pub mod console;
64mod content;
65mod control_api;
66#[cfg(feature = "email")]
69mod email_spool;
70#[cfg(feature = "session")]
72mod session_driver;
73#[cfg(feature = "session")]
76mod session_store;
77#[cfg(feature = "admin")]
78pub use admin_controller::ServerAdminController;
79#[cfg(feature = "compression")]
80pub(crate) use content::maybe_compress;
81pub(crate) use content::multipart_byteranges;
82pub(crate) use content::{
83 negotiate_encoding, parse_ranges, response_headers, set_content_encoding, MAX_RANGES,
84};
85pub(crate) use control_api::{
86 add_root_anchor, auth_whoami, bootstrap_token, cluster_join, cluster_members, cluster_promote,
87 cluster_revoke, cluster_rotate_key, create_join_token, create_token, delete_email_profile,
88 delete_secret, get_authz_policy, list_email_profiles, list_root_anchors, list_secrets,
89 list_tokens, put_authz_policy, remove_root_anchor, revoke_token, set_email_profile, set_secret,
90 show_email_profile,
91};
92#[cfg(all(test, feature = "handlers"))]
93use control_api::{BootstrapRequest, CreateJoinTokenRequest, JoinRequest};
94#[cfg(feature = "email")]
95pub use email_spool::NodeEmailSpool;
96mod domain_verify;
97pub use domain_verify::{spawn_domain_verify_reconcile, verification_pending_page};
98pub mod envelope;
99#[cfg(feature = "handlers")]
100mod graphql_apq;
101#[cfg(feature = "handlers")]
102mod graphql_cache;
103#[cfg(feature = "handlers")]
104mod graphql_data;
105#[cfg(feature = "handlers")]
106mod graphql_federation;
107#[cfg(feature = "handlers")]
108mod graphql_gateway;
109#[cfg(feature = "handlers")]
110mod graphql_graphiql;
111#[cfg(feature = "handlers")]
112mod graphql_guard;
113#[cfg(feature = "handlers")]
114mod graphql_plan;
115#[cfg(feature = "handlers")]
116mod graphql_registry;
117#[cfg(feature = "handlers")]
118mod graphql_subscription;
119#[cfg(feature = "handlers")]
120mod handler_cache;
121#[cfg(feature = "handlers")]
122mod handler_dispatch;
123#[cfg(feature = "handlers")]
124pub(crate) use handler_dispatch::{
125 build_bindings, dispatch_consumer_batch, dispatch_handler, precheck_component, read_blob_bytes,
126 read_blob_fully, resolve_secret_env,
127};
128#[cfg(all(feature = "handlers", test))]
129use handler_dispatch::{resolve_env, set_forwarded_headers};
130mod function_api;
131pub(crate) use function_api::{
132 alias_function, deploy_function, list_functions, remove_function, rollback_function,
133};
134#[cfg(feature = "handlers")]
139pub use function_api::{
140 component_requires, host_capability_features, host_capability_features_detailed, unmet_requires,
141};
142#[cfg(all(test, feature = "handlers"))]
143use function_api::{AliasBody, DeployFunctionQuery, FunctionUpsert, RollbackBody};
144pub use function_api::{CapabilityFeature, Lifecycle};
147mod gateway;
148mod host;
149pub(crate) use host::{is_local_host, parse_deploy_host, strip_port};
150#[cfg(feature = "http3")]
151mod http3;
152mod limits;
153#[cfg(feature = "handlers")]
154mod logs;
155#[cfg(feature = "handlers")]
156mod metrics;
157#[cfg(feature = "oidc")]
158mod oidc;
159mod operator;
160pub(crate) use operator::prometheus_metrics;
161#[cfg(feature = "handlers")]
162pub(crate) use operator::{
163 operator_dlq, operator_function_logs, operator_function_logs_stream, operator_handler_stats,
164 operator_logs, operator_logs_stream,
165};
166mod proxy;
167pub use proxy::spawn_compute_reconcile;
168pub(crate) use proxy::{
169 await_warm, compute_endpoint_regions, compute_endpoints, dispatch_gateway, has_parked_replica,
170 proxy, COMPUTE_WAKE_TIMEOUT,
171};
172mod splice;
173mod http_serve;
176pub use http_serve::{
177 alpn_h1_h2, serve_plaintext, serve_plaintext_listener, serve_router_conn, serve_tls,
178 serve_tls_listener, ReloadableTls, ServeInput,
179};
180#[cfg(feature = "handlers")]
182pub(crate) use proxy::is_upgrade_request;
183#[cfg(all(test, feature = "handlers"))]
184use proxy::{gateway_addr_allowed, CLOUD_METADATA_IPV4};
185mod project_api;
186pub(crate) use project_api::{create_project, delete_project, get_project, list_projects};
187mod project_scope;
188pub(crate) use project_scope::{project_scope, OriginalPath, ProjectContext};
189mod ratelimit;
190mod routes;
191pub use routes::{router, router_with, router_with_fast};
192#[cfg(feature = "mcp")]
193mod mcp_http;
194#[cfg(feature = "handlers")]
195mod scheduler;
196mod serve_pipeline;
197#[cfg(feature = "handlers")]
198mod tenant_resolve;
199pub use serve_pipeline::{http_redirect_router, FastServe};
200#[cfg(test)]
201mod hotpath_test;
202#[cfg(all(test, feature = "handlers"))]
203use serve_pipeline::{apply_vary, parse_cookie_header, parse_query_string};
204pub(crate) use serve_pipeline::{
205 serve_bootstrap_identity, serve_by_host, serve_domain_challenge, serve_preview, serve_sites,
206 BootstrapAttestation,
207};
208pub mod signer;
211mod srvmetrics;
212#[cfg(all(feature = "handlers", test))]
213use scheduler::run_scheduler_tick;
214#[cfg(feature = "handlers")]
215pub(crate) use scheduler::{
216 acquire_site_permit, effective_limits, handler_error_response, handler_unavailable, CronNow,
217};
218#[cfg(feature = "handlers")]
219use scheduler::{CONSUMER_BATCH, CONSUMER_LEASE, CONSUMER_MAX_ATTEMPTS};
220#[cfg(feature = "handlers")]
221mod function_runtime;
222#[cfg(feature = "handlers")]
223pub(crate) use function_runtime::{
224 b64_decode, b64_encode, blob_storage_prefix, capture_response, delete_trigger_handler,
225 dispatch_function_triggers, drain_function_invocations, execute_function, get_function_usage,
226 get_invocation_record, invoke_function, list_triggers_handler, new_invocation_id,
227 put_trigger_handler, webhook_ingress,
228};
229#[cfg(feature = "session")]
233mod session_serve;
234#[cfg(feature = "handlers")]
235mod stream;
236#[cfg(feature = "handlers")]
237mod workflow;
238pub use auth::{require_auth, Auth};
239#[cfg(feature = "http3")]
240pub use http3::{
241 advertise_http3, http3_endpoint, quinn_server_config, serve_http3, serve_http3_endpoint,
242 Http3Error,
243};
244pub use limits::{ServerLimits, UploadGuard};
245#[cfg(feature = "oidc")]
246pub use oidc::{OidcConfig, OidcError, OidcVerifier};
247use ratelimit::{KvRateLimiter, RateLimitStore, RateLimiter};
248#[cfg(feature = "handlers")]
249pub(crate) use stream::{route_matches, serve_stream, serve_ws_stream};
250#[cfg(feature = "handlers")]
251pub(crate) use workflow::{
252 define_workflow, delete_workflow_handler, drain_workflow_runs, get_workflow_handler,
253 get_workflow_run_handler, list_workflows_handler, start_workflow_run,
254};
255pub use srvmetrics::{server_metrics, ServerMetrics};
258
259#[derive(Clone, Default)]
264pub struct HandlerRuntime {
265 #[cfg(feature = "handlers")]
266 inner: Option<Arc<HandlerRuntimeInner>>,
267}
268
269#[cfg(feature = "handlers")]
270struct HandlerRuntimeInner {
271 engine: boatramp_handlers::HandlerEngine,
272 async_drain_gate: Arc<tokio::sync::Semaphore>,
278 kv: Arc<dyn boatramp_core::kv::KvStore>,
279 storage: Arc<dyn boatramp_core::Storage>,
280 sql: Option<Arc<dyn boatramp_core::sql::SqlBackends>>,
284 messaging: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
287 site_semaphores:
290 std::sync::Mutex<std::collections::HashMap<String, Arc<tokio::sync::Semaphore>>>,
291 stream_semaphores:
295 std::sync::Mutex<std::collections::HashMap<String, Arc<tokio::sync::Semaphore>>>,
296 stream_ip_counts: Arc<std::sync::Mutex<std::collections::HashMap<(String, IpAddr), u32>>>,
299 metrics: metrics::Metrics,
302 logs: Arc<logs::LogStore>,
304 #[cfg(feature = "handlers")]
308 graphql_cache: graphql_cache::GraphqlCache,
309 cron_leader_gate: std::sync::OnceLock<CronLeaderGate>,
315 max_blob_bytes: std::sync::OnceLock<u64>,
319 max_component_bytes: std::sync::OnceLock<u64>,
324 allow_env_secret_refs: std::sync::OnceLock<bool>,
333 require_tenancy_declaration: std::sync::OnceLock<bool>,
339 allow_cross_tenant_db: std::sync::OnceLock<bool>,
344 secret_store: std::sync::OnceLock<Arc<boatramp_core::secret_store::SecretStore>>,
350 function_meter_locks:
354 std::sync::Mutex<std::collections::HashMap<String, Arc<tokio::sync::Mutex<()>>>>,
355 function_semaphores:
358 std::sync::Mutex<std::collections::HashMap<String, Arc<tokio::sync::Semaphore>>>,
359 watch_provider: std::sync::OnceLock<Arc<dyn boatramp_core::blob_provision::WatchProvider>>,
364 provision_tier: std::sync::OnceLock<boatramp_core::blob_notify::ProvisionTier>,
368 invoker: std::sync::OnceLock<Arc<function_runtime::FunctionInvoker>>,
376 federation_runner: std::sync::OnceLock<Arc<graphql_gateway::FederationRunner>>,
381 #[cfg(feature = "email")]
386 email_profile_store: std::sync::OnceLock<Arc<boatramp_core::email_config::EmailProfileStore>>,
387 #[cfg(feature = "email")]
392 email_spool: std::sync::OnceLock<Arc<dyn boatramp_handlers::EmailSpool>>,
393 #[cfg(feature = "admin")]
397 admin_controller: std::sync::OnceLock<Arc<admin_controller::ServerAdminController>>,
398 #[cfg(feature = "admin")]
402 admin_surfaces:
403 std::sync::OnceLock<std::collections::BTreeSet<boatramp_handlers::AdminSurface>>,
404 #[cfg(feature = "session")]
411 session_store: std::sync::OnceLock<session_store::SessionStore>,
412 session_signer: std::sync::OnceLock<Arc<dyn Signer>>,
417}
418
419pub type CronLeaderGate = Arc<dyn Fn() -> bool + Send + Sync>;
422
423impl HandlerRuntime {
424 pub fn disabled() -> Self {
426 Self::default()
427 }
428
429 #[cfg(feature = "handlers")]
435 pub fn new(
436 engine: boatramp_handlers::HandlerEngine,
437 kv: Arc<dyn boatramp_core::kv::KvStore>,
438 storage: Arc<dyn boatramp_core::Storage>,
439 sql: Option<Arc<dyn boatramp_core::sql::SqlBackends>>,
440 messaging: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
441 ) -> Self {
442 let async_drain_slots = engine.async_max_concurrency().max(1);
445 Self {
446 inner: Some(Arc::new(HandlerRuntimeInner {
447 engine,
448 async_drain_gate: Arc::new(tokio::sync::Semaphore::new(async_drain_slots)),
449 kv,
450 storage,
451 sql,
452 messaging,
453 site_semaphores: std::sync::Mutex::new(std::collections::HashMap::new()),
454 stream_semaphores: std::sync::Mutex::new(std::collections::HashMap::new()),
455 stream_ip_counts: Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())),
456 metrics: metrics::Metrics::default(),
457 logs: Arc::new(logs::LogStore::default()),
458 #[cfg(feature = "handlers")]
459 graphql_cache: graphql_cache::GraphqlCache::default(),
460 cron_leader_gate: std::sync::OnceLock::new(),
461 max_blob_bytes: std::sync::OnceLock::new(),
462 max_component_bytes: std::sync::OnceLock::new(),
463 allow_env_secret_refs: std::sync::OnceLock::new(),
464 require_tenancy_declaration: std::sync::OnceLock::new(),
465 allow_cross_tenant_db: std::sync::OnceLock::new(),
466 secret_store: std::sync::OnceLock::new(),
467 function_meter_locks: std::sync::Mutex::new(std::collections::HashMap::new()),
468 function_semaphores: std::sync::Mutex::new(std::collections::HashMap::new()),
469 watch_provider: std::sync::OnceLock::new(),
470 provision_tier: std::sync::OnceLock::new(),
471 invoker: std::sync::OnceLock::new(),
472 federation_runner: std::sync::OnceLock::new(),
473 #[cfg(feature = "email")]
474 email_profile_store: std::sync::OnceLock::new(),
475 #[cfg(feature = "email")]
476 email_spool: std::sync::OnceLock::new(),
477 #[cfg(feature = "admin")]
478 admin_controller: std::sync::OnceLock::new(),
479 #[cfg(feature = "admin")]
480 admin_surfaces: std::sync::OnceLock::new(),
481 #[cfg(feature = "session")]
482 session_store: std::sync::OnceLock::new(),
483 session_signer: std::sync::OnceLock::new(),
484 })),
485 }
486 }
487
488 #[cfg(feature = "handlers")]
491 pub(crate) fn sql_provider(&self) -> Option<Arc<dyn boatramp_core::sql::SqlBackends>> {
492 self.inner.as_ref().and_then(|inner| inner.sql.clone())
493 }
494
495 #[cfg(feature = "handlers")]
499 pub(crate) fn invoker(&self) -> Option<Arc<function_runtime::FunctionInvoker>> {
500 self.inner
501 .as_ref()
502 .and_then(|inner| inner.invoker.get().cloned())
503 }
504
505 #[cfg(feature = "handlers")]
509 pub(crate) async fn introspect_subgraph_sdl(
510 &self,
511 deploy: &DeployStore,
512 project: boatramp_core::project::ProjectRef<'_>,
513 function: &boatramp_core::function::Function,
514 component: &str,
515 ) -> Result<String, function_runtime::SubgraphSdlError> {
516 match self.inner.as_ref() {
517 Some(inner) => {
518 function_runtime::introspect_service_sdl(
519 inner, deploy, project, function, component,
520 )
521 .await
522 }
523 None => Err(function_runtime::SubgraphSdlError::Unavailable),
524 }
525 }
526
527 #[cfg(feature = "handlers")]
533 pub fn set_invoker(&self, deploy: DeployStore) {
534 if let Some(inner) = self.inner.as_ref() {
535 let invoker = Arc::new(function_runtime::FunctionInvoker::new(
536 deploy,
537 Arc::downgrade(inner),
538 ));
539 let _ = inner.invoker.set(invoker);
540 let runner = Arc::new(graphql_gateway::FederationRunner::new(Arc::downgrade(
543 inner,
544 )));
545 let _ = inner.federation_runner.set(runner);
546 }
547 }
548
549 #[cfg(feature = "handlers")]
553 pub fn sql_backends(&self) -> Option<Arc<dyn boatramp_core::sql::SqlBackends>> {
554 self.inner.as_ref().and_then(|inner| inner.sql.clone())
555 }
556
557 #[cfg(feature = "handlers")]
561 pub fn set_watch_provider(
562 &self,
563 provider: Arc<dyn boatramp_core::blob_provision::WatchProvider>,
564 ) {
565 if let Some(inner) = self.inner.as_ref() {
566 let _ = inner.watch_provider.set(provider);
567 }
568 }
569
570 #[cfg(feature = "handlers")]
574 pub fn set_provision_tier(&self, tier: boatramp_core::blob_notify::ProvisionTier) {
575 if let Some(inner) = self.inner.as_ref() {
576 let _ = inner.provision_tier.set(tier);
577 }
578 }
579
580 #[cfg(feature = "handlers")]
584 pub fn set_max_blob_bytes(&self, max_bytes: u64) {
585 if let Some(inner) = self.inner.as_ref() {
586 let _ = inner.max_blob_bytes.set(max_bytes);
587 }
588 }
589
590 #[cfg(feature = "handlers")]
593 pub fn set_max_component_bytes(&self, max_bytes: u64) {
594 if let Some(inner) = self.inner.as_ref() {
595 let _ = inner.max_component_bytes.set(max_bytes);
596 }
597 }
598
599 #[cfg(feature = "handlers")]
607 pub fn set_allow_env_secret_refs(&self, allow: bool) {
608 if let Some(inner) = self.inner.as_ref() {
609 let _ = inner.allow_env_secret_refs.set(allow);
610 }
611 }
612
613 #[cfg(feature = "handlers")]
618 pub fn set_tenancy_posture(&self, require_declaration: bool, allow_cross_tenant: bool) {
619 if let Some(inner) = self.inner.as_ref() {
620 let _ = inner.require_tenancy_declaration.set(require_declaration);
621 let _ = inner.allow_cross_tenant_db.set(allow_cross_tenant);
622 }
623 }
624
625 #[cfg(feature = "handlers")]
631 pub fn set_session_signer(&self, signer: Arc<dyn Signer>) {
632 if let Some(inner) = self.inner.as_ref() {
633 let _ = inner.session_signer.set(signer);
634 }
635 }
636
637 #[cfg(feature = "handlers")]
642 pub fn set_secret_store(&self, store: Arc<boatramp_core::secret_store::SecretStore>) {
643 if let Some(inner) = self.inner.as_ref() {
644 let _ = inner.secret_store.set(store);
645 }
646 }
647
648 #[cfg(feature = "email")]
652 pub fn set_email_profile_store(
653 &self,
654 store: Arc<boatramp_core::email_config::EmailProfileStore>,
655 ) {
656 if let Some(inner) = self.inner.as_ref() {
657 let _ = inner.email_profile_store.set(store);
658 }
659 }
660
661 #[cfg(feature = "email")]
666 pub fn set_email_spool(&self, spool: Arc<dyn boatramp_handlers::EmailSpool>) {
667 if let Some(inner) = self.inner.as_ref() {
668 let _ = inner.email_spool.set(spool);
669 }
670 }
671
672 #[cfg(feature = "admin")]
677 pub fn set_admin(
678 &self,
679 controller: Arc<admin_controller::ServerAdminController>,
680 surfaces: std::collections::BTreeSet<boatramp_handlers::AdminSurface>,
681 ) {
682 if let Some(inner) = self.inner.as_ref() {
683 let _ = inner.admin_controller.set(controller);
684 let _ = inner.admin_surfaces.set(surfaces);
685 }
686 }
687
688 #[cfg(feature = "handlers")]
693 pub(crate) fn allow_env_secret_refs(&self) -> bool {
694 self.inner
695 .as_ref()
696 .and_then(|inner| inner.allow_env_secret_refs.get().copied())
697 .unwrap_or(false)
698 }
699
700 #[cfg(feature = "handlers")]
705 pub fn set_cron_leader_gate(&self, gate: CronLeaderGate) {
706 if let Some(inner) = self.inner.as_ref() {
707 let _ = inner.cron_leader_gate.set(gate);
708 }
709 }
710
711 #[cfg(feature = "handlers")]
718 async fn precheck_activation(
719 &self,
720 deploy: &DeployStore,
721 manifest: &Manifest,
722 site_config: Option<&SiteConfig>,
723 ) -> Result<(), String> {
724 let Some(inner) = self.inner.as_ref() else {
725 return Ok(());
726 };
727 if manifest.config.handlers.is_empty() && manifest.config.consumers.is_empty() {
730 return Ok(());
731 }
732 let site_handlers = site_config
734 .and_then(|c| c.handlers.as_ref())
735 .filter(|h| h.enabled)
736 .ok_or_else(|| {
737 "deployment ships handlers/consumers but the site has them disabled".to_string()
738 })?;
739 let max_component = inner.max_component_bytes.get().copied().unwrap_or(0);
740
741 let allow_env_secret_refs = inner.allow_env_secret_refs.get().copied().unwrap_or(false);
748 crate::handler_dispatch::admit_secret_refs(&site_handlers.secrets, allow_env_secret_refs)
749 .map_err(|err| format!("handler secrets: {err}"))?;
750
751 let sync_ceiling = inner.engine.sync_timeout_ms();
757 let async_ceiling = inner.engine.async_timeout_ms();
758 if let Some(ms) = site_handlers.max_timeout_ms {
759 if u64::from(ms) > sync_ceiling {
760 tracing::warn!(
761 "site max_timeout_ms={ms} exceeds sync_max_timeout_ms={sync_ceiling}: \
762 synchronous HTTP handlers are capped at {sync_ceiling}ms; the extra time \
763 applies only to async calls (?mode=async / triggers), capped at \
764 async_max_timeout_ms={async_ceiling}"
765 );
766 }
767 }
768
769 for handler in &manifest.config.handlers {
771 if let Some(ms) = handler.limits.as_ref().and_then(|l| l.timeout_ms) {
772 if u64::from(ms) > sync_ceiling {
773 let route = &handler.route;
774 tracing::warn!(
775 "route {route:?} declares limits.timeout_ms={ms}, above \
776 sync_max_timeout_ms={sync_ceiling}: synchronous HTTP calls to this route \
777 are capped at {sync_ceiling}ms; the {ms}ms only applies to async calls \
778 (?mode=async / a queue trigger / a #[consumer]), capped at \
779 async_max_timeout_ms={async_ceiling}. If you need {ms}ms synchronously, \
780 that isn't possible — move the work to the async lane"
781 );
782 }
783 }
784 if !handler.streaming {
789 if let Some(entry) = manifest.files.get(&handler.component) {
790 if let Ok(bytes) = read_blob_bytes(deploy, &entry.hash).await {
791 if crate::function_api::component_declares_streaming_route(
792 &bytes,
793 &handler.route,
794 ) {
795 let route = &handler.route;
796 tracing::warn!(
797 "route {route:?} is a streaming handler (#[handler(stream)]) but \
798 its config lacks streaming = true: it will run on the sync request \
799 lane and be cut at sync_max_timeout_ms={sync_ceiling}ms. Set \
800 streaming = true so it serves on the dedicated streaming lane (its \
801 own concurrency budget + a much larger wall-clock)."
802 );
803 }
804 }
805 }
806 }
807 if let Some(entry) = manifest.files.get(&handler.component) {
811 if let Ok(bytes) = read_blob_bytes(deploy, &entry.hash).await {
812 let unmet = crate::function_api::unmet_requires(&bytes);
813 if !unmet.is_empty() {
814 return Err(format!(
815 "route {:?} [{}] requires capabilities this host does not implement: \
816 {}. Upgrade boatramp or enable those features — see `boatramp \
817 capabilities`.",
818 handler.route,
819 handler.methods.join(","),
820 unmet.join(", ")
821 ));
822 }
823 }
824 }
825 precheck_component(
826 deploy,
827 manifest,
828 site_handlers,
829 inner,
830 max_component,
831 &handler.imports,
832 &handler.component,
833 &format!("route {:?} [{}]", handler.route, handler.methods.join(",")),
836 false,
837 )
838 .await?;
839 }
840 for consumer in &manifest.config.consumers {
841 precheck_component(
842 deploy,
843 manifest,
844 site_handlers,
845 inner,
846 max_component,
847 &consumer.imports,
848 &consumer.component,
849 &format!("consumer {:?}", consumer.topic),
850 true,
851 )
852 .await?;
853 }
854 Ok(())
855 }
856
857 #[cfg(not(feature = "handlers"))]
858 async fn precheck_activation(
859 &self,
860 _deploy: &DeployStore,
861 _manifest: &Manifest,
862 _site_config: Option<&SiteConfig>,
863 ) -> Result<(), String> {
864 Ok(())
865 }
866}
867
868#[derive(Default, Clone)]
874pub struct ServerOptions {
875 pub limits: ServerLimits,
877 pub probe: Option<Arc<dyn boatramp_core::domain_verify::DomainProbe>>,
880 pub default_site: Option<String>,
883 pub implicit_routing: bool,
889 pub protect_previews: bool,
894 pub cluster_rate_limit_kv: Option<Arc<dyn boatramp_core::kv::KvStore>>,
898 pub issuer: Option<Arc<dyn Signer>>,
902 pub bootstrap_secret: Option<String>,
907 pub bootstrap_attestation: Option<String>,
913 pub mesh_control: Option<Arc<dyn MeshControl>>,
917 pub cors_allowed_origins: Vec<String>,
927 #[cfg(feature = "oidc")]
930 pub oidc_verifier: Option<Arc<oidc::OidcVerifier>>,
931 pub posture: boatramp_core::security::SecurityPosture,
935 pub served_over_tls: bool,
940 pub pop_origin: Option<String>,
947 pub daemon_runtime: Option<Arc<DaemonRuntime>>,
951 pub operator_sql: Option<Arc<dyn boatramp_core::sql::OperatorSql>>,
955 pub tenant_deprovisioner: Option<Arc<dyn boatramp_core::sql::TenantDeprovisioner>>,
960 pub compute_exec: Option<Arc<dyn boatramp_core::compute::ComputeExec>>,
964 pub compute_volumes: Option<Arc<dyn boatramp_core::compute::ComputeVolumes>>,
969 pub compute_control: Option<Arc<dyn boatramp_core::compute::ComputeControl>>,
973 pub secret_store: Option<Arc<boatramp_core::secret_store::SecretStore>>,
981 pub email_profile_store: Option<Arc<boatramp_core::email_config::EmailProfileStore>>,
989 #[cfg(feature = "console")]
993 pub console: Option<console::ConsoleMount>,
994}
995
996#[derive(Clone, Copy)]
1000struct ServedOverTls(bool);
1001
1002#[derive(Clone, Copy, Default)]
1007struct ImplicitRouting(bool);
1008
1009const DAEMON_RELOAD_BACKSTOP: std::time::Duration = std::time::Duration::from_secs(300);
1021
1022pub struct DaemonRuntime {
1023 baseline: boatramp_core::daemon_config::ConfigBaseline,
1024 state: std::sync::RwLock<DaemonState>,
1025 reload: tokio::sync::Notify,
1028}
1029
1030struct DaemonState {
1031 effective: Arc<boatramp_core::daemon_config::EffectiveConfig>,
1032 generation: Option<String>,
1033}
1034
1035pub fn config_baseline(options: &ServerOptions) -> boatramp_core::daemon_config::ConfigBaseline {
1040 #[cfg(feature = "console")]
1044 let (console_enabled, console_host, console_path) = match options.console.as_ref() {
1045 Some(m) => (true, Some(m.host.clone()), Some(m.path.clone())),
1046 None => (false, None, None),
1047 };
1048 #[cfg(not(feature = "console"))]
1049 let (console_enabled, console_host, console_path) = (false, None, None);
1050 boatramp_core::daemon_config::ConfigBaseline {
1051 default_site: options.default_site.clone(),
1052 protect_previews: options.protect_previews,
1053 max_upload_bytes: options.limits.max_upload_bytes.unwrap_or(0),
1054 upload_idle_timeout_secs: options.limits.upload_idle_timeout.map(|d| d.as_secs()),
1055 max_concurrent_uploads: options.limits.max_concurrent_uploads.map(|n| n as u64),
1056 cluster_rate_limit: options.cluster_rate_limit_kv.is_some(),
1057 compute_vcpus: 0,
1058 compute_mem_mib: 0,
1059 console_enabled,
1060 console_host,
1061 console_path,
1062 max_upload_ceiling: options.posture.max_upload_bytes,
1063 max_concurrent_uploads_ceiling: None,
1064 posture: options.posture,
1065 }
1066}
1067
1068impl DaemonRuntime {
1069 pub fn new(baseline: boatramp_core::daemon_config::ConfigBaseline) -> Self {
1073 let effective =
1074 Arc::new(boatramp_core::daemon_config::DaemonConfig::default().resolve(&baseline));
1075 Self {
1076 baseline,
1077 state: std::sync::RwLock::new(DaemonState {
1078 effective,
1079 generation: None,
1080 }),
1081 reload: tokio::sync::Notify::new(),
1082 }
1083 }
1084
1085 pub fn notify_reload(&self) {
1089 self.reload.notify_one();
1090 }
1091
1092 pub fn effective(&self) -> Arc<boatramp_core::daemon_config::EffectiveConfig> {
1094 self.state
1095 .read()
1096 .expect("daemon config lock")
1097 .effective
1098 .clone()
1099 }
1100
1101 pub fn generation(&self) -> Option<String> {
1104 self.state
1105 .read()
1106 .expect("daemon config lock")
1107 .generation
1108 .clone()
1109 }
1110
1111 pub fn baseline(&self) -> &boatramp_core::daemon_config::ConfigBaseline {
1113 &self.baseline
1114 }
1115
1116 pub async fn reload(&self, deploy: &DeployStore) -> Result<(), DeployError> {
1119 let cfg = deploy.get_daemon_config().await?.unwrap_or_default();
1120 let generation = deploy.daemon_config_generation().await?;
1121 let effective = Arc::new(cfg.resolve(&self.baseline));
1122 *self.state.write().expect("daemon config lock") = DaemonState {
1123 effective,
1124 generation,
1125 };
1126 Ok(())
1127 }
1128}
1129
1130#[derive(Clone, Copy, Default)]
1133struct PreviewPolicy {
1134 protect: bool,
1135}
1136
1137#[derive(Clone, Default)]
1142struct Issuer(Option<Arc<dyn Signer>>);
1143
1144#[derive(Clone, Default)]
1149struct BootstrapGate(Option<Arc<BootstrapInner>>);
1150
1151struct BootstrapInner {
1152 secret_hash: String,
1155 lock: tokio::sync::Mutex<()>,
1158}
1159
1160impl BootstrapGate {
1161 fn new(secret: Option<&str>) -> Self {
1162 Self(secret.filter(|s| !s.is_empty()).map(|s| {
1163 Arc::new(BootstrapInner {
1164 secret_hash: boatramp_core::deploy::sha256_hex(s.as_bytes()),
1165 lock: tokio::sync::Mutex::new(()),
1166 })
1167 }))
1168 }
1169}
1170
1171#[async_trait::async_trait]
1175pub trait MeshControl: Send + Sync {
1176 async fn admit(
1184 &self,
1185 mesh_pubkey_hex: &str,
1186 jti: &str,
1187 possession_proof: &[u8],
1188 proof_iat: u64,
1189 now: u64,
1190 advertise_addr: Option<&str>,
1191 ) -> Result<JoinOutcome, String>;
1192
1193 async fn rotate_key(&self) -> Result<String, String>;
1197
1198 async fn revoke(&self, node: u64) -> Result<(), String>;
1202
1203 async fn members(&self) -> Result<Vec<MeshMember>, String>;
1207
1208 async fn promote(&self, node: u64) -> Result<(), String>;
1211}
1212
1213pub enum JoinOutcome {
1215 Admitted {
1218 members: Vec<String>,
1220 addrs: std::collections::BTreeMap<u64, String>,
1222 },
1223 TokenSpent,
1225 ProofInvalid,
1227 Revoked,
1230}
1231
1232#[derive(Debug, Clone, Serialize)]
1234pub struct MeshMember {
1235 pub node: u64,
1237 pub voter: bool,
1239 pub caught_up: bool,
1241 pub leader: bool,
1243 #[serde(default, skip_serializing_if = "Option::is_none")]
1247 pub addr: Option<String>,
1248}
1249
1250#[derive(Clone, Default)]
1253struct MeshControlHandle(Option<Arc<dyn MeshControl>>);
1254
1255#[cfg(feature = "oidc")]
1257#[derive(Clone, Default)]
1258struct OidcState(Option<Arc<oidc::OidcVerifier>>);
1259
1260#[cfg(feature = "oidc")]
1263const EXCHANGE_TTL_SECS: u64 = 3600;
1264
1265use boatramp_core::time::now_unix;
1266
1267#[derive(Clone)]
1269struct CorsState(Arc<Vec<String>>);
1270
1271const CORS_ALLOW_METHODS: &str = "GET, POST, PUT, DELETE, OPTIONS";
1273const CORS_ALLOW_HEADERS: &str = "authorization, content-type";
1276const CORS_MAX_AGE: &str = "600";
1278
1279fn cors_origin_allowed(allowed: &[String], origin: &str) -> bool {
1283 allowed.iter().any(|a| a == "*" || a == origin)
1284}
1285
1286async fn cors(
1293 State(allowed): State<CorsState>,
1294 request: Request,
1295 next: axum::middleware::Next,
1296) -> Response {
1297 let origin = request
1298 .headers()
1299 .get(header::ORIGIN)
1300 .and_then(|v| v.to_str().ok())
1301 .filter(|o| cors_origin_allowed(&allowed.0, o))
1302 .map(str::to_string);
1303 let is_preflight = request.method() == Method::OPTIONS
1305 && request
1306 .headers()
1307 .contains_key(header::ACCESS_CONTROL_REQUEST_METHOD);
1308 if is_preflight {
1309 let allow_headers = request
1311 .headers()
1312 .get(header::ACCESS_CONTROL_REQUEST_HEADERS)
1313 .and_then(|v| v.to_str().ok())
1314 .map(str::to_string)
1315 .unwrap_or_else(|| CORS_ALLOW_HEADERS.to_string());
1316 let mut response = Response::new(Body::empty());
1317 *response.status_mut() = StatusCode::NO_CONTENT;
1318 if let Some(origin) = origin {
1319 let headers = response.headers_mut();
1320 set_header(headers, header::ACCESS_CONTROL_ALLOW_ORIGIN, &origin);
1321 set_header(headers, header::VARY, "Origin");
1322 set_header(
1323 headers,
1324 header::ACCESS_CONTROL_ALLOW_METHODS,
1325 CORS_ALLOW_METHODS,
1326 );
1327 set_header(
1328 headers,
1329 header::ACCESS_CONTROL_ALLOW_HEADERS,
1330 &allow_headers,
1331 );
1332 set_header(headers, header::ACCESS_CONTROL_MAX_AGE, CORS_MAX_AGE);
1333 }
1334 return response;
1335 }
1336 let mut response = next.run(request).await;
1337 if let Some(origin) = origin {
1338 let headers = response.headers_mut();
1339 set_header(headers, header::ACCESS_CONTROL_ALLOW_ORIGIN, &origin);
1340 if let Ok(value) = HeaderValue::from_str("Origin") {
1343 headers.append(header::VARY, value);
1344 }
1345 }
1346 response
1347}
1348
1349const DRAIN_DEADLINE: Duration = Duration::from_secs(30);
1354
1355#[derive(Debug, thiserror::Error)]
1357pub enum ServeError {
1358 #[error("server I/O: {0}")]
1360 Io(#[from] std::io::Error),
1361}
1362
1363pub async fn serve(
1366 addr: SocketAddr,
1367 deploy: DeployStore,
1368 auth: Auth,
1369 handlers: HandlerRuntime,
1370) -> Result<(), ServeError> {
1371 serve_with(addr, deploy, auth, handlers, ServerOptions::default()).await
1372}
1373
1374pub(crate) fn disable_nagle(stream: &mut tokio::net::TcpStream) {
1384 if let Err(err) = stream.set_nodelay(true) {
1385 tracing::debug!(%err, "failed to set TCP_NODELAY on an accepted connection");
1386 }
1387}
1388
1389pub async fn serve_with(
1391 addr: SocketAddr,
1392 deploy: DeployStore,
1393 auth: Auth,
1394 handlers: HandlerRuntime,
1395 options: ServerOptions,
1396) -> Result<(), ServeError> {
1397 let tcp = tokio::net::TcpListener::bind(addr).await?;
1398 tracing::info!(%addr, auth = !auth.is_disabled(), "boatramp server listening");
1399 let splice_ctx = splice::SpliceCtx {
1405 deploy: deploy.clone(),
1406 posture: options.posture,
1407 daemon: options.daemon_runtime.clone(),
1408 };
1409 #[cfg(feature = "handlers")]
1412 let scheduler = handlers.spawn_scheduler(deploy.clone());
1413 let gateway_prober = gateway::spawn_active_health_prober();
1417 let (router, fast) = router_with_fast(deploy, auth, handlers, options);
1420
1421 let (signalled_tx, signalled_rx) = tokio::sync::watch::channel(false);
1425 let server = splice::serve(tcp, splice_ctx, (router, fast), async move {
1429 shutdown_signal().await;
1430 let _ = signalled_tx.send(true);
1431 });
1432 let signalled = {
1433 let mut rx = signalled_rx;
1434 async move {
1435 let _ = rx.wait_for(|fired| *fired).await;
1436 }
1437 };
1438 let result = serve_with_drain_deadline(
1439 async move { server.await.map_err(ServeError::from) },
1440 signalled,
1441 DRAIN_DEADLINE,
1442 )
1443 .await;
1444 #[cfg(feature = "handlers")]
1446 if let Some(handle) = scheduler {
1447 handle.abort();
1448 }
1449 gateway_prober.abort();
1450 result
1451}
1452
1453async fn serve_with_drain_deadline<Srv, Sig>(
1458 server: Srv,
1459 signalled: Sig,
1460 deadline: Duration,
1461) -> Result<(), ServeError>
1462where
1463 Srv: Future<Output = Result<(), ServeError>>,
1464 Sig: Future<Output = ()>,
1465{
1466 tokio::pin!(server);
1467 let drain_cap = async move {
1468 signalled.await;
1469 tokio::time::sleep(deadline).await;
1470 };
1471 tokio::select! {
1472 result = &mut server => result,
1473 _ = drain_cap => {
1474 tracing::warn!(
1475 deadline_s = deadline.as_secs(),
1476 "drain deadline exceeded; forcing shutdown with requests still in flight"
1477 );
1478 Ok(())
1479 }
1480 }
1481}
1482
1483pub async fn shutdown_signal() {
1486 let ctrl_c = async {
1487 let _ = tokio::signal::ctrl_c().await;
1488 };
1489 #[cfg(unix)]
1490 let terminate = async {
1491 if let Ok(mut sig) =
1492 tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
1493 {
1494 sig.recv().await;
1495 }
1496 };
1497 #[cfg(not(unix))]
1498 let terminate = std::future::pending::<()>();
1499
1500 tokio::select! {
1501 _ = ctrl_c => {}
1502 _ = terminate => {}
1503 }
1504 tracing::info!("shutdown signal received; draining");
1505}
1506
1507async fn healthz(Extension(daemon): Extension<Arc<DaemonRuntime>>) -> String {
1511 match daemon.generation() {
1512 Some(gen) => format!("ok gen={gen}"),
1513 None => "ok".to_string(),
1514 }
1515}
1516
1517async fn readyz(State(deploy): State<DeployStore>) -> Response {
1519 match deploy.ready().await {
1520 Ok(()) => (StatusCode::OK, "ready\n").into_response(),
1521 Err(err) => {
1522 tracing::warn!(error = %err, "readiness probe failed");
1523 (StatusCode::SERVICE_UNAVAILABLE, "not ready\n").into_response()
1524 }
1525 }
1526}
1527
1528#[derive(Clone)]
1533pub struct RequestId(pub String);
1534
1535#[cfg(feature = "handlers")]
1541#[derive(Clone)]
1542pub struct DomainContext(pub String);
1543
1544fn request_id_for(headers: &HeaderMap) -> String {
1547 if let Some(id) = headers
1548 .get("x-request-id")
1549 .and_then(|v| v.to_str().ok())
1550 .map(str::trim)
1551 .filter(|s| !s.is_empty())
1552 {
1553 return id.chars().filter(|c| !c.is_control()).take(128).collect();
1554 }
1555 use std::sync::atomic::{AtomicU64, Ordering};
1556 static SEQ: AtomicU64 = AtomicU64::new(0);
1557 let n = SEQ.fetch_add(1, Ordering::Relaxed);
1558 format!("{:x}-{:x}", boatramp_core::time::now_unix_ms(), n)
1559}
1560
1561struct AccessLog {
1565 request_id: String,
1566 method: Method,
1567 path: String,
1568 host: String,
1569 client: String,
1570 status: u16,
1571 encoding: String,
1573 start: std::time::Instant,
1574 bytes: std::sync::atomic::AtomicU64,
1575}
1576
1577impl Drop for AccessLog {
1578 fn drop(&mut self) {
1579 let bytes = self.bytes.load(std::sync::atomic::Ordering::Relaxed);
1580 srvmetrics::server_metrics().record_request(self.status, bytes);
1583 tracing::info!(
1584 target: "boatramp::access",
1585 request_id = %self.request_id,
1586 method = %self.method,
1587 path = %self.path,
1588 host = %self.host,
1589 client = %self.client,
1590 status = self.status,
1591 bytes = bytes,
1592 encoding = %self.encoding,
1593 cache_result = srvmetrics::cache_result(self.status),
1594 elapsed_ms = self.start.elapsed().as_millis() as u64,
1595 "request"
1596 );
1597 }
1598}
1599
1600pub(crate) fn assign_request_id(request: &mut axum::extract::Request) -> String {
1606 let request_id = request_id_for(request.headers());
1607 request
1608 .extensions_mut()
1609 .insert(RequestId(request_id.clone()));
1610 request_id
1611}
1612
1613pub(crate) struct AccessLogCtx {
1618 request_id: String,
1619 method: Method,
1620 path: String,
1621 host: String,
1622 client: String,
1623 start: std::time::Instant,
1624}
1625
1626impl AccessLogCtx {
1627 pub(crate) fn capture(request: &axum::extract::Request, request_id: String) -> Option<Self> {
1633 if !tracing::enabled!(target: "boatramp::access", tracing::Level::INFO) {
1634 return None;
1635 }
1636 Some(Self {
1637 request_id,
1638 method: request.method().clone(),
1639 path: request.uri().path().to_string(),
1640 host: request
1641 .headers()
1642 .get(header::HOST)
1643 .and_then(|value| value.to_str().ok())
1644 .or_else(|| request.uri().host()) .unwrap_or("-")
1646 .to_string(),
1647 client: request
1648 .extensions()
1649 .get::<axum::extract::ConnectInfo<SocketAddr>>()
1650 .map(|info| info.0.ip().to_string())
1651 .unwrap_or_else(|| "-".to_string()),
1652 start: std::time::Instant::now(),
1653 })
1654 }
1655
1656 pub(crate) fn finish(self, response: Response) -> Response {
1660 let encoding = response
1661 .headers()
1662 .get(header::CONTENT_ENCODING)
1663 .and_then(|v| v.to_str().ok())
1664 .unwrap_or("identity")
1665 .to_string();
1666 let log = AccessLog {
1667 request_id: self.request_id,
1668 method: self.method,
1669 path: self.path,
1670 host: self.host,
1671 client: self.client,
1672 status: response.status().as_u16(),
1673 encoding,
1674 start: self.start,
1675 bytes: std::sync::atomic::AtomicU64::new(0),
1676 };
1677 let (parts, body) = response.into_parts();
1678 let counted = body.into_data_stream().map(move |chunk| {
1679 if let Ok(bytes) = &chunk {
1680 log.bytes
1681 .fetch_add(bytes.len() as u64, std::sync::atomic::Ordering::Relaxed);
1682 }
1683 chunk
1684 });
1685 Response::from_parts(parts, Body::from_stream(counted))
1686 }
1687}
1688
1689async fn access_log(mut request: axum::extract::Request, next: axum::middleware::Next) -> Response {
1694 let request_id = assign_request_id(&mut request);
1695 match AccessLogCtx::capture(&request, request_id) {
1696 None => next.run(request).await,
1697 Some(ctx) => ctx.finish(next.run(request).await),
1698 }
1699}
1700
1701fn if_none_match(req_headers: &HeaderMap, etag: &str) -> bool {
1703 req_headers
1704 .get(header::IF_NONE_MATCH)
1705 .and_then(|value| value.to_str().ok())
1706 .is_some_and(|value| {
1707 value
1708 .split(',')
1709 .map(str::trim)
1710 .any(|tag| tag == "*" || tag == etag || tag.trim_start_matches("W/") == etag)
1711 })
1712}
1713
1714fn set_header(headers: &mut HeaderMap, name: header::HeaderName, value: &str) {
1715 if let Ok(value) = HeaderValue::from_str(value) {
1716 headers.insert(name, value);
1717 }
1718}
1719
1720fn not_found() -> Response {
1721 (StatusCode::NOT_FOUND, "not found\n").into_response()
1722}
1723
1724fn redirect(status: u16, location: &str) -> Response {
1725 let status = StatusCode::from_u16(status).unwrap_or(StatusCode::FOUND);
1726 match HeaderValue::from_str(location) {
1727 Ok(location) => {
1728 let mut headers = HeaderMap::new();
1729 headers.insert(header::LOCATION, location);
1730 (status, headers).into_response()
1731 }
1732 Err(_) => (StatusCode::INTERNAL_SERVER_ERROR, "bad redirect target\n").into_response(),
1733 }
1734}
1735
1736fn deploy_error_response(err: DeployError) -> Response {
1738 let status = match &err {
1739 DeployError::NotFound(_) | DeployError::Storage(StorageError::NotFound(_)) => {
1740 StatusCode::NOT_FOUND
1741 }
1742 DeployError::HashMismatch { .. } => StatusCode::BAD_REQUEST,
1743 DeployError::Incomplete(_) => StatusCode::CONFLICT,
1744 DeployError::Conflict(_) => StatusCode::CONFLICT,
1746 DeployError::Ambiguous(_) => StatusCode::NOT_FOUND,
1748 DeployError::Invalid(_) => StatusCode::BAD_REQUEST,
1750 _ => StatusCode::INTERNAL_SERVER_ERROR,
1751 };
1752 tracing::warn!(error = %err, "request failed");
1753 (status, format!("{err}\n")).into_response()
1754}
1755
1756fn reject_invalid_name(kind: &'static str, value: &str) -> Option<Response> {
1761 boatramp_core::project::validate_resource_name(kind, value)
1762 .err()
1763 .map(|err| (StatusCode::UNPROCESSABLE_ENTITY, format!("{err}\n")).into_response())
1764}
1765
1766#[cfg(test)]
1767mod drain_tests {
1768 use super::*;
1769
1770 #[tokio::test]
1771 async fn deadline_forces_shutdown_after_signal() {
1772 let server = std::future::pending::<Result<(), ServeError>>();
1775 let signalled = async {}; let result = serve_with_drain_deadline(server, signalled, Duration::from_millis(20)).await;
1777 assert!(result.is_ok());
1778 }
1779
1780 #[tokio::test]
1781 async fn server_finishing_first_wins() {
1782 let server = async { Ok(()) };
1785 let signalled = std::future::pending::<()>();
1786 let result = serve_with_drain_deadline(server, signalled, Duration::from_secs(30)).await;
1787 assert!(result.is_ok());
1788 }
1789
1790 #[tokio::test]
1791 async fn deadline_does_not_trip_before_signal() {
1792 let server = async {
1796 tokio::time::sleep(Duration::from_millis(40)).await;
1797 Err(ServeError::Io(std::io::Error::other("server error")))
1798 };
1799 let signalled = std::future::pending::<()>();
1800 let result = serve_with_drain_deadline(server, signalled, Duration::from_millis(10)).await;
1801 assert!(result.is_err());
1802 }
1803}
1804
1805#[cfg(all(test, feature = "handlers"))]
1806mod tests {
1807 use super::*;
1808 use boatramp_core::cose::{LocalSigner, TokenAlg};
1809 use boatramp_core::project::ProjectRef;
1810
1811 #[test]
1812 fn query_string_parses_and_url_decodes() {
1813 let q = parse_query_string("lang=fr&city=S%C3%A3o+Paulo&flag&dup=1&dup=2");
1814 assert_eq!(q.get("lang").map(String::as_str), Some("fr"));
1815 assert_eq!(q.get("city").map(String::as_str), Some("São Paulo")); assert_eq!(q.get("flag").map(String::as_str), Some("")); assert_eq!(q.get("dup").map(String::as_str), Some("1")); }
1819
1820 #[test]
1821 fn cookie_header_parses_pairs() {
1822 let c = parse_cookie_header("beta=1; sid = abc ; empty=");
1823 assert_eq!(c.get("beta").map(String::as_str), Some("1"));
1824 assert_eq!(c.get("sid").map(String::as_str), Some("abc"));
1825 assert_eq!(c.get("empty").map(String::as_str), Some(""));
1826 }
1827
1828 #[test]
1829 fn apply_vary_merges_without_duplicates() {
1830 let base = (StatusCode::OK, "x").into_response();
1831 let r = apply_vary(base, &["accept-language".into()]);
1832 assert_eq!(r.headers().get(header::VARY).unwrap(), "accept-language");
1833 let r = apply_vary(r, &["cookie".into(), "accept-language".into()]);
1835 let v = r.headers().get(header::VARY).unwrap().to_str().unwrap();
1836 assert!(v.contains("accept-language") && v.contains("cookie"));
1837 assert_eq!(v.matches("accept-language").count(), 1);
1838 let plain = apply_vary((StatusCode::OK, "y").into_response(), &[]);
1840 assert!(plain.headers().get(header::VARY).is_none());
1841 }
1842
1843 #[tokio::test]
1847 async fn join_token_endpoint_mints_a_verifiable_bearer_token() {
1848 let keys: Arc<dyn Signer> = Arc::new(LocalSigner::generate(TokenAlg::Es256));
1849 let public = keys.public_key();
1850
1851 let resp = create_join_token(
1853 Extension(Issuer(Some(keys.clone()))),
1854 Json(CreateJoinTokenRequest {
1855 ttl_secs: Some(600),
1856 }),
1857 )
1858 .await;
1859 assert_eq!(resp.status(), StatusCode::CREATED);
1860 let body = axum::body::to_bytes(resp.into_body(), usize::MAX)
1861 .await
1862 .unwrap();
1863 let parsed: serde_json::Value = serde_json::from_slice(&body).unwrap();
1864 let token = parsed["token"].as_str().unwrap();
1865 let jti = cose::verify_join(token, &public, now_unix()).unwrap();
1866 assert!(!jti.is_empty());
1867
1868 let no_issuer = create_join_token(
1870 Extension(Issuer(None)),
1871 Json(CreateJoinTokenRequest { ttl_secs: None }),
1872 )
1873 .await;
1874 assert_eq!(no_issuer.status(), StatusCode::NOT_IMPLEMENTED);
1875 }
1876
1877 #[tokio::test]
1882 async fn function_write_path_deploy_rollback_alias_remove() {
1883 use boatramp_core::function::Lifecycle;
1884 use boatramp_core::kv::MemoryKv;
1885 use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
1886
1887 struct FakeStorage {
1890 present: bool,
1891 }
1892 #[async_trait::async_trait]
1893 impl Storage for FakeStorage {
1894 async fn get(&self, _: &str) -> Result<GetObject, StorageError> {
1895 Err(StorageError::NotFound(String::new()))
1896 }
1897 async fn get_range(
1898 &self,
1899 _: &str,
1900 _: u64,
1901 _: Option<u64>,
1902 ) -> Result<GetObject, StorageError> {
1903 Err(StorageError::NotFound(String::new()))
1904 }
1905 async fn put(
1906 &self,
1907 _: &str,
1908 _: ByteStream,
1909 _: PutMeta,
1910 ) -> Result<ObjectMeta, StorageError> {
1911 Err(StorageError::unsupported("fake"))
1912 }
1913 async fn head(&self, key: &str) -> Result<ObjectMeta, StorageError> {
1914 if self.present {
1915 Ok(ObjectMeta {
1916 key: key.to_string(),
1917 ..Default::default()
1918 })
1919 } else {
1920 Err(StorageError::NotFound(key.to_string()))
1921 }
1922 }
1923 async fn delete(&self, _: &str) -> Result<(), StorageError> {
1924 Ok(())
1925 }
1926 async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
1927 Ok(Vec::new())
1928 }
1929 }
1930
1931 async fn body_json(resp: Response) -> (StatusCode, serde_json::Value) {
1932 let status = resp.status();
1933 let bytes = axum::body::to_bytes(resp.into_body(), usize::MAX)
1934 .await
1935 .unwrap();
1936 let value = if bytes.is_empty() {
1937 serde_json::Value::Null
1938 } else {
1939 serde_json::from_slice(&bytes).unwrap()
1940 };
1941 (status, value)
1942 }
1943
1944 let deploy = DeployStore::new(
1945 Arc::new(FakeStorage { present: true }),
1946 Arc::new(MemoryKv::new()),
1947 );
1948 let v1 = "a".repeat(64);
1949 let v2 = "b".repeat(64);
1950
1951 let (st, body) = body_json(
1953 deploy_function(
1954 State(deploy.clone()),
1955 axum::extract::Extension(crate::ProjectContext::default()),
1956 axum::extract::Extension(Arc::new(crate::HandlerRuntime::disabled())),
1957 axum::extract::Query(DeployFunctionQuery::default()),
1958 Path("greeter".to_string()),
1959 Json(FunctionUpsert {
1960 component: v1.clone(),
1961 config: Default::default(),
1962 lifecycle: Lifecycle::Independent,
1963 }),
1964 )
1965 .await,
1966 )
1967 .await;
1968 assert_eq!(st, StatusCode::OK);
1969 assert_eq!(body["active"], v1);
1970
1971 let (_, body) = body_json(
1973 deploy_function(
1974 State(deploy.clone()),
1975 axum::extract::Extension(crate::ProjectContext::default()),
1976 axum::extract::Extension(Arc::new(crate::HandlerRuntime::disabled())),
1977 axum::extract::Query(DeployFunctionQuery::default()),
1978 Path("greeter".to_string()),
1979 Json(FunctionUpsert {
1980 component: v2.clone(),
1981 config: Default::default(),
1982 lifecycle: Lifecycle::Independent,
1983 }),
1984 )
1985 .await,
1986 )
1987 .await;
1988 assert_eq!(body["active"], v2);
1989 assert_eq!(body["versions"].as_array().unwrap().len(), 2);
1990
1991 let (st, body) = body_json(
1993 rollback_function(
1994 State(deploy.clone()),
1995 axum::extract::Extension(crate::ProjectContext::default()),
1996 Path("greeter".to_string()),
1997 Json(RollbackBody { to: v1.clone() }),
1998 )
1999 .await,
2000 )
2001 .await;
2002 assert_eq!(st, StatusCode::OK);
2003 assert_eq!(body["active"], v1);
2004
2005 let resp = rollback_function(
2007 State(deploy.clone()),
2008 axum::extract::Extension(crate::ProjectContext::default()),
2009 Path("greeter".to_string()),
2010 Json(RollbackBody { to: "c".repeat(64) }),
2011 )
2012 .await;
2013 assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
2014
2015 let (st, body) = body_json(
2017 alias_function(
2018 State(deploy.clone()),
2019 axum::extract::Extension(crate::ProjectContext::default()),
2020 Path(("greeter".to_string(), "prod".to_string())),
2021 Json(AliasBody {
2022 version: v2.clone(),
2023 }),
2024 )
2025 .await,
2026 )
2027 .await;
2028 assert_eq!(st, StatusCode::OK);
2029 assert_eq!(body["aliases"]["prod"], v2);
2030
2031 let (st, _) = body_json(
2033 remove_function(
2034 State(deploy.clone()),
2035 axum::extract::Extension(crate::ProjectContext::default()),
2036 Path("greeter".to_string()),
2037 )
2038 .await,
2039 )
2040 .await;
2041 assert_eq!(st, StatusCode::NO_CONTENT);
2042 assert!(deploy
2043 .get_function(ProjectRef::DEFAULT, "greeter")
2044 .await
2045 .unwrap()
2046 .is_none());
2047
2048 let empty = DeployStore::new(
2050 Arc::new(FakeStorage { present: false }),
2051 Arc::new(MemoryKv::new()),
2052 );
2053 let resp = deploy_function(
2054 State(empty),
2055 axum::extract::Extension(crate::ProjectContext::default()),
2056 axum::extract::Extension(Arc::new(crate::HandlerRuntime::disabled())),
2057 axum::extract::Query(DeployFunctionQuery::default()),
2058 Path("orphan".to_string()),
2059 Json(FunctionUpsert {
2060 component: v1.clone(),
2061 config: Default::default(),
2062 lifecycle: Lifecycle::default(),
2063 }),
2064 )
2065 .await;
2066 assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
2067 }
2068
2069 struct StubControl {
2073 admits: std::sync::Mutex<Vec<(String, String)>>,
2074 respond: StubJoin,
2075 }
2076 #[derive(Clone, Copy)]
2077 enum StubJoin {
2078 Admit,
2079 Spent,
2080 Invalid,
2081 Revoked,
2082 }
2083
2084 #[async_trait::async_trait]
2085 impl MeshControl for StubControl {
2086 async fn admit(
2087 &self,
2088 mesh_pubkey_hex: &str,
2089 jti: &str,
2090 _proof: &[u8],
2091 _proof_iat: u64,
2092 _now: u64,
2093 _advertise_addr: Option<&str>,
2094 ) -> Result<JoinOutcome, String> {
2095 self.admits
2096 .lock()
2097 .unwrap()
2098 .push((mesh_pubkey_hex.to_string(), jti.to_string()));
2099 Ok(match self.respond {
2100 StubJoin::Admit => JoinOutcome::Admitted {
2101 members: vec!["signed-member".to_string()],
2102 addrs: std::collections::BTreeMap::from([(7u64, "https://x:7000".to_string())]),
2103 },
2104 StubJoin::Spent => JoinOutcome::TokenSpent,
2105 StubJoin::Invalid => JoinOutcome::ProofInvalid,
2106 StubJoin::Revoked => JoinOutcome::Revoked,
2107 })
2108 }
2109 async fn rotate_key(&self) -> Result<String, String> {
2110 Ok("cafe".to_string())
2111 }
2112 async fn revoke(&self, _node: u64) -> Result<(), String> {
2113 Ok(())
2114 }
2115 async fn members(&self) -> Result<Vec<MeshMember>, String> {
2116 Ok(Vec::new())
2117 }
2118 async fn promote(&self, _node: u64) -> Result<(), String> {
2119 Ok(())
2120 }
2121 }
2122
2123 #[tokio::test]
2127 async fn cluster_join_dispatches_and_maps_outcomes() {
2128 let keys: Arc<dyn Signer> = Arc::new(LocalSigner::generate(TokenAlg::Es256));
2129 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
2130 let auth = Auth::with_key(keys.public_key(), kv);
2131 let token = cose::mint_join(600, now_unix(), &*keys).await.unwrap();
2132 let req = |proof: &str| JoinRequest {
2133 token: token.clone(),
2134 mesh_pubkey: "302a300506032b6570032100feed".into(),
2135 possession_proof: proof.to_string(),
2136 proof_iat: now_unix(),
2137 advertise_addr: Some("https://joiner:7000".into()),
2138 };
2139
2140 let admitter = Arc::new(StubControl {
2142 admits: std::sync::Mutex::new(Vec::new()),
2143 respond: StubJoin::Admit,
2144 });
2145 let resp = cluster_join(
2146 Extension(auth.clone()),
2147 Extension(MeshControlHandle(Some(admitter.clone()))),
2148 Json(req("aa01")),
2149 )
2150 .await;
2151 assert_eq!(resp.status(), StatusCode::OK);
2152 assert_eq!(admitter.admits.lock().unwrap().len(), 1);
2153
2154 let spent = Arc::new(StubControl {
2156 admits: std::sync::Mutex::new(Vec::new()),
2157 respond: StubJoin::Spent,
2158 });
2159 assert_eq!(
2160 cluster_join(
2161 Extension(auth.clone()),
2162 Extension(MeshControlHandle(Some(spent))),
2163 Json(req("aa01")),
2164 )
2165 .await
2166 .status(),
2167 StatusCode::CONFLICT
2168 );
2169 let invalid = Arc::new(StubControl {
2170 admits: std::sync::Mutex::new(Vec::new()),
2171 respond: StubJoin::Invalid,
2172 });
2173 assert_eq!(
2174 cluster_join(
2175 Extension(auth.clone()),
2176 Extension(MeshControlHandle(Some(invalid))),
2177 Json(req("aa01")),
2178 )
2179 .await
2180 .status(),
2181 StatusCode::FORBIDDEN
2182 );
2183 let revoked = Arc::new(StubControl {
2185 admits: std::sync::Mutex::new(Vec::new()),
2186 respond: StubJoin::Revoked,
2187 });
2188 assert_eq!(
2189 cluster_join(
2190 Extension(auth.clone()),
2191 Extension(MeshControlHandle(Some(revoked))),
2192 Json(req("aa01")),
2193 )
2194 .await
2195 .status(),
2196 StatusCode::FORBIDDEN
2197 );
2198
2199 let ok = Arc::new(StubControl {
2201 admits: std::sync::Mutex::new(Vec::new()),
2202 respond: StubJoin::Admit,
2203 });
2204 assert_eq!(
2205 cluster_join(
2206 Extension(auth.clone()),
2207 Extension(MeshControlHandle(Some(ok))),
2208 Json(req("not-hex")),
2209 )
2210 .await
2211 .status(),
2212 StatusCode::BAD_REQUEST
2213 );
2214
2215 let none = cluster_join(
2217 Extension(auth),
2218 Extension(MeshControlHandle(None)),
2219 Json(req("aa01")),
2220 )
2221 .await;
2222 assert_eq!(none.status(), StatusCode::NOT_IMPLEMENTED);
2223 }
2224
2225 #[tokio::test]
2229 async fn bootstrap_mints_the_first_token_once() {
2230 use axum::http::{header::AUTHORIZATION, HeaderMap, HeaderValue};
2231 let keys: Arc<dyn Signer> = Arc::new(LocalSigner::generate(TokenAlg::Es256));
2232 let public = keys.public_key();
2233 let deploy = DeployStore::new(
2234 Arc::new(MemStorage::default()),
2235 Arc::new(MemoryKv::new()) as Arc<dyn KvStore>,
2236 );
2237 let secret = "s3cr3t-bootstrap-value";
2238 let gate = BootstrapGate::new(Some(secret));
2239 let issuer = Issuer(Some(keys.clone()));
2240 let bearer = |s: &str| {
2241 let mut h = HeaderMap::new();
2242 h.insert(
2243 AUTHORIZATION,
2244 HeaderValue::from_str(&format!("Bearer {s}")).unwrap(),
2245 );
2246 h
2247 };
2248 let req = || BootstrapRequest {
2249 roles: vec!["admin".to_string()],
2250 ttl_secs: None,
2251 };
2252
2253 let bad = bootstrap_token(
2255 State(deploy.clone()),
2256 Extension(issuer.clone()),
2257 Extension(gate.clone()),
2258 bearer("wrong"),
2259 Json(req()),
2260 )
2261 .await;
2262 assert_eq!(bad.status(), StatusCode::UNAUTHORIZED);
2263
2264 let ok = bootstrap_token(
2266 State(deploy.clone()),
2267 Extension(issuer.clone()),
2268 Extension(gate.clone()),
2269 bearer(secret),
2270 Json(req()),
2271 )
2272 .await;
2273 assert_eq!(ok.status(), StatusCode::CREATED);
2274 let body = axum::body::to_bytes(ok.into_body(), usize::MAX)
2275 .await
2276 .unwrap();
2277 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
2278 let token = json["token"].as_str().unwrap();
2279 let id = json["id"].as_str().unwrap();
2280 let verified = cose::verify(token, &public, now_unix()).unwrap();
2281 assert!(verified.roles.iter().any(|r| r.name == "admin"));
2282 assert!(deploy
2283 .list_token_meta()
2284 .await
2285 .unwrap()
2286 .iter()
2287 .any(|m| m.revocation_id == id));
2288
2289 let reuse = bootstrap_token(
2291 State(deploy.clone()),
2292 Extension(issuer.clone()),
2293 Extension(gate),
2294 bearer(secret),
2295 Json(req()),
2296 )
2297 .await;
2298 assert_eq!(reuse.status(), StatusCode::CONFLICT);
2299
2300 let disabled = bootstrap_token(
2302 State(deploy),
2303 Extension(issuer),
2304 Extension(BootstrapGate(None)),
2305 bearer(secret),
2306 Json(req()),
2307 )
2308 .await;
2309 assert_eq!(disabled.status(), StatusCode::NOT_IMPLEMENTED);
2310 }
2311
2312 #[tokio::test]
2315 async fn cluster_rotate_key_returns_the_new_pubkey_or_501() {
2316 let control = Arc::new(StubControl {
2317 admits: std::sync::Mutex::new(Vec::new()),
2318 respond: StubJoin::Admit,
2319 });
2320 let resp = cluster_rotate_key(Extension(MeshControlHandle(Some(control)))).await;
2321 assert_eq!(resp.status(), StatusCode::OK);
2322 let body = axum::body::to_bytes(resp.into_body(), usize::MAX)
2323 .await
2324 .unwrap();
2325 let parsed: serde_json::Value = serde_json::from_slice(&body).unwrap();
2326 assert_eq!(parsed["pubkey"].as_str(), Some("cafe"));
2327
2328 let none = cluster_rotate_key(Extension(MeshControlHandle(None))).await;
2329 assert_eq!(none.status(), StatusCode::NOT_IMPLEMENTED);
2330 }
2331
2332 #[test]
2333 fn gateway_addr_gate_refuses_metadata_and_private_per_posture() {
2334 use boatramp_core::security::SecurityProfile;
2335 let strict = SecurityProfile::MultiTenant.preset();
2336 let loose = SecurityProfile::SingleTenant.preset(); let public: IpAddr = "93.184.216.34".parse().unwrap(); let private: IpAddr = "10.1.2.3".parse().unwrap();
2340 let loopback: IpAddr = "127.0.0.1".parse().unwrap();
2341 let metadata: IpAddr = IpAddr::V4(CLOUD_METADATA_IPV4);
2342
2343 assert!(gateway_addr_allowed(public, &strict));
2345 assert!(!gateway_addr_allowed(private, &strict));
2346 assert!(!gateway_addr_allowed(loopback, &strict));
2347 assert!(!gateway_addr_allowed(metadata, &strict));
2348
2349 assert!(gateway_addr_allowed(public, &loose));
2352 assert!(gateway_addr_allowed(private, &loose));
2353 assert!(gateway_addr_allowed(loopback, &loose));
2354 assert!(!gateway_addr_allowed(metadata, &loose));
2355 }
2356
2357 #[tokio::test]
2358 async fn resolve_env_merges_static_and_host_secrets() {
2359 use boatramp_core::config::HandlersSiteConfig;
2360
2361 std::env::set_var("BOATRAMP_TEST_RESOLVE_SECRET", "topsecret");
2363
2364 let deploy_env = std::collections::BTreeMap::from([
2365 ("GREETING".to_string(), "hi".to_string()),
2366 ("OVERRIDE_ME".to_string(), "static".to_string()),
2367 ]);
2368 let site_handlers = HandlersSiteConfig {
2369 enabled: true,
2370 secrets: std::collections::BTreeMap::from([
2371 (
2373 "SECRET_TOKEN".to_string(),
2374 "BOATRAMP_TEST_RESOLVE_SECRET".to_string(),
2375 ),
2376 (
2377 "OVERRIDE_ME".to_string(),
2378 "BOATRAMP_TEST_RESOLVE_SECRET".to_string(),
2379 ),
2380 (
2381 "MISSING".to_string(),
2382 "BOATRAMP_TEST_NOT_SET_VAR".to_string(),
2383 ),
2384 ]),
2385 ..Default::default()
2386 };
2387 let env = resolve_env(
2390 "blog",
2391 boatramp_core::project::ProjectRef::DEFAULT,
2392 &deploy_env,
2393 &site_handlers,
2394 true,
2395 None,
2396 )
2397 .await
2398 .expect("resolves");
2399
2400 assert!(env.contains(&("GREETING".to_string(), "hi".to_string())));
2404 assert!(env.contains(&("SECRET_TOKEN".to_string(), "topsecret".to_string())));
2405 assert!(env.contains(&("OVERRIDE_ME".to_string(), "topsecret".to_string())));
2406 assert!(!env.iter().any(|(k, _)| k == "MISSING"));
2407
2408 std::env::remove_var("BOATRAMP_TEST_RESOLVE_SECRET");
2409 }
2410
2411 #[tokio::test]
2412 async fn multi_tenant_posture_refuses_a_host_env_handler_secret() {
2413 use boatramp_core::config::HandlersSiteConfig;
2414
2415 std::env::set_var("BOATRAMP_TEST_OTHER_TENANT_SECRET", "leak-me");
2420 let deploy_env = std::collections::BTreeMap::new();
2421 let bare = HandlersSiteConfig {
2422 enabled: true,
2423 secrets: std::collections::BTreeMap::from([(
2424 "STOLEN".to_string(),
2425 "BOATRAMP_TEST_OTHER_TENANT_SECRET".to_string(),
2426 )]),
2427 ..Default::default()
2428 };
2429 let err = resolve_env(
2430 "evil",
2431 boatramp_core::project::ProjectRef::DEFAULT,
2432 &deploy_env,
2433 &bare,
2434 false,
2435 None,
2436 )
2437 .await
2438 .expect_err("multi-tenant must refuse a bare host-env ref");
2439 assert!(
2440 err.contains("STOLEN"),
2441 "error names the offending guest var: {err}"
2442 );
2443 assert!(
2444 err.contains("multi-tenant"),
2445 "error steers the tenant: {err}"
2446 );
2447 assert!(
2448 !err.contains("leak-me"),
2449 "the host value must never appear (never read): {err}"
2450 );
2451
2452 let explicit = HandlersSiteConfig {
2454 enabled: true,
2455 secrets: std::collections::BTreeMap::from([(
2456 "STOLEN".to_string(),
2457 "env:BOATRAMP_TEST_OTHER_TENANT_SECRET".to_string(),
2458 )]),
2459 ..Default::default()
2460 };
2461 assert!(resolve_env(
2462 "evil",
2463 boatramp_core::project::ProjectRef::DEFAULT,
2464 &deploy_env,
2465 &explicit,
2466 false,
2467 None,
2468 )
2469 .await
2470 .is_err());
2471
2472 let reserved = HandlersSiteConfig {
2474 enabled: true,
2475 secrets: std::collections::BTreeMap::from([(
2476 "TOKEN".to_string(),
2477 "vault:kv/data/app#token".to_string(),
2478 )]),
2479 ..Default::default()
2480 };
2481 let err = resolve_env(
2482 "evil",
2483 boatramp_core::project::ProjectRef::DEFAULT,
2484 &deploy_env,
2485 &reserved,
2486 true,
2487 None,
2488 )
2489 .await
2490 .expect_err("a reserved scheme is not yet supported, even under single-tenant");
2491 assert!(err.contains("not yet supported"), "{err}");
2492
2493 let arbitrary = HandlersSiteConfig {
2497 enabled: true,
2498 secrets: std::collections::BTreeMap::from([(
2499 "KEY".to_string(),
2500 "aws:sm/prod/apikey".to_string(),
2501 )]),
2502 ..Default::default()
2503 };
2504 let err = resolve_env(
2505 "evil",
2506 boatramp_core::project::ProjectRef::DEFAULT,
2507 &deploy_env,
2508 &arbitrary,
2509 true,
2510 None,
2511 )
2512 .await
2513 .expect_err("any unknown scheme is reserved, even under single-tenant");
2514 assert!(
2515 err.contains("not yet supported") && err.contains("aws"),
2516 "provider-neutral reservation names the scheme: {err}"
2517 );
2518
2519 std::env::remove_var("BOATRAMP_TEST_OTHER_TENANT_SECRET");
2520 }
2521
2522 #[tokio::test]
2523 async fn function_resolve_secret_env_reads_host_and_matches_handler_semantics() {
2524 std::env::set_var("BOATRAMP_TEST_FN_SECRET", "fnsecret");
2529
2530 let static_env = std::collections::BTreeMap::from([
2531 ("STAGE".to_string(), "prod".to_string()),
2532 ("OVERRIDE_ME".to_string(), "static".to_string()),
2533 ]);
2534 let secrets = std::collections::BTreeMap::from([
2535 ("DB_URL".to_string(), "BOATRAMP_TEST_FN_SECRET".to_string()),
2537 (
2539 "OVERRIDE_ME".to_string(),
2540 "BOATRAMP_TEST_FN_SECRET".to_string(),
2541 ),
2542 (
2544 "MISSING".to_string(),
2545 "BOATRAMP_TEST_FN_NOT_SET".to_string(),
2546 ),
2547 ]);
2548 let env = resolve_secret_env(
2550 "fn/api",
2551 boatramp_core::project::ProjectRef::DEFAULT,
2552 &static_env,
2553 &secrets,
2554 true,
2555 None,
2556 )
2557 .await
2558 .expect("resolves");
2559
2560 assert!(env.contains(&("STAGE".to_string(), "prod".to_string())));
2561 assert!(env.contains(&("DB_URL".to_string(), "fnsecret".to_string())));
2563 assert!(env.contains(&("OVERRIDE_ME".to_string(), "fnsecret".to_string())));
2565 assert!(!env.iter().any(|(k, _)| k == "MISSING"));
2567
2568 std::env::remove_var("BOATRAMP_TEST_FN_SECRET");
2569 }
2570
2571 #[tokio::test]
2572 async fn multi_tenant_posture_refuses_a_host_env_function_secret() {
2573 std::env::set_var("BOATRAMP_TEST_FN_LEAK", "leak-me");
2578 let static_env = std::collections::BTreeMap::new();
2579 let secrets = std::collections::BTreeMap::from([(
2580 "DB_URL".to_string(),
2581 "BOATRAMP_TEST_FN_LEAK".to_string(),
2582 )]);
2583
2584 let err = resolve_secret_env(
2586 "fn/api",
2587 boatramp_core::project::ProjectRef::DEFAULT,
2588 &static_env,
2589 &secrets,
2590 false,
2591 None,
2592 )
2593 .await
2594 .expect_err("multi-tenant must refuse a function host-env ref");
2595 assert!(
2596 err.contains("DB_URL"),
2597 "error names the offending guest var: {err}"
2598 );
2599 assert!(
2600 !err.contains("leak-me"),
2601 "host value must never appear: {err}"
2602 );
2603
2604 let env = resolve_secret_env(
2606 "fn/api",
2607 boatramp_core::project::ProjectRef::DEFAULT,
2608 &static_env,
2609 &secrets,
2610 true,
2611 None,
2612 )
2613 .await
2614 .expect("resolves");
2615 assert!(env.contains(&("DB_URL".to_string(), "leak-me".to_string())));
2616
2617 std::env::remove_var("BOATRAMP_TEST_FN_LEAK");
2618 }
2619
2620 #[tokio::test]
2621 async fn boatramp_scheme_resolves_from_the_project_scoped_store() {
2622 use boatramp_core::project::ProjectRef;
2623 use boatramp_core::secret_store::SecretStore;
2624 use std::sync::Arc;
2625
2626 struct XorEnvelope;
2628 #[async_trait::async_trait]
2629 impl boatramp_core::envelope::KeyEnvelope for XorEnvelope {
2630 async fn wrap(
2631 &self,
2632 p: &[u8],
2633 ) -> Result<Vec<u8>, boatramp_core::envelope::EnvelopeError> {
2634 Ok(p.iter().map(|b| b ^ 0x5a).collect())
2635 }
2636 async fn unwrap(
2637 &self,
2638 c: &[u8],
2639 ) -> Result<Vec<u8>, boatramp_core::envelope::EnvelopeError> {
2640 Ok(c.iter().map(|b| b ^ 0x5a).collect())
2641 }
2642 }
2643
2644 let store = SecretStore::new(
2645 Arc::new(boatramp_core::kv::MemoryKv::new()),
2646 Arc::new(XorEnvelope),
2647 );
2648 store
2649 .set(ProjectRef::new("acme"), "api-key", b"s3cr3t")
2650 .await
2651 .unwrap();
2652
2653 let static_env = std::collections::BTreeMap::new();
2654 let secrets = std::collections::BTreeMap::from([
2655 ("API_KEY".to_string(), "boatramp:api-key".to_string()),
2656 ("MISSING".to_string(), "boatramp:not-set".to_string()),
2657 ]);
2658
2659 let env = resolve_secret_env(
2662 "site",
2663 ProjectRef::new("acme"),
2664 &static_env,
2665 &secrets,
2666 false,
2667 Some(&store),
2668 )
2669 .await
2670 .expect("boatramp refs resolve without the host-env gate");
2671 assert!(env.contains(&("API_KEY".to_string(), "s3cr3t".to_string())));
2672 assert!(!env.iter().any(|(k, _)| k == "MISSING"));
2674
2675 let other_secrets = std::collections::BTreeMap::from([(
2677 "API_KEY".to_string(),
2678 "boatramp:api-key".to_string(),
2679 )]);
2680 let other = resolve_secret_env(
2681 "site",
2682 ProjectRef::new("globex"),
2683 &static_env,
2684 &other_secrets,
2685 false,
2686 Some(&store),
2687 )
2688 .await
2689 .expect("resolves (a foreign project's secret is simply absent → skipped)");
2690 assert!(
2691 !other.iter().any(|(k, _)| k == "API_KEY"),
2692 "a tenant must not read another project's secret"
2693 );
2694
2695 let err = resolve_secret_env(
2697 "site",
2698 ProjectRef::new("acme"),
2699 &static_env,
2700 &other_secrets,
2701 false,
2702 None,
2703 )
2704 .await
2705 .expect_err("no store configured must fail closed");
2706 assert!(err.contains("no internal secret store"), "{err}");
2707 }
2708
2709 fn req() -> Request {
2710 Request::builder()
2711 .uri("/")
2712 .header(header::HOST, "example.com")
2713 .body(Body::empty())
2714 .unwrap()
2715 }
2716
2717 #[test]
2718 fn forwarded_headers_set_standard_triple() {
2719 let mut request = req();
2720 set_forwarded_headers(&mut request, "203.0.113.7".parse().unwrap());
2721 let h = request.headers();
2722 assert_eq!(h.get("x-forwarded-for").unwrap(), "203.0.113.7");
2723 assert_eq!(h.get("x-forwarded-host").unwrap(), "example.com");
2724 assert_eq!(h.get("x-forwarded-proto").unwrap(), "http");
2725 }
2726
2727 #[test]
2728 fn forwarded_for_overwrites_spoofed_value() {
2729 let mut request = Request::builder()
2732 .uri("/")
2733 .header(header::HOST, "example.com")
2734 .header("x-forwarded-for", "10.0.0.1, 1.2.3.4")
2735 .body(Body::empty())
2736 .unwrap();
2737 set_forwarded_headers(&mut request, "203.0.113.7".parse().unwrap());
2738 let values: Vec<_> = request
2739 .headers()
2740 .get_all("x-forwarded-for")
2741 .iter()
2742 .collect();
2743 assert_eq!(values.len(), 1);
2744 assert_eq!(values[0], "203.0.113.7");
2745 }
2746
2747 #[test]
2748 fn forwarded_proto_preserves_upstream_tls() {
2749 let mut request = Request::builder()
2751 .uri("/")
2752 .header(header::HOST, "example.com")
2753 .header("x-forwarded-proto", "https")
2754 .body(Body::empty())
2755 .unwrap();
2756 set_forwarded_headers(&mut request, "203.0.113.7".parse().unwrap());
2757 assert_eq!(request.headers().get("x-forwarded-proto").unwrap(), "https");
2758 }
2759
2760 #[test]
2761 fn forwarded_host_absent_when_no_host_header() {
2762 let mut request = Request::builder().uri("/").body(Body::empty()).unwrap();
2763 set_forwarded_headers(&mut request, "203.0.113.7".parse().unwrap());
2764 assert!(request.headers().get("x-forwarded-host").is_none());
2765 assert_eq!(
2766 request.headers().get("x-forwarded-for").unwrap(),
2767 "203.0.113.7"
2768 );
2769 }
2770
2771 use boatramp_core::kv::{KvStore, MemoryKv};
2774 use boatramp_core::messaging::{LogMessaging, Messaging};
2775 use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, StorageError};
2776
2777 const EVENT_CONSUMER: &[u8] =
2778 include_bytes!("../../boatramp-handlers/tests/fixtures/event-consumer.wasm");
2779
2780 #[derive(Default)]
2781 struct MemStorage {
2782 objects: std::sync::Mutex<std::collections::HashMap<String, Vec<u8>>>,
2783 }
2784
2785 #[async_trait::async_trait]
2786 impl boatramp_core::Storage for MemStorage {
2787 async fn get(&self, key: &str) -> Result<GetObject, StorageError> {
2788 let bytes = self
2789 .objects
2790 .lock()
2791 .unwrap()
2792 .get(key)
2793 .cloned()
2794 .ok_or_else(|| StorageError::NotFound(key.to_string()))?;
2795 let body: ByteStream =
2796 futures::stream::once(async move { Ok(bytes::Bytes::from(bytes)) }).boxed();
2797 Ok(GetObject {
2798 meta: ObjectMeta {
2799 key: key.to_string(),
2800 ..Default::default()
2801 },
2802 body,
2803 })
2804 }
2805 async fn get_range(
2806 &self,
2807 key: &str,
2808 _: u64,
2809 _: Option<u64>,
2810 ) -> Result<GetObject, StorageError> {
2811 self.get(key).await
2812 }
2813 async fn put(
2814 &self,
2815 key: &str,
2816 mut body: ByteStream,
2817 _: PutMeta,
2818 ) -> Result<ObjectMeta, StorageError> {
2819 use futures::StreamExt;
2820 let mut buf = Vec::new();
2821 while let Some(chunk) = body.next().await {
2822 buf.extend_from_slice(&chunk?);
2823 }
2824 self.objects.lock().unwrap().insert(key.to_string(), buf);
2825 Ok(ObjectMeta {
2826 key: key.to_string(),
2827 ..Default::default()
2828 })
2829 }
2830 async fn head(&self, key: &str) -> Result<ObjectMeta, StorageError> {
2831 self.objects
2832 .lock()
2833 .unwrap()
2834 .get(key)
2835 .map(|_| ObjectMeta {
2836 key: key.to_string(),
2837 ..Default::default()
2838 })
2839 .ok_or_else(|| StorageError::NotFound(key.to_string()))
2840 }
2841 async fn delete(&self, key: &str) -> Result<(), StorageError> {
2842 self.objects.lock().unwrap().remove(key);
2843 Ok(())
2844 }
2845 async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
2846 Ok(Vec::new())
2847 }
2848 }
2849
2850 fn observed_state_in(
2853 project: &str,
2854 workload: &str,
2855 host: &str,
2856 healthy: bool,
2857 phase: boatramp_core::compute::ReplicaPhase,
2858 ) -> boatramp_core::compute::ObservedInstance {
2859 use boatramp_core::compute::{Endpoint, InstanceHandle, ReplicaPhase, Scheme, Snapshot};
2860 boatramp_core::compute::ObservedInstance {
2861 handle: InstanceHandle {
2862 project: project.into(),
2863 workload: workload.into(),
2864 replica: 0,
2865 backend_ref: "ref-0".into(),
2866 },
2867 node: 1,
2868 backend: "vmm".into(),
2869 endpoint: Endpoint {
2870 scheme: Scheme::Http,
2871 host: host.into(),
2872 port: 80,
2873 },
2874 region: None,
2875 healthy,
2876 started_at: None,
2877 phase,
2878 snapshot: matches!(phase, ReplicaPhase::Zero).then(|| Snapshot {
2879 project: project.into(),
2880 workload: workload.into(),
2881 replica: 0,
2882 data_ref: "snap-0".into(),
2883 }),
2884 }
2885 }
2886
2887 fn observed_state(
2889 workload: &str,
2890 healthy: bool,
2891 phase: boatramp_core::compute::ReplicaPhase,
2892 ) -> boatramp_core::compute::ObservedInstance {
2893 observed_state_in("default", workload, "10.0.0.2", healthy, phase)
2894 }
2895
2896 #[tokio::test]
2897 async fn has_parked_replica_detects_a_zeroed_replica() {
2898 use boatramp_core::compute::ReplicaPhase;
2899 let storage = Arc::new(MemStorage::default());
2900 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
2901 let deploy = DeployStore::new(storage, kv);
2902
2903 assert!(!has_parked_replica(&deploy, "default", "w").await);
2905 deploy
2907 .set_replica_state(
2908 ProjectRef::DEFAULT,
2909 &observed_state("w", true, ReplicaPhase::Running),
2910 )
2911 .await
2912 .unwrap();
2913 assert!(!has_parked_replica(&deploy, "default", "w").await);
2914 deploy
2916 .set_replica_state(
2917 ProjectRef::DEFAULT,
2918 &observed_state("w", false, ReplicaPhase::Zero),
2919 )
2920 .await
2921 .unwrap();
2922 assert!(has_parked_replica(&deploy, "default", "w").await);
2923 }
2924
2925 #[tokio::test]
2926 async fn await_warm_returns_immediately_when_healthy_and_times_out_otherwise() {
2927 use boatramp_core::compute::ReplicaPhase;
2928 let storage = Arc::new(MemStorage::default());
2929 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
2930 let deploy = DeployStore::new(storage, kv);
2931
2932 let empty = await_warm(
2934 &deploy,
2935 "default",
2936 "w",
2937 std::time::Duration::from_millis(150),
2938 )
2939 .await;
2940 assert!(empty.is_empty());
2941
2942 deploy
2944 .set_replica_state(
2945 ProjectRef::DEFAULT,
2946 &observed_state("w", true, ReplicaPhase::Running),
2947 )
2948 .await
2949 .unwrap();
2950 let warm = await_warm(&deploy, "default", "w", std::time::Duration::from_secs(5)).await;
2951 assert_eq!(warm, vec!["http://10.0.0.2:80".to_string()]);
2952 }
2953
2954 #[tokio::test]
2960 async fn compute_endpoints_are_project_scoped() {
2961 use boatramp_core::compute::ReplicaPhase;
2962 let storage = Arc::new(MemStorage::default());
2963 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
2964 let deploy = DeployStore::new(storage, kv);
2965
2966 deploy
2968 .set_replica_state(
2969 ProjectRef::new("acme"),
2970 &observed_state_in("acme", "web", "10.0.0.5", true, ReplicaPhase::Running),
2971 )
2972 .await
2973 .unwrap();
2974 deploy
2975 .set_replica_state(
2976 ProjectRef::DEFAULT,
2977 &observed_state_in("default", "web", "10.0.0.9", true, ReplicaPhase::Running),
2978 )
2979 .await
2980 .unwrap();
2981
2982 assert_eq!(
2984 compute_endpoints(&deploy, "acme", "web").await,
2985 vec!["http://10.0.0.5:80".to_string()],
2986 "acme's web resolves against acme, not default"
2987 );
2988 assert_eq!(
2990 compute_endpoints(&deploy, "default", "web").await,
2991 vec!["http://10.0.0.9:80".to_string()]
2992 );
2993 assert!(compute_endpoints(&deploy, "beta", "web").await.is_empty());
2996 }
2997
2998 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3002 async fn dispatcher_delivers_at_least_once_then_dead_letters() {
3003 use boatramp_handlers::{Bindings, HandlerEngine, Limits};
3004 let storage = Arc::new(MemStorage::default());
3005 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3006 let mq = LogMessaging::new(storage, kv.clone());
3007 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3008 let hash = boatramp_core::deploy::sha256_hex(EVENT_CONSUMER);
3009 let bindings = Bindings::new("blog").with_keyvalue("blog", kv.clone());
3010 let topic = "blog/orders/created";
3011
3012 for _ in 0..3 {
3014 mq.publish(topic, b"ok").await.unwrap();
3015 }
3016 loop {
3017 let acked = dispatch_consumer_batch(
3018 &engine,
3019 &mq,
3020 &metrics::Metrics::default(),
3021 "blog",
3022 topic,
3023 "blog/",
3024 "",
3025 boatramp_core::messaging::StartPosition::Latest,
3026 &hash,
3027 EVENT_CONSUMER,
3028 &bindings,
3029 Limits::default(),
3030 Duration::from_secs(30),
3031 5,
3032 10,
3033 )
3034 .await;
3035 if acked == 0 {
3036 break;
3037 }
3038 }
3039 assert_eq!(
3040 kv.get("hkv/blog/delivered/orders/created").await.unwrap(),
3041 Some(b"3".to_vec())
3042 );
3043
3044 mq.publish(topic, b"fail").await.unwrap();
3047 for _ in 0..5 {
3048 dispatch_consumer_batch(
3049 &engine,
3050 &mq,
3051 &metrics::Metrics::default(),
3052 "blog",
3053 topic,
3054 "blog/",
3055 "",
3056 boatramp_core::messaging::StartPosition::Latest,
3057 &hash,
3058 EVENT_CONSUMER,
3059 &bindings,
3060 Limits::default(),
3061 Duration::ZERO,
3062 2,
3063 10,
3064 )
3065 .await;
3066 }
3067 assert_eq!(mq.dead_letter_count(topic).await.unwrap(), 1);
3068 assert_eq!(
3070 kv.get("hkv/blog/delivered/orders/created").await.unwrap(),
3071 Some(b"3".to_vec())
3072 );
3073 }
3074
3075 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3080 async fn consumer_groups_fan_out_through_the_dispatcher() {
3081 use boatramp_handlers::{Bindings, HandlerEngine, Limits};
3082 let storage = Arc::new(MemStorage::default());
3083 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3084 let mq = LogMessaging::new(storage, kv.clone());
3085 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3086 let hash = boatramp_core::deploy::sha256_hex(EVENT_CONSUMER);
3087 let bindings = Bindings::new("blog").with_keyvalue("blog", kv.clone());
3088 let topic = "blog/orders/created";
3089 let start = boatramp_core::messaging::StartPosition::Latest;
3090
3091 for g in ["billing", "audit"] {
3094 let n = dispatch_consumer_batch(
3095 &engine,
3096 &mq,
3097 &metrics::Metrics::default(),
3098 "blog",
3099 topic,
3100 "blog/",
3101 g,
3102 start,
3103 &hash,
3104 EVENT_CONSUMER,
3105 &bindings,
3106 Limits::default(),
3107 Duration::from_secs(30),
3108 5,
3109 10,
3110 )
3111 .await;
3112 assert_eq!(n, 0, "no events yet for group {g}");
3113 }
3114 mq.publish(topic, b"ok").await.unwrap();
3115
3116 for g in ["billing", "audit"] {
3118 let n = dispatch_consumer_batch(
3119 &engine,
3120 &mq,
3121 &metrics::Metrics::default(),
3122 "blog",
3123 topic,
3124 "blog/",
3125 g,
3126 start,
3127 &hash,
3128 EVENT_CONSUMER,
3129 &bindings,
3130 Limits::default(),
3131 Duration::from_secs(30),
3132 5,
3133 10,
3134 )
3135 .await;
3136 assert_eq!(n, 1, "group {g} should receive the message");
3137 }
3138 assert_eq!(
3140 kv.get("hkv/blog/delivered/orders/created").await.unwrap(),
3141 Some(b"2".to_vec())
3142 );
3143 }
3144
3145 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3149 async fn scheduler_runs_current_consumers_not_previews() {
3150 use boatramp_core::config::{ConsumerConfig, DeployConfig, HandlersSiteConfig, SiteConfig};
3151 use boatramp_core::deploy::{DeployStore, FileEntry, Manifest};
3152 use boatramp_handlers::{HandlerEngine, Limits};
3153 use futures::StreamExt;
3154
3155 let storage = Arc::new(MemStorage::default());
3156 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3157 let deploy = DeployStore::new(storage.clone(), kv.clone());
3158 let messaging: Arc<dyn Messaging> =
3159 Arc::new(LogMessaging::new(storage.clone(), kv.clone()));
3160
3161 let hash = boatramp_core::deploy::sha256_hex(EVENT_CONSUMER);
3163 let stream: ByteStream =
3164 futures::stream::once(async move { Ok(bytes::Bytes::from_static(EVENT_CONSUMER)) })
3165 .boxed();
3166 deploy.put_blob(&hash, stream).await.unwrap();
3167 let mut files = std::collections::BTreeMap::new();
3168 files.insert(
3169 "consumer.wasm".to_string(),
3170 FileEntry {
3171 hash: hash.clone(),
3172 size: EVENT_CONSUMER.len() as u64,
3173 content_type: None,
3174 variants: std::collections::BTreeMap::new(),
3175 },
3176 );
3177 let manifest = Manifest {
3178 files,
3179 config: DeployConfig {
3180 consumers: vec![ConsumerConfig {
3181 topic: "orders/created".into(),
3182 component: "consumer.wasm".into(),
3183 imports: vec!["wasi:keyvalue".into()],
3184 group: String::new(),
3185 start: Default::default(),
3186 }],
3187 ..Default::default()
3188 },
3189 ..Default::default()
3190 };
3191 let id = deploy.put_manifest(&manifest).await.unwrap();
3192 deploy
3193 .activate(ProjectRef::DEFAULT, "blog", &id)
3194 .await
3195 .unwrap();
3196 deploy
3197 .set_site_config(
3198 ProjectRef::DEFAULT,
3199 "blog",
3200 &SiteConfig {
3201 handlers: Some(HandlersSiteConfig {
3202 enabled: true,
3203 allow_imports: vec!["wasi:keyvalue".into()],
3204 ..Default::default()
3205 }),
3206 ..Default::default()
3207 },
3208 )
3209 .await
3210 .unwrap();
3211
3212 messaging
3214 .publish("blog/orders/created", b"live")
3215 .await
3216 .unwrap();
3217 messaging
3218 .publish("blog/_preview/abc/orders/created", b"preview")
3219 .await
3220 .unwrap();
3221
3222 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3223 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, Some(messaging));
3224 let inner = rt.inner.clone().unwrap();
3225 let mut cache = std::collections::HashMap::new();
3226 let mut crons = std::collections::HashMap::new();
3227 let mut sweep = std::collections::HashMap::new();
3228 let now = CronNow {
3229 minute: 0,
3230 hour: 0,
3231 dom: 1,
3232 month: 1,
3233 dow: 0,
3234 minute_stamp: 0,
3235 };
3236 for _ in 0..3 {
3237 run_scheduler_tick(&inner, &deploy, &mut cache, &mut crons, &mut sweep, now)
3238 .await
3239 .unwrap();
3240 }
3241
3242 assert_eq!(
3244 kv.get("hkv/blog/delivered/orders/created").await.unwrap(),
3245 Some(b"1".to_vec())
3246 );
3247 assert_eq!(
3250 kv.get("hkv/blog/_preview/abc/delivered/orders/created")
3251 .await
3252 .unwrap(),
3253 None
3254 );
3255 }
3256
3257 const KV_COUNTER: &[u8] =
3262 include_bytes!("../../boatramp-handlers/tests/fixtures/kv-counter.wasm");
3263
3264 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3269 async fn function_invoker_runs_target_buffers_and_meters() {
3270 use boatramp_core::deploy::DeployStore;
3271 use boatramp_core::function::{Function, FunctionVersion, Lifecycle, Owner};
3272 use boatramp_handlers::{HandlerEngine, InvokeError, InvokeRequest, Invoker, Limits};
3273 use futures::StreamExt;
3274
3275 const HTTP_200: &[u8] =
3278 include_bytes!("../../boatramp-handlers/tests/fixtures/http-200.wasm");
3279
3280 let storage = Arc::new(MemStorage::default());
3281 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3282 let deploy = DeployStore::new(storage.clone(), kv.clone());
3283
3284 let hash = boatramp_core::deploy::sha256_hex(HTTP_200);
3285 let stream: ByteStream =
3286 futures::stream::once(async move { Ok(bytes::Bytes::from_static(HTTP_200)) }).boxed();
3287 deploy.put_blob(&hash, stream).await.unwrap();
3288 let function = Function {
3289 name: "target".into(),
3290 owner: Owner::Project("default".into()),
3291 versions: vec![FunctionVersion {
3292 id: "v1".into(),
3293 component: hash.clone(),
3294 created: 0,
3295 lifecycle: Lifecycle::Independent,
3296 }],
3297 active: "v1".into(),
3298 aliases: Default::default(),
3299 config: Default::default(),
3300 };
3301 deploy
3302 .put_function(ProjectRef::DEFAULT, &function)
3303 .await
3304 .unwrap();
3305
3306 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3307 let rt = HandlerRuntime::new(engine, kv, storage, None, None);
3308 rt.set_invoker(deploy.clone());
3309 let invoker = rt.inner.as_ref().unwrap().invoker.get().unwrap().clone();
3310
3311 let request = || InvokeRequest {
3312 method: "GET".into(),
3313 path: "/".into(),
3314 headers: vec![],
3315 body: vec![],
3316 };
3317
3318 let response = invoker.invoke("target", request(), 1).await.unwrap();
3320 assert_eq!(response.status, 200);
3321
3322 let metering = deploy
3324 .get_metering(ProjectRef::DEFAULT, "target")
3325 .await
3326 .unwrap()
3327 .unwrap();
3328 assert_eq!(metering.invocations, 1);
3329
3330 let err = invoker.invoke("ghost", request(), 1).await.unwrap_err();
3332 assert!(matches!(err, InvokeError::NotFound));
3333 }
3334
3335 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3341 async fn federation_runner_enforces_the_safelist_before_planning() {
3342 use boatramp_core::deploy::DeployStore;
3343 use boatramp_core::project::ProjectRef;
3344 use boatramp_handlers::{GraphqlRequest, HandlerEngine, Limits, SupergraphRunError};
3345
3346 let storage = Arc::new(MemStorage::default());
3347 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3348 let deploy = DeployStore::new(storage.clone(), kv.clone());
3349 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3350 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
3351 rt.set_invoker(deploy.clone());
3352 let runner = rt
3353 .inner
3354 .as_ref()
3355 .unwrap()
3356 .federation_runner
3357 .get()
3358 .unwrap()
3359 .scoped(ProjectRef::new("default"));
3360
3361 let req = |query: &str| GraphqlRequest {
3362 query: Some(query.to_string()),
3363 persisted_hash: None,
3364 variables: "{}".to_string(),
3365 operation_name: None,
3366 authorization: None,
3367 };
3368
3369 assert!(matches!(
3371 runner.run(req("{ me { id } }"), 1).await,
3372 Err(SupergraphRunError::NotSafelisted)
3373 ));
3374
3375 let query = "{ me { id } }";
3378 let hash = crate::graphql_apq::sha256_hex(query);
3379 kv.put(&format!("hapq/default/{hash}"), query.as_bytes().to_vec())
3380 .await
3381 .unwrap();
3382 assert!(matches!(
3383 runner.run(req(query), 1).await,
3384 Err(SupergraphRunError::PlanFailed(_))
3385 ));
3386
3387 let persisted = GraphqlRequest {
3389 query: None,
3390 persisted_hash: Some("deadbeef".to_string()),
3391 variables: "{}".to_string(),
3392 operation_name: None,
3393 authorization: None,
3394 };
3395 assert!(matches!(
3396 runner.run(persisted, 1).await,
3397 Err(SupergraphRunError::NotSafelisted)
3398 ));
3399 }
3400
3401 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3407 async fn scheduler_drains_a_non_default_projects_invocation_in_its_own_tenant() {
3408 use crate::scheduler::{run_scheduler_tick, CronNow};
3409 use boatramp_core::deploy::DeployStore;
3410 use boatramp_core::function::{
3411 Function, FunctionVersion, Invocation, InvocationStatus, InvokeMode, Lifecycle, Owner,
3412 };
3413 use boatramp_handlers::{HandlerEngine, Limits};
3414 use futures::StreamExt;
3415
3416 const HTTP_200: &[u8] =
3417 include_bytes!("../../boatramp-handlers/tests/fixtures/http-200.wasm");
3418
3419 let storage = Arc::new(MemStorage::default());
3420 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3421 let deploy = DeployStore::new(storage.clone(), kv.clone());
3422
3423 let hash = boatramp_core::deploy::sha256_hex(HTTP_200);
3424 let stream: ByteStream =
3425 futures::stream::once(async move { Ok(bytes::Bytes::from_static(HTTP_200)) }).boxed();
3426 deploy.put_blob(&hash, stream).await.unwrap();
3427
3428 let acme = ProjectRef::new("acme");
3430 let function = Function {
3431 name: "worker".into(),
3432 owner: Owner::Project("acme".into()),
3433 versions: vec![FunctionVersion {
3434 id: "v1".into(),
3435 component: hash.clone(),
3436 created: 0,
3437 lifecycle: Lifecycle::Independent,
3438 }],
3439 active: "v1".into(),
3440 aliases: Default::default(),
3441 config: Default::default(),
3442 };
3443 deploy.put_function(acme, &function).await.unwrap();
3444 let inv = Invocation {
3445 id: "inv1".into(),
3446 function: "worker".into(),
3447 version: "v1".into(),
3448 mode: InvokeMode::Async,
3449 status: InvocationStatus::Queued,
3450 idempotency_key: None,
3451 attempts: 0,
3452 lease_expires: None,
3453 request_b64: None,
3454 request_content_type: None,
3455 result: None,
3456 created: 0,
3457 updated: 0,
3458 };
3459 deploy.put_invocation(acme, &inv).await.unwrap();
3460
3461 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3462 let rt = HandlerRuntime::new(engine, kv.clone(), storage.clone(), None, None);
3463 let inner = rt.inner.as_ref().unwrap();
3464
3465 let mut wasm_cache = std::collections::HashMap::new();
3468 let mut cron_state = std::collections::HashMap::new();
3469 let mut sweep = std::collections::HashMap::new();
3470 let now = CronNow {
3471 minute: 0,
3472 hour: 0,
3473 dom: 1,
3474 month: 1,
3475 dow: 0,
3476 minute_stamp: 0,
3477 };
3478 run_scheduler_tick(
3479 inner,
3480 &deploy,
3481 &mut wasm_cache,
3482 &mut cron_state,
3483 &mut sweep,
3484 now,
3485 )
3486 .await
3487 .unwrap();
3488
3489 let settled = poll_invocation_settled(&deploy, acme, "worker", "inv1").await;
3492 assert_eq!(settled.status, InvocationStatus::Succeeded);
3494 let metering = deploy.get_metering(acme, "worker").await.unwrap().unwrap();
3496 assert_eq!(metering.invocations, 1);
3497 assert!(deploy
3499 .get_invocation(ProjectRef::DEFAULT, "worker", "inv1")
3500 .await
3501 .unwrap()
3502 .is_none());
3503 assert!(deploy
3504 .get_metering(ProjectRef::DEFAULT, "worker")
3505 .await
3506 .unwrap()
3507 .is_none());
3508 }
3509
3510 #[cfg(feature = "handlers")]
3514 async fn poll_invocation_settled(
3515 deploy: &boatramp_core::deploy::DeployStore,
3516 project: ProjectRef<'_>,
3517 function: &str,
3518 id: &str,
3519 ) -> boatramp_core::function::Invocation {
3520 use boatramp_core::function::InvocationStatus;
3521 for _ in 0..200 {
3522 if let Some(inv) = deploy.get_invocation(project, function, id).await.unwrap() {
3523 if matches!(
3524 inv.status,
3525 InvocationStatus::Succeeded | InvocationStatus::Failed
3526 ) {
3527 return inv;
3528 }
3529 }
3530 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
3531 }
3532 panic!("invocation {function}/{id} never settled");
3533 }
3534
3535 #[cfg(feature = "handlers")]
3540 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3541 async fn drain_reclaims_an_expired_lease_and_skips_a_live_one() {
3542 use crate::scheduler::{run_scheduler_tick, CronNow};
3543 use boatramp_core::deploy::DeployStore;
3544 use boatramp_core::function::{
3545 Function, FunctionVersion, Invocation, InvocationStatus, InvokeMode, Lifecycle, Owner,
3546 };
3547 use boatramp_handlers::{HandlerEngine, Limits};
3548 use futures::StreamExt;
3549
3550 const HTTP_200: &[u8] =
3551 include_bytes!("../../boatramp-handlers/tests/fixtures/http-200.wasm");
3552
3553 let storage = Arc::new(MemStorage::default());
3554 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3555 let deploy = DeployStore::new(storage.clone(), kv.clone());
3556 let hash = boatramp_core::deploy::sha256_hex(HTTP_200);
3557 let stream: ByteStream =
3558 futures::stream::once(async move { Ok(bytes::Bytes::from_static(HTTP_200)) }).boxed();
3559 deploy.put_blob(&hash, stream).await.unwrap();
3560
3561 let function = Function {
3562 name: "worker".into(),
3563 owner: Owner::Project("default".into()),
3564 versions: vec![FunctionVersion {
3565 id: "v1".into(),
3566 component: hash.clone(),
3567 created: 0,
3568 lifecycle: Lifecycle::Independent,
3569 }],
3570 active: "v1".into(),
3571 aliases: Default::default(),
3572 config: Default::default(),
3573 };
3574 deploy
3575 .put_function(ProjectRef::DEFAULT, &function)
3576 .await
3577 .unwrap();
3578
3579 let base = Invocation {
3582 id: String::new(),
3583 function: "worker".into(),
3584 version: "v1".into(),
3585 mode: InvokeMode::Async,
3586 status: InvocationStatus::Running,
3587 idempotency_key: None,
3588 attempts: 1,
3589 lease_expires: None,
3590 request_b64: None,
3591 request_content_type: None,
3592 result: None,
3593 created: 0,
3594 updated: 0,
3595 };
3596 let orphan = Invocation {
3597 id: "orphan".into(),
3598 lease_expires: Some(1),
3599 ..base.clone()
3600 };
3601 deploy
3602 .put_invocation(ProjectRef::DEFAULT, &orphan)
3603 .await
3604 .unwrap();
3605 let live = Invocation {
3606 id: "live".into(),
3607 lease_expires: Some(u64::MAX),
3608 ..base.clone()
3609 };
3610 deploy
3611 .put_invocation(ProjectRef::DEFAULT, &live)
3612 .await
3613 .unwrap();
3614
3615 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3616 let rt = HandlerRuntime::new(engine, kv.clone(), storage.clone(), None, None);
3617 let inner = rt.inner.as_ref().unwrap();
3618
3619 let now = CronNow {
3620 minute: 0,
3621 hour: 0,
3622 dom: 1,
3623 month: 1,
3624 dow: 0,
3625 minute_stamp: 0,
3626 };
3627 let mut wasm_cache = std::collections::HashMap::new();
3628 let mut cron_state = std::collections::HashMap::new();
3629 let mut sweep = std::collections::HashMap::new();
3630 run_scheduler_tick(
3631 inner,
3632 &deploy,
3633 &mut wasm_cache,
3634 &mut cron_state,
3635 &mut sweep,
3636 now,
3637 )
3638 .await
3639 .unwrap();
3640
3641 let settled =
3643 poll_invocation_settled(&deploy, ProjectRef::DEFAULT, "worker", "orphan").await;
3644 assert_eq!(settled.status, InvocationStatus::Succeeded);
3645 assert_eq!(settled.attempts, 2, "a reclaim counts as another attempt");
3646 assert_eq!(
3647 settled.lease_expires, None,
3648 "a settled invocation drops its lease"
3649 );
3650 let live_after = deploy
3652 .get_invocation(ProjectRef::DEFAULT, "worker", "live")
3653 .await
3654 .unwrap()
3655 .unwrap();
3656 assert_eq!(live_after.status, InvocationStatus::Running);
3657 assert_eq!(live_after.attempts, 1, "a live lease is never reclaimed");
3658 assert_eq!(live_after.lease_expires, Some(u64::MAX));
3659 }
3660
3661 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3676 async fn guest_kv_is_isolated_between_same_named_functions_in_two_projects() {
3677 use boatramp_core::deploy::DeployStore;
3678 use boatramp_core::function::{
3679 Function, FunctionConfig, FunctionVersion, Lifecycle, Owner,
3680 };
3681 use boatramp_handlers::{HandlerEngine, Limits};
3682 use futures::StreamExt;
3683
3684 let storage = Arc::new(MemStorage::default());
3685 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3686 let deploy = DeployStore::new(storage.clone(), kv.clone());
3687
3688 let hash = boatramp_core::deploy::sha256_hex(KV_COUNTER);
3691 let stream: ByteStream =
3692 futures::stream::once(async move { Ok(bytes::Bytes::from_static(KV_COUNTER)) }).boxed();
3693 deploy.put_blob(&hash, stream).await.unwrap();
3694
3695 let store = Function {
3700 name: "store".into(),
3701 owner: Owner::Project("default".into()),
3702 versions: vec![FunctionVersion {
3703 id: "v1".into(),
3704 component: hash.clone(),
3705 created: 0,
3706 lifecycle: Lifecycle::Independent,
3707 }],
3708 active: "v1".into(),
3709 aliases: Default::default(),
3710 config: FunctionConfig {
3711 imports: vec!["wasi:keyvalue".into()],
3712 ..Default::default()
3713 },
3714 };
3715 let acme = ProjectRef::new("acme");
3716 let globex = ProjectRef::new("globex");
3717
3718 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3719 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
3720 let inner = rt.inner.as_ref().unwrap();
3721
3722 let request = || {
3723 axum::http::Request::builder()
3724 .method("GET")
3725 .uri("/")
3726 .body(axum::body::Body::empty())
3727 .unwrap()
3728 };
3729
3730 let component = store.resolve(&store.active).unwrap().to_owned();
3733 for project in [acme, globex, ProjectRef::DEFAULT] {
3734 let (response, _) = execute_function(
3735 inner,
3736 &deploy,
3737 project,
3738 &store,
3739 &component,
3740 request(),
3741 0,
3742 boatramp_handlers::Lane::Sync,
3743 crate::function_runtime::FnTenant::Request,
3744 )
3745 .await;
3746 assert!(response.status().is_success(), "invocation should succeed");
3747 }
3748
3749 assert_eq!(
3752 kv.get("hkv/acme/fn/store/hits").await.unwrap(),
3753 Some(b"1".to_vec()),
3754 "acme's write must be tenant-qualified"
3755 );
3756 assert_eq!(
3757 kv.get("hkv/globex/fn/store/hits").await.unwrap(),
3758 Some(b"1".to_vec()),
3759 "globex's write must be tenant-qualified"
3760 );
3761 assert_eq!(
3762 kv.get("hkv/fn/store/hits").await.unwrap(),
3763 Some(b"1".to_vec()),
3764 "the default project must keep the byte-identical pre-project key"
3765 );
3766 }
3769
3770 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3774 async fn scheduler_fires_crons_with_dedup_and_overlap_skip() {
3775 use boatramp_core::config::{
3776 CronConfig, DeployConfig, HandlerConfig, HandlersSiteConfig, Overlap, SiteConfig,
3777 };
3778 use boatramp_core::deploy::{DeployStore, FileEntry, Manifest};
3779 use boatramp_handlers::{HandlerEngine, Limits};
3780 use futures::StreamExt;
3781 use std::sync::atomic::Ordering;
3782
3783 let storage = Arc::new(MemStorage::default());
3784 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3785 let deploy = DeployStore::new(storage.clone(), kv.clone());
3786
3787 let hash = boatramp_core::deploy::sha256_hex(KV_COUNTER);
3788 let stream: ByteStream =
3789 futures::stream::once(async move { Ok(bytes::Bytes::from_static(KV_COUNTER)) }).boxed();
3790 deploy.put_blob(&hash, stream).await.unwrap();
3791 let mut files = std::collections::BTreeMap::new();
3792 files.insert(
3793 "counter.wasm".to_string(),
3794 FileEntry {
3795 hash: hash.clone(),
3796 size: KV_COUNTER.len() as u64,
3797 content_type: None,
3798 variants: std::collections::BTreeMap::new(),
3799 },
3800 );
3801 let manifest = Manifest {
3802 files,
3803 config: DeployConfig {
3804 handlers: vec![HandlerConfig {
3805 route: "/".into(),
3806 methods: Vec::new(),
3807 component: "counter.wasm".into(),
3808 imports: vec!["wasi:keyvalue".into()],
3809 streaming: false,
3810 limits: None,
3811 env: std::collections::BTreeMap::new(),
3812 invoke_targets: Vec::new(),
3813 }],
3814 crons: vec![CronConfig {
3815 schedule: "* * * * *".into(),
3816 route: "/".into(),
3817 overlap: Overlap::Skip,
3818 }],
3819 ..Default::default()
3820 },
3821 ..Default::default()
3822 };
3823 let id = deploy.put_manifest(&manifest).await.unwrap();
3824 deploy
3825 .activate(ProjectRef::DEFAULT, "blog", &id)
3826 .await
3827 .unwrap();
3828 deploy
3829 .set_site_config(
3830 ProjectRef::DEFAULT,
3831 "blog",
3832 &SiteConfig {
3833 handlers: Some(HandlersSiteConfig {
3834 enabled: true,
3835 allow_imports: vec!["wasi:keyvalue".into()],
3836 ..Default::default()
3837 }),
3838 ..Default::default()
3839 },
3840 )
3841 .await
3842 .unwrap();
3843
3844 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3845 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
3846 let inner = rt.inner.clone().unwrap();
3847 let mut wasm = std::collections::HashMap::new();
3848 let mut crons = std::collections::HashMap::new();
3849 let mut sweep = std::collections::HashMap::new();
3850 let at = |stamp| CronNow {
3851 minute: 0,
3852 hour: 0,
3853 dom: 1,
3854 month: 1,
3855 dow: 0,
3856 minute_stamp: stamp,
3857 };
3858
3859 let (_, handles) =
3861 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, at(100))
3862 .await
3863 .unwrap();
3864 for h in handles {
3865 h.await.unwrap();
3866 }
3867 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), Some(b"1".to_vec()));
3868
3869 let (_, handles) =
3871 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, at(100))
3872 .await
3873 .unwrap();
3874 assert!(handles.is_empty());
3875 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), Some(b"1".to_vec()));
3876
3877 let (_, handles) =
3879 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, at(101))
3880 .await
3881 .unwrap();
3882 for h in handles {
3883 h.await.unwrap();
3884 }
3885 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), Some(b"2".to_vec()));
3886
3887 crons
3891 .get("default|blog|cron|0")
3892 .unwrap()
3893 .running
3894 .store(true, Ordering::Release);
3895 let (_, handles) =
3896 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, at(102))
3897 .await
3898 .unwrap();
3899 assert!(handles.is_empty());
3900 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), Some(b"2".to_vec()));
3901 }
3902
3903 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3907 async fn cron_leader_gate_suppresses_crons_off_leader() {
3908 use boatramp_core::config::{
3909 CronConfig, DeployConfig, HandlerConfig, HandlersSiteConfig, Overlap, SiteConfig,
3910 };
3911 use boatramp_core::deploy::{DeployStore, FileEntry, Manifest};
3912 use boatramp_handlers::{HandlerEngine, Limits};
3913 use futures::StreamExt;
3914
3915 let storage = Arc::new(MemStorage::default());
3916 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3917 let deploy = DeployStore::new(storage.clone(), kv.clone());
3918
3919 let hash = boatramp_core::deploy::sha256_hex(KV_COUNTER);
3920 let stream: ByteStream =
3921 futures::stream::once(async move { Ok(bytes::Bytes::from_static(KV_COUNTER)) }).boxed();
3922 deploy.put_blob(&hash, stream).await.unwrap();
3923 let mut files = std::collections::BTreeMap::new();
3924 files.insert(
3925 "counter.wasm".to_string(),
3926 FileEntry {
3927 hash: hash.clone(),
3928 size: KV_COUNTER.len() as u64,
3929 content_type: None,
3930 variants: std::collections::BTreeMap::new(),
3931 },
3932 );
3933 let manifest = Manifest {
3934 files,
3935 config: DeployConfig {
3936 handlers: vec![HandlerConfig {
3937 route: "/".into(),
3938 methods: Vec::new(),
3939 component: "counter.wasm".into(),
3940 imports: vec!["wasi:keyvalue".into()],
3941 streaming: false,
3942 limits: None,
3943 env: std::collections::BTreeMap::new(),
3944 invoke_targets: Vec::new(),
3945 }],
3946 crons: vec![CronConfig {
3947 schedule: "* * * * *".into(),
3948 route: "/".into(),
3949 overlap: Overlap::Skip,
3950 }],
3951 ..Default::default()
3952 },
3953 ..Default::default()
3954 };
3955 let id = deploy.put_manifest(&manifest).await.unwrap();
3956 deploy
3957 .activate(ProjectRef::DEFAULT, "blog", &id)
3958 .await
3959 .unwrap();
3960 deploy
3961 .set_site_config(
3962 ProjectRef::DEFAULT,
3963 "blog",
3964 &SiteConfig {
3965 handlers: Some(HandlersSiteConfig {
3966 enabled: true,
3967 allow_imports: vec!["wasi:keyvalue".into()],
3968 ..Default::default()
3969 }),
3970 ..Default::default()
3971 },
3972 )
3973 .await
3974 .unwrap();
3975
3976 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3977 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
3978 rt.set_cron_leader_gate(Arc::new(|| false));
3980 let inner = rt.inner.clone().unwrap();
3981 let mut wasm = std::collections::HashMap::new();
3982 let mut crons = std::collections::HashMap::new();
3983 let mut sweep = std::collections::HashMap::new();
3984 let now = CronNow {
3985 minute: 0,
3986 hour: 0,
3987 dom: 1,
3988 month: 1,
3989 dow: 0,
3990 minute_stamp: 100,
3991 };
3992
3993 let (_, handles) =
3994 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, now)
3995 .await
3996 .unwrap();
3997 assert!(handles.is_empty(), "a non-leader must not fire crons");
3999 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), None);
4000 }
4001
4002 #[tokio::test]
4009 async fn build_bindings_dispatches_named_sql_databases_with_least_privilege() {
4010 use boatramp_core::config::HandlersSiteConfig;
4011 use boatramp_core::project::ProjectRef;
4012 use boatramp_handlers::{HandlerEngine, Limits};
4013
4014 let kv: Arc<dyn boatramp_core::kv::KvStore> = Arc::new(boatramp_core::kv::MemoryKv::new());
4015 let storage: Arc<dyn boatramp_core::Storage> = Arc::new(MemStorage::default());
4016 let sql_dir =
4018 std::env::temp_dir().join(format!("boatramp-named-sql-{}", std::process::id()));
4019 let _ = std::fs::remove_dir_all(&sql_dir);
4020 let sql: Arc<dyn boatramp_core::sql::SqlBackends> =
4021 Arc::new(boatramp_storage::LibsqlSqlBackends::local(&sql_dir));
4022
4023 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4024 let rt = HandlerRuntime::new(engine, kv, storage, Some(sql), None);
4025 let inner = rt.inner.as_ref().unwrap();
4026
4027 let site = HandlersSiteConfig {
4031 enabled: true,
4032 allow_imports: vec!["sql".into(), "sql:product".into(), "sql:privileged".into()],
4033 tenancy: Some(boatramp_core::tenancy::Tenancy::Disabled),
4034 ..Default::default()
4035 };
4036 let env = std::collections::BTreeMap::new();
4037 let build = |imports: &[&str]| {
4038 let imports: Vec<String> = imports.iter().copied().map(String::from).collect();
4039 let site = &site;
4040 let env = &env;
4041 async move {
4042 crate::handler_dispatch::build_bindings(
4043 inner,
4044 ProjectRef::new("default"),
4045 "shop",
4046 "shop",
4047 None,
4048 &imports,
4049 site,
4050 env,
4051 &[],
4052 0,
4053 None,
4054 None,
4055 None,
4056 None,
4057 None,
4058 )
4059 .await
4060 .expect("no secrets → resolves")
4061 .sql_database_names()
4062 }
4063 };
4064
4065 assert_eq!(build(&["sql", "sql:product"]).await, vec!["", "product"]);
4068 assert_eq!(
4070 build(&["sql", "sql:*"]).await,
4071 vec!["", "privileged", "product"]
4072 );
4073 assert!(build(&["sql:secret"]).await.is_empty());
4075 assert_eq!(build(&["sql:product"]).await, vec!["product"]);
4077
4078 let _ = std::fs::remove_dir_all(&sql_dir);
4079 }
4080}