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 #[cfg(feature = "capability")]
423 capability_max_ttl_secs: std::sync::OnceLock<u64>,
424 #[cfg(feature = "handlers")]
429 tenancy_posture_overrides: std::sync::OnceLock<
430 Arc<std::collections::BTreeMap<String, boatramp_core::security::ResolvedProjectTenancy>>,
431 >,
432}
433
434#[cfg(feature = "handlers")]
435impl HandlerRuntimeInner {
436 pub(crate) fn project_tenancy_knobs(
442 &self,
443 project: &str,
444 ) -> boatramp_core::security::ResolvedProjectTenancy {
445 if let Some(map) = self.tenancy_posture_overrides.get() {
446 if let Some(knobs) = map.get(project) {
447 return *knobs;
448 }
449 }
450 boatramp_core::security::ResolvedProjectTenancy {
451 require_tenancy_declaration: self
452 .require_tenancy_declaration
453 .get()
454 .copied()
455 .unwrap_or(true),
456 allow_cross_tenant_db: self.allow_cross_tenant_db.get().copied().unwrap_or(false),
457 #[cfg(feature = "capability")]
458 capability_max_ttl_secs: self.capability_max_ttl_secs.get().copied(),
459 #[cfg(not(feature = "capability"))]
460 capability_max_ttl_secs: None,
461 }
462 }
463}
464
465pub type CronLeaderGate = Arc<dyn Fn() -> bool + Send + Sync>;
468
469impl HandlerRuntime {
470 pub fn disabled() -> Self {
472 Self::default()
473 }
474
475 #[cfg(feature = "handlers")]
481 pub fn new(
482 engine: boatramp_handlers::HandlerEngine,
483 kv: Arc<dyn boatramp_core::kv::KvStore>,
484 storage: Arc<dyn boatramp_core::Storage>,
485 sql: Option<Arc<dyn boatramp_core::sql::SqlBackends>>,
486 messaging: Option<Arc<dyn boatramp_core::messaging::Messaging>>,
487 ) -> Self {
488 let async_drain_slots = engine.async_max_concurrency().max(1);
491 Self {
492 inner: Some(Arc::new(HandlerRuntimeInner {
493 engine,
494 async_drain_gate: Arc::new(tokio::sync::Semaphore::new(async_drain_slots)),
495 kv,
496 storage,
497 sql,
498 messaging,
499 site_semaphores: std::sync::Mutex::new(std::collections::HashMap::new()),
500 stream_semaphores: std::sync::Mutex::new(std::collections::HashMap::new()),
501 stream_ip_counts: Arc::new(std::sync::Mutex::new(std::collections::HashMap::new())),
502 metrics: metrics::Metrics::default(),
503 logs: Arc::new(logs::LogStore::default()),
504 #[cfg(feature = "handlers")]
505 graphql_cache: graphql_cache::GraphqlCache::default(),
506 cron_leader_gate: std::sync::OnceLock::new(),
507 max_blob_bytes: std::sync::OnceLock::new(),
508 max_component_bytes: std::sync::OnceLock::new(),
509 allow_env_secret_refs: std::sync::OnceLock::new(),
510 require_tenancy_declaration: std::sync::OnceLock::new(),
511 allow_cross_tenant_db: std::sync::OnceLock::new(),
512 secret_store: std::sync::OnceLock::new(),
513 function_meter_locks: std::sync::Mutex::new(std::collections::HashMap::new()),
514 function_semaphores: std::sync::Mutex::new(std::collections::HashMap::new()),
515 watch_provider: std::sync::OnceLock::new(),
516 provision_tier: std::sync::OnceLock::new(),
517 invoker: std::sync::OnceLock::new(),
518 federation_runner: std::sync::OnceLock::new(),
519 #[cfg(feature = "email")]
520 email_profile_store: std::sync::OnceLock::new(),
521 #[cfg(feature = "email")]
522 email_spool: std::sync::OnceLock::new(),
523 #[cfg(feature = "admin")]
524 admin_controller: std::sync::OnceLock::new(),
525 #[cfg(feature = "admin")]
526 admin_surfaces: std::sync::OnceLock::new(),
527 #[cfg(feature = "session")]
528 session_store: std::sync::OnceLock::new(),
529 session_signer: std::sync::OnceLock::new(),
530 #[cfg(feature = "capability")]
531 capability_max_ttl_secs: std::sync::OnceLock::new(),
532 #[cfg(feature = "handlers")]
533 tenancy_posture_overrides: std::sync::OnceLock::new(),
534 })),
535 }
536 }
537
538 #[cfg(feature = "handlers")]
541 pub(crate) fn sql_provider(&self) -> Option<Arc<dyn boatramp_core::sql::SqlBackends>> {
542 self.inner.as_ref().and_then(|inner| inner.sql.clone())
543 }
544
545 #[cfg(feature = "handlers")]
549 pub(crate) fn invoker(&self) -> Option<Arc<function_runtime::FunctionInvoker>> {
550 self.inner
551 .as_ref()
552 .and_then(|inner| inner.invoker.get().cloned())
553 }
554
555 #[cfg(feature = "handlers")]
559 pub(crate) async fn introspect_subgraph_sdl(
560 &self,
561 deploy: &DeployStore,
562 project: boatramp_core::project::ProjectRef<'_>,
563 function: &boatramp_core::function::Function,
564 component: &str,
565 ) -> Result<String, function_runtime::SubgraphSdlError> {
566 match self.inner.as_ref() {
567 Some(inner) => {
568 function_runtime::introspect_service_sdl(
569 inner, deploy, project, function, component,
570 )
571 .await
572 }
573 None => Err(function_runtime::SubgraphSdlError::Unavailable),
574 }
575 }
576
577 #[cfg(feature = "handlers")]
583 pub fn set_invoker(&self, deploy: DeployStore) {
584 if let Some(inner) = self.inner.as_ref() {
585 let invoker = Arc::new(function_runtime::FunctionInvoker::new(
586 deploy,
587 Arc::downgrade(inner),
588 ));
589 let _ = inner.invoker.set(invoker);
590 let runner = Arc::new(graphql_gateway::FederationRunner::new(Arc::downgrade(
593 inner,
594 )));
595 let _ = inner.federation_runner.set(runner);
596 }
597 }
598
599 #[cfg(feature = "handlers")]
603 pub fn sql_backends(&self) -> Option<Arc<dyn boatramp_core::sql::SqlBackends>> {
604 self.inner.as_ref().and_then(|inner| inner.sql.clone())
605 }
606
607 #[cfg(feature = "handlers")]
611 pub fn set_watch_provider(
612 &self,
613 provider: Arc<dyn boatramp_core::blob_provision::WatchProvider>,
614 ) {
615 if let Some(inner) = self.inner.as_ref() {
616 let _ = inner.watch_provider.set(provider);
617 }
618 }
619
620 #[cfg(feature = "handlers")]
624 pub fn set_provision_tier(&self, tier: boatramp_core::blob_notify::ProvisionTier) {
625 if let Some(inner) = self.inner.as_ref() {
626 let _ = inner.provision_tier.set(tier);
627 }
628 }
629
630 #[cfg(feature = "handlers")]
634 pub fn set_max_blob_bytes(&self, max_bytes: u64) {
635 if let Some(inner) = self.inner.as_ref() {
636 let _ = inner.max_blob_bytes.set(max_bytes);
637 }
638 }
639
640 #[cfg(feature = "handlers")]
643 pub fn set_max_component_bytes(&self, max_bytes: u64) {
644 if let Some(inner) = self.inner.as_ref() {
645 let _ = inner.max_component_bytes.set(max_bytes);
646 }
647 }
648
649 #[cfg(feature = "handlers")]
657 pub fn set_allow_env_secret_refs(&self, allow: bool) {
658 if let Some(inner) = self.inner.as_ref() {
659 let _ = inner.allow_env_secret_refs.set(allow);
660 }
661 }
662
663 #[cfg(feature = "handlers")]
668 pub fn set_tenancy_posture(&self, require_declaration: bool, allow_cross_tenant: bool) {
669 if let Some(inner) = self.inner.as_ref() {
670 let _ = inner.require_tenancy_declaration.set(require_declaration);
671 let _ = inner.allow_cross_tenant_db.set(allow_cross_tenant);
672 }
673 }
674
675 #[cfg(feature = "handlers")]
682 pub fn set_project_tenancy_overrides(
683 &self,
684 overrides: std::collections::BTreeMap<
685 String,
686 boatramp_core::security::ResolvedProjectTenancy,
687 >,
688 ) {
689 if overrides.is_empty() {
690 return;
691 }
692 if let Some(inner) = self.inner.as_ref() {
693 let _ = inner.tenancy_posture_overrides.set(Arc::new(overrides));
694 }
695 }
696
697 #[cfg(feature = "handlers")]
703 pub fn set_session_signer(&self, signer: Arc<dyn Signer>) {
704 if let Some(inner) = self.inner.as_ref() {
705 let _ = inner.session_signer.set(signer);
706 }
707 }
708
709 #[cfg(feature = "capability")]
716 pub fn set_capability_minting(&self, max_ttl_secs: u64) {
717 if max_ttl_secs == 0 {
718 return;
719 }
720 if let Some(inner) = self.inner.as_ref() {
721 let _ = inner.capability_max_ttl_secs.set(max_ttl_secs);
722 }
723 }
724
725 #[cfg(feature = "handlers")]
730 pub fn set_secret_store(&self, store: Arc<boatramp_core::secret_store::SecretStore>) {
731 if let Some(inner) = self.inner.as_ref() {
732 let _ = inner.secret_store.set(store);
733 }
734 }
735
736 #[cfg(feature = "email")]
740 pub fn set_email_profile_store(
741 &self,
742 store: Arc<boatramp_core::email_config::EmailProfileStore>,
743 ) {
744 if let Some(inner) = self.inner.as_ref() {
745 let _ = inner.email_profile_store.set(store);
746 }
747 }
748
749 #[cfg(feature = "email")]
754 pub fn set_email_spool(&self, spool: Arc<dyn boatramp_handlers::EmailSpool>) {
755 if let Some(inner) = self.inner.as_ref() {
756 let _ = inner.email_spool.set(spool);
757 }
758 }
759
760 #[cfg(feature = "admin")]
765 pub fn set_admin(
766 &self,
767 controller: Arc<admin_controller::ServerAdminController>,
768 surfaces: std::collections::BTreeSet<boatramp_handlers::AdminSurface>,
769 ) {
770 if let Some(inner) = self.inner.as_ref() {
771 let _ = inner.admin_controller.set(controller);
772 let _ = inner.admin_surfaces.set(surfaces);
773 }
774 }
775
776 #[cfg(feature = "handlers")]
781 pub(crate) fn allow_env_secret_refs(&self) -> bool {
782 self.inner
783 .as_ref()
784 .and_then(|inner| inner.allow_env_secret_refs.get().copied())
785 .unwrap_or(false)
786 }
787
788 #[cfg(feature = "handlers")]
793 pub fn set_cron_leader_gate(&self, gate: CronLeaderGate) {
794 if let Some(inner) = self.inner.as_ref() {
795 let _ = inner.cron_leader_gate.set(gate);
796 }
797 }
798
799 #[cfg(feature = "handlers")]
806 async fn precheck_activation(
807 &self,
808 deploy: &DeployStore,
809 manifest: &Manifest,
810 site_config: Option<&SiteConfig>,
811 ) -> Result<(), String> {
812 let Some(inner) = self.inner.as_ref() else {
813 return Ok(());
814 };
815 if manifest.config.handlers.is_empty() && manifest.config.consumers.is_empty() {
818 return Ok(());
819 }
820 let site_handlers = site_config
822 .and_then(|c| c.handlers.as_ref())
823 .filter(|h| h.enabled)
824 .ok_or_else(|| {
825 "deployment ships handlers/consumers but the site has them disabled".to_string()
826 })?;
827 let max_component = inner.max_component_bytes.get().copied().unwrap_or(0);
828
829 let allow_env_secret_refs = inner.allow_env_secret_refs.get().copied().unwrap_or(false);
836 crate::handler_dispatch::admit_secret_refs(&site_handlers.secrets, allow_env_secret_refs)
837 .map_err(|err| format!("handler secrets: {err}"))?;
838
839 let sync_ceiling = inner.engine.sync_timeout_ms();
845 let async_ceiling = inner.engine.async_timeout_ms();
846 if let Some(ms) = site_handlers.max_timeout_ms {
847 if u64::from(ms) > sync_ceiling {
848 tracing::warn!(
849 "site max_timeout_ms={ms} exceeds sync_max_timeout_ms={sync_ceiling}: \
850 synchronous HTTP handlers are capped at {sync_ceiling}ms; the extra time \
851 applies only to async calls (?mode=async / triggers), capped at \
852 async_max_timeout_ms={async_ceiling}"
853 );
854 }
855 }
856
857 for handler in &manifest.config.handlers {
859 if let Some(ms) = handler.limits.as_ref().and_then(|l| l.timeout_ms) {
860 if u64::from(ms) > sync_ceiling {
861 let route = &handler.route;
862 tracing::warn!(
863 "route {route:?} declares limits.timeout_ms={ms}, above \
864 sync_max_timeout_ms={sync_ceiling}: synchronous HTTP calls to this route \
865 are capped at {sync_ceiling}ms; the {ms}ms only applies to async calls \
866 (?mode=async / a queue trigger / a #[consumer]), capped at \
867 async_max_timeout_ms={async_ceiling}. If you need {ms}ms synchronously, \
868 that isn't possible — move the work to the async lane"
869 );
870 }
871 }
872 if !handler.streaming {
877 if let Some(entry) = manifest.files.get(&handler.component) {
878 if let Ok(bytes) = read_blob_bytes(deploy, &entry.hash).await {
879 if crate::function_api::component_declares_streaming_route(
880 &bytes,
881 &handler.route,
882 ) {
883 let route = &handler.route;
884 tracing::warn!(
885 "route {route:?} is a streaming handler (#[handler(stream)]) but \
886 its config lacks streaming = true: it will run on the sync request \
887 lane and be cut at sync_max_timeout_ms={sync_ceiling}ms. Set \
888 streaming = true so it serves on the dedicated streaming lane (its \
889 own concurrency budget + a much larger wall-clock)."
890 );
891 }
892 }
893 }
894 }
895 if let Some(entry) = manifest.files.get(&handler.component) {
899 if let Ok(bytes) = read_blob_bytes(deploy, &entry.hash).await {
900 let unmet = crate::function_api::unmet_requires(&bytes);
901 if !unmet.is_empty() {
902 return Err(format!(
903 "route {:?} [{}] requires capabilities this host does not implement: \
904 {}. Upgrade boatramp or enable those features — see `boatramp \
905 capabilities`.",
906 handler.route,
907 handler.methods.join(","),
908 unmet.join(", ")
909 ));
910 }
911 }
912 }
913 precheck_component(
914 deploy,
915 manifest,
916 site_handlers,
917 inner,
918 max_component,
919 &handler.imports,
920 &handler.component,
921 &format!("route {:?} [{}]", handler.route, handler.methods.join(",")),
924 false,
925 )
926 .await?;
927 }
928 for consumer in &manifest.config.consumers {
929 precheck_component(
930 deploy,
931 manifest,
932 site_handlers,
933 inner,
934 max_component,
935 &consumer.imports,
936 &consumer.component,
937 &format!("consumer {:?}", consumer.topic),
938 true,
939 )
940 .await?;
941 }
942 Ok(())
943 }
944
945 #[cfg(not(feature = "handlers"))]
946 async fn precheck_activation(
947 &self,
948 _deploy: &DeployStore,
949 _manifest: &Manifest,
950 _site_config: Option<&SiteConfig>,
951 ) -> Result<(), String> {
952 Ok(())
953 }
954}
955
956#[derive(Default, Clone)]
962pub struct ServerOptions {
963 pub limits: ServerLimits,
965 pub probe: Option<Arc<dyn boatramp_core::domain_verify::DomainProbe>>,
968 pub default_site: Option<String>,
971 pub implicit_routing: bool,
977 pub protect_previews: bool,
982 pub cluster_rate_limit_kv: Option<Arc<dyn boatramp_core::kv::KvStore>>,
986 pub issuer: Option<Arc<dyn Signer>>,
990 pub bootstrap_secret: Option<String>,
995 pub bootstrap_attestation: Option<String>,
1001 pub mesh_control: Option<Arc<dyn MeshControl>>,
1005 pub cors_allowed_origins: Vec<String>,
1015 #[cfg(feature = "oidc")]
1018 pub oidc_verifier: Option<Arc<oidc::OidcVerifier>>,
1019 pub posture: boatramp_core::security::SecurityPosture,
1023 pub served_over_tls: bool,
1028 pub pop_origin: Option<String>,
1035 pub daemon_runtime: Option<Arc<DaemonRuntime>>,
1039 pub operator_sql: Option<Arc<dyn boatramp_core::sql::OperatorSql>>,
1043 pub tenant_deprovisioner: Option<Arc<dyn boatramp_core::sql::TenantDeprovisioner>>,
1048 pub compute_exec: Option<Arc<dyn boatramp_core::compute::ComputeExec>>,
1052 pub compute_volumes: Option<Arc<dyn boatramp_core::compute::ComputeVolumes>>,
1057 pub compute_control: Option<Arc<dyn boatramp_core::compute::ComputeControl>>,
1061 pub secret_store: Option<Arc<boatramp_core::secret_store::SecretStore>>,
1069 pub email_profile_store: Option<Arc<boatramp_core::email_config::EmailProfileStore>>,
1077 #[cfg(feature = "console")]
1081 pub console: Option<console::ConsoleMount>,
1082}
1083
1084#[derive(Clone, Copy)]
1088struct ServedOverTls(bool);
1089
1090#[derive(Clone, Copy, Default)]
1095struct ImplicitRouting(bool);
1096
1097const DAEMON_RELOAD_BACKSTOP: std::time::Duration = std::time::Duration::from_secs(300);
1109
1110pub struct DaemonRuntime {
1111 baseline: boatramp_core::daemon_config::ConfigBaseline,
1112 state: std::sync::RwLock<DaemonState>,
1113 reload: tokio::sync::Notify,
1116}
1117
1118struct DaemonState {
1119 effective: Arc<boatramp_core::daemon_config::EffectiveConfig>,
1120 generation: Option<String>,
1121}
1122
1123pub fn config_baseline(options: &ServerOptions) -> boatramp_core::daemon_config::ConfigBaseline {
1128 #[cfg(feature = "console")]
1132 let (console_enabled, console_host, console_path) = match options.console.as_ref() {
1133 Some(m) => (true, Some(m.host.clone()), Some(m.path.clone())),
1134 None => (false, None, None),
1135 };
1136 #[cfg(not(feature = "console"))]
1137 let (console_enabled, console_host, console_path) = (false, None, None);
1138 boatramp_core::daemon_config::ConfigBaseline {
1139 default_site: options.default_site.clone(),
1140 protect_previews: options.protect_previews,
1141 max_upload_bytes: options.limits.max_upload_bytes.unwrap_or(0),
1142 upload_idle_timeout_secs: options.limits.upload_idle_timeout.map(|d| d.as_secs()),
1143 max_concurrent_uploads: options.limits.max_concurrent_uploads.map(|n| n as u64),
1144 cluster_rate_limit: options.cluster_rate_limit_kv.is_some(),
1145 compute_vcpus: 0,
1146 compute_mem_mib: 0,
1147 console_enabled,
1148 console_host,
1149 console_path,
1150 max_upload_ceiling: options.posture.max_upload_bytes,
1151 max_concurrent_uploads_ceiling: None,
1152 posture: options.posture,
1153 }
1154}
1155
1156impl DaemonRuntime {
1157 pub fn new(baseline: boatramp_core::daemon_config::ConfigBaseline) -> Self {
1161 let effective =
1162 Arc::new(boatramp_core::daemon_config::DaemonConfig::default().resolve(&baseline));
1163 Self {
1164 baseline,
1165 state: std::sync::RwLock::new(DaemonState {
1166 effective,
1167 generation: None,
1168 }),
1169 reload: tokio::sync::Notify::new(),
1170 }
1171 }
1172
1173 pub fn notify_reload(&self) {
1177 self.reload.notify_one();
1178 }
1179
1180 pub fn effective(&self) -> Arc<boatramp_core::daemon_config::EffectiveConfig> {
1182 self.state
1183 .read()
1184 .expect("daemon config lock")
1185 .effective
1186 .clone()
1187 }
1188
1189 pub fn generation(&self) -> Option<String> {
1192 self.state
1193 .read()
1194 .expect("daemon config lock")
1195 .generation
1196 .clone()
1197 }
1198
1199 pub fn baseline(&self) -> &boatramp_core::daemon_config::ConfigBaseline {
1201 &self.baseline
1202 }
1203
1204 pub async fn reload(&self, deploy: &DeployStore) -> Result<(), DeployError> {
1207 let cfg = deploy.get_daemon_config().await?.unwrap_or_default();
1208 let generation = deploy.daemon_config_generation().await?;
1209 let effective = Arc::new(cfg.resolve(&self.baseline));
1210 *self.state.write().expect("daemon config lock") = DaemonState {
1211 effective,
1212 generation,
1213 };
1214 Ok(())
1215 }
1216}
1217
1218#[derive(Clone, Copy, Default)]
1221struct PreviewPolicy {
1222 protect: bool,
1223}
1224
1225#[derive(Clone, Default)]
1230struct Issuer(Option<Arc<dyn Signer>>);
1231
1232#[derive(Clone, Default)]
1237struct BootstrapGate(Option<Arc<BootstrapInner>>);
1238
1239struct BootstrapInner {
1240 secret_hash: String,
1243 lock: tokio::sync::Mutex<()>,
1246}
1247
1248impl BootstrapGate {
1249 fn new(secret: Option<&str>) -> Self {
1250 Self(secret.filter(|s| !s.is_empty()).map(|s| {
1251 Arc::new(BootstrapInner {
1252 secret_hash: boatramp_core::deploy::sha256_hex(s.as_bytes()),
1253 lock: tokio::sync::Mutex::new(()),
1254 })
1255 }))
1256 }
1257}
1258
1259#[async_trait::async_trait]
1263pub trait MeshControl: Send + Sync {
1264 async fn admit(
1272 &self,
1273 mesh_pubkey_hex: &str,
1274 jti: &str,
1275 possession_proof: &[u8],
1276 proof_iat: u64,
1277 now: u64,
1278 advertise_addr: Option<&str>,
1279 ) -> Result<JoinOutcome, String>;
1280
1281 async fn rotate_key(&self) -> Result<String, String>;
1285
1286 async fn revoke(&self, node: u64) -> Result<(), String>;
1290
1291 async fn members(&self) -> Result<Vec<MeshMember>, String>;
1295
1296 async fn promote(&self, node: u64) -> Result<(), String>;
1299}
1300
1301pub enum JoinOutcome {
1303 Admitted {
1306 members: Vec<String>,
1308 addrs: std::collections::BTreeMap<u64, String>,
1310 },
1311 TokenSpent,
1313 ProofInvalid,
1315 Revoked,
1318}
1319
1320#[derive(Debug, Clone, Serialize)]
1322pub struct MeshMember {
1323 pub node: u64,
1325 pub voter: bool,
1327 pub caught_up: bool,
1329 pub leader: bool,
1331 #[serde(default, skip_serializing_if = "Option::is_none")]
1335 pub addr: Option<String>,
1336}
1337
1338#[derive(Clone, Default)]
1341struct MeshControlHandle(Option<Arc<dyn MeshControl>>);
1342
1343#[cfg(feature = "oidc")]
1345#[derive(Clone, Default)]
1346struct OidcState(Option<Arc<oidc::OidcVerifier>>);
1347
1348#[cfg(feature = "oidc")]
1351const EXCHANGE_TTL_SECS: u64 = 3600;
1352
1353use boatramp_core::time::now_unix;
1354
1355#[derive(Clone)]
1357struct CorsState(Arc<Vec<String>>);
1358
1359const CORS_ALLOW_METHODS: &str = "GET, POST, PUT, DELETE, OPTIONS";
1361const CORS_ALLOW_HEADERS: &str = "authorization, content-type";
1364const CORS_MAX_AGE: &str = "600";
1366
1367fn cors_origin_allowed(allowed: &[String], origin: &str) -> bool {
1371 allowed.iter().any(|a| a == "*" || a == origin)
1372}
1373
1374async fn cors(
1381 State(allowed): State<CorsState>,
1382 request: Request,
1383 next: axum::middleware::Next,
1384) -> Response {
1385 let origin = request
1386 .headers()
1387 .get(header::ORIGIN)
1388 .and_then(|v| v.to_str().ok())
1389 .filter(|o| cors_origin_allowed(&allowed.0, o))
1390 .map(str::to_string);
1391 let is_preflight = request.method() == Method::OPTIONS
1393 && request
1394 .headers()
1395 .contains_key(header::ACCESS_CONTROL_REQUEST_METHOD);
1396 if is_preflight {
1397 let allow_headers = request
1399 .headers()
1400 .get(header::ACCESS_CONTROL_REQUEST_HEADERS)
1401 .and_then(|v| v.to_str().ok())
1402 .map(str::to_string)
1403 .unwrap_or_else(|| CORS_ALLOW_HEADERS.to_string());
1404 let mut response = Response::new(Body::empty());
1405 *response.status_mut() = StatusCode::NO_CONTENT;
1406 if let Some(origin) = origin {
1407 let headers = response.headers_mut();
1408 set_header(headers, header::ACCESS_CONTROL_ALLOW_ORIGIN, &origin);
1409 set_header(headers, header::VARY, "Origin");
1410 set_header(
1411 headers,
1412 header::ACCESS_CONTROL_ALLOW_METHODS,
1413 CORS_ALLOW_METHODS,
1414 );
1415 set_header(
1416 headers,
1417 header::ACCESS_CONTROL_ALLOW_HEADERS,
1418 &allow_headers,
1419 );
1420 set_header(headers, header::ACCESS_CONTROL_MAX_AGE, CORS_MAX_AGE);
1421 }
1422 return response;
1423 }
1424 let mut response = next.run(request).await;
1425 if let Some(origin) = origin {
1426 let headers = response.headers_mut();
1427 set_header(headers, header::ACCESS_CONTROL_ALLOW_ORIGIN, &origin);
1428 if let Ok(value) = HeaderValue::from_str("Origin") {
1431 headers.append(header::VARY, value);
1432 }
1433 }
1434 response
1435}
1436
1437const DRAIN_DEADLINE: Duration = Duration::from_secs(30);
1442
1443#[derive(Debug, thiserror::Error)]
1445pub enum ServeError {
1446 #[error("server I/O: {0}")]
1448 Io(#[from] std::io::Error),
1449}
1450
1451pub async fn serve(
1454 addr: SocketAddr,
1455 deploy: DeployStore,
1456 auth: Auth,
1457 handlers: HandlerRuntime,
1458) -> Result<(), ServeError> {
1459 serve_with(addr, deploy, auth, handlers, ServerOptions::default()).await
1460}
1461
1462pub(crate) fn disable_nagle(stream: &mut tokio::net::TcpStream) {
1472 if let Err(err) = stream.set_nodelay(true) {
1473 tracing::debug!(%err, "failed to set TCP_NODELAY on an accepted connection");
1474 }
1475}
1476
1477pub async fn serve_with(
1479 addr: SocketAddr,
1480 deploy: DeployStore,
1481 auth: Auth,
1482 handlers: HandlerRuntime,
1483 options: ServerOptions,
1484) -> Result<(), ServeError> {
1485 let tcp = tokio::net::TcpListener::bind(addr).await?;
1486 tracing::info!(%addr, auth = !auth.is_disabled(), "boatramp server listening");
1487 let splice_ctx = splice::SpliceCtx {
1493 deploy: deploy.clone(),
1494 posture: options.posture,
1495 daemon: options.daemon_runtime.clone(),
1496 };
1497 #[cfg(feature = "handlers")]
1500 let scheduler = handlers.spawn_scheduler(deploy.clone());
1501 let gateway_prober = gateway::spawn_active_health_prober();
1505 let (router, fast) = router_with_fast(deploy, auth, handlers, options);
1508
1509 let (signalled_tx, signalled_rx) = tokio::sync::watch::channel(false);
1513 let server = splice::serve(tcp, splice_ctx, (router, fast), async move {
1517 shutdown_signal().await;
1518 let _ = signalled_tx.send(true);
1519 });
1520 let signalled = {
1521 let mut rx = signalled_rx;
1522 async move {
1523 let _ = rx.wait_for(|fired| *fired).await;
1524 }
1525 };
1526 let result = serve_with_drain_deadline(
1527 async move { server.await.map_err(ServeError::from) },
1528 signalled,
1529 DRAIN_DEADLINE,
1530 )
1531 .await;
1532 #[cfg(feature = "handlers")]
1534 if let Some(handle) = scheduler {
1535 handle.abort();
1536 }
1537 gateway_prober.abort();
1538 result
1539}
1540
1541async fn serve_with_drain_deadline<Srv, Sig>(
1546 server: Srv,
1547 signalled: Sig,
1548 deadline: Duration,
1549) -> Result<(), ServeError>
1550where
1551 Srv: Future<Output = Result<(), ServeError>>,
1552 Sig: Future<Output = ()>,
1553{
1554 tokio::pin!(server);
1555 let drain_cap = async move {
1556 signalled.await;
1557 tokio::time::sleep(deadline).await;
1558 };
1559 tokio::select! {
1560 result = &mut server => result,
1561 _ = drain_cap => {
1562 tracing::warn!(
1563 deadline_s = deadline.as_secs(),
1564 "drain deadline exceeded; forcing shutdown with requests still in flight"
1565 );
1566 Ok(())
1567 }
1568 }
1569}
1570
1571pub async fn shutdown_signal() {
1574 let ctrl_c = async {
1575 let _ = tokio::signal::ctrl_c().await;
1576 };
1577 #[cfg(unix)]
1578 let terminate = async {
1579 if let Ok(mut sig) =
1580 tokio::signal::unix::signal(tokio::signal::unix::SignalKind::terminate())
1581 {
1582 sig.recv().await;
1583 }
1584 };
1585 #[cfg(not(unix))]
1586 let terminate = std::future::pending::<()>();
1587
1588 tokio::select! {
1589 _ = ctrl_c => {}
1590 _ = terminate => {}
1591 }
1592 tracing::info!("shutdown signal received; draining");
1593}
1594
1595async fn healthz(Extension(daemon): Extension<Arc<DaemonRuntime>>) -> String {
1599 match daemon.generation() {
1600 Some(gen) => format!("ok gen={gen}"),
1601 None => "ok".to_string(),
1602 }
1603}
1604
1605async fn readyz(State(deploy): State<DeployStore>) -> Response {
1607 match deploy.ready().await {
1608 Ok(()) => (StatusCode::OK, "ready\n").into_response(),
1609 Err(err) => {
1610 tracing::warn!(error = %err, "readiness probe failed");
1611 (StatusCode::SERVICE_UNAVAILABLE, "not ready\n").into_response()
1612 }
1613 }
1614}
1615
1616#[derive(Clone)]
1621pub struct RequestId(pub String);
1622
1623#[cfg(feature = "handlers")]
1629#[derive(Clone)]
1630pub struct DomainContext(pub String);
1631
1632fn request_id_for(headers: &HeaderMap) -> String {
1635 if let Some(id) = headers
1636 .get("x-request-id")
1637 .and_then(|v| v.to_str().ok())
1638 .map(str::trim)
1639 .filter(|s| !s.is_empty())
1640 {
1641 return id.chars().filter(|c| !c.is_control()).take(128).collect();
1642 }
1643 use std::sync::atomic::{AtomicU64, Ordering};
1644 static SEQ: AtomicU64 = AtomicU64::new(0);
1645 let n = SEQ.fetch_add(1, Ordering::Relaxed);
1646 format!("{:x}-{:x}", boatramp_core::time::now_unix_ms(), n)
1647}
1648
1649struct AccessLog {
1653 request_id: String,
1654 method: Method,
1655 path: String,
1656 host: String,
1657 client: String,
1658 status: u16,
1659 encoding: String,
1661 start: std::time::Instant,
1662 bytes: std::sync::atomic::AtomicU64,
1663}
1664
1665impl Drop for AccessLog {
1666 fn drop(&mut self) {
1667 let bytes = self.bytes.load(std::sync::atomic::Ordering::Relaxed);
1668 srvmetrics::server_metrics().record_request(self.status, bytes);
1671 tracing::info!(
1672 target: "boatramp::access",
1673 request_id = %self.request_id,
1674 method = %self.method,
1675 path = %self.path,
1676 host = %self.host,
1677 client = %self.client,
1678 status = self.status,
1679 bytes = bytes,
1680 encoding = %self.encoding,
1681 cache_result = srvmetrics::cache_result(self.status),
1682 elapsed_ms = self.start.elapsed().as_millis() as u64,
1683 "request"
1684 );
1685 }
1686}
1687
1688pub(crate) fn assign_request_id(request: &mut axum::extract::Request) -> String {
1694 let request_id = request_id_for(request.headers());
1695 request
1696 .extensions_mut()
1697 .insert(RequestId(request_id.clone()));
1698 request_id
1699}
1700
1701pub(crate) struct AccessLogCtx {
1706 request_id: String,
1707 method: Method,
1708 path: String,
1709 host: String,
1710 client: String,
1711 start: std::time::Instant,
1712}
1713
1714impl AccessLogCtx {
1715 pub(crate) fn capture(request: &axum::extract::Request, request_id: String) -> Option<Self> {
1721 if !tracing::enabled!(target: "boatramp::access", tracing::Level::INFO) {
1722 return None;
1723 }
1724 Some(Self {
1725 request_id,
1726 method: request.method().clone(),
1727 path: request.uri().path().to_string(),
1728 host: request
1729 .headers()
1730 .get(header::HOST)
1731 .and_then(|value| value.to_str().ok())
1732 .or_else(|| request.uri().host()) .unwrap_or("-")
1734 .to_string(),
1735 client: request
1736 .extensions()
1737 .get::<axum::extract::ConnectInfo<SocketAddr>>()
1738 .map(|info| info.0.ip().to_string())
1739 .unwrap_or_else(|| "-".to_string()),
1740 start: std::time::Instant::now(),
1741 })
1742 }
1743
1744 pub(crate) fn finish(self, response: Response) -> Response {
1748 let encoding = response
1749 .headers()
1750 .get(header::CONTENT_ENCODING)
1751 .and_then(|v| v.to_str().ok())
1752 .unwrap_or("identity")
1753 .to_string();
1754 let log = AccessLog {
1755 request_id: self.request_id,
1756 method: self.method,
1757 path: self.path,
1758 host: self.host,
1759 client: self.client,
1760 status: response.status().as_u16(),
1761 encoding,
1762 start: self.start,
1763 bytes: std::sync::atomic::AtomicU64::new(0),
1764 };
1765 let (parts, body) = response.into_parts();
1766 let counted = body.into_data_stream().map(move |chunk| {
1767 if let Ok(bytes) = &chunk {
1768 log.bytes
1769 .fetch_add(bytes.len() as u64, std::sync::atomic::Ordering::Relaxed);
1770 }
1771 chunk
1772 });
1773 Response::from_parts(parts, Body::from_stream(counted))
1774 }
1775}
1776
1777async fn access_log(mut request: axum::extract::Request, next: axum::middleware::Next) -> Response {
1782 let request_id = assign_request_id(&mut request);
1783 match AccessLogCtx::capture(&request, request_id) {
1784 None => next.run(request).await,
1785 Some(ctx) => ctx.finish(next.run(request).await),
1786 }
1787}
1788
1789fn if_none_match(req_headers: &HeaderMap, etag: &str) -> bool {
1791 req_headers
1792 .get(header::IF_NONE_MATCH)
1793 .and_then(|value| value.to_str().ok())
1794 .is_some_and(|value| {
1795 value
1796 .split(',')
1797 .map(str::trim)
1798 .any(|tag| tag == "*" || tag == etag || tag.trim_start_matches("W/") == etag)
1799 })
1800}
1801
1802fn set_header(headers: &mut HeaderMap, name: header::HeaderName, value: &str) {
1803 if let Ok(value) = HeaderValue::from_str(value) {
1804 headers.insert(name, value);
1805 }
1806}
1807
1808fn not_found() -> Response {
1809 (StatusCode::NOT_FOUND, "not found\n").into_response()
1810}
1811
1812fn redirect(status: u16, location: &str) -> Response {
1813 let status = StatusCode::from_u16(status).unwrap_or(StatusCode::FOUND);
1814 match HeaderValue::from_str(location) {
1815 Ok(location) => {
1816 let mut headers = HeaderMap::new();
1817 headers.insert(header::LOCATION, location);
1818 (status, headers).into_response()
1819 }
1820 Err(_) => (StatusCode::INTERNAL_SERVER_ERROR, "bad redirect target\n").into_response(),
1821 }
1822}
1823
1824fn deploy_error_response(err: DeployError) -> Response {
1826 let status = match &err {
1827 DeployError::NotFound(_) | DeployError::Storage(StorageError::NotFound(_)) => {
1828 StatusCode::NOT_FOUND
1829 }
1830 DeployError::HashMismatch { .. } => StatusCode::BAD_REQUEST,
1831 DeployError::Incomplete(_) => StatusCode::CONFLICT,
1832 DeployError::Conflict(_) => StatusCode::CONFLICT,
1834 DeployError::Ambiguous(_) => StatusCode::NOT_FOUND,
1836 DeployError::Invalid(_) => StatusCode::BAD_REQUEST,
1838 _ => StatusCode::INTERNAL_SERVER_ERROR,
1839 };
1840 tracing::warn!(error = %err, "request failed");
1841 (status, format!("{err}\n")).into_response()
1842}
1843
1844fn reject_invalid_name(kind: &'static str, value: &str) -> Option<Response> {
1849 boatramp_core::project::validate_resource_name(kind, value)
1850 .err()
1851 .map(|err| (StatusCode::UNPROCESSABLE_ENTITY, format!("{err}\n")).into_response())
1852}
1853
1854#[cfg(test)]
1855mod drain_tests {
1856 use super::*;
1857
1858 #[tokio::test]
1859 async fn deadline_forces_shutdown_after_signal() {
1860 let server = std::future::pending::<Result<(), ServeError>>();
1863 let signalled = async {}; let result = serve_with_drain_deadline(server, signalled, Duration::from_millis(20)).await;
1865 assert!(result.is_ok());
1866 }
1867
1868 #[tokio::test]
1869 async fn server_finishing_first_wins() {
1870 let server = async { Ok(()) };
1873 let signalled = std::future::pending::<()>();
1874 let result = serve_with_drain_deadline(server, signalled, Duration::from_secs(30)).await;
1875 assert!(result.is_ok());
1876 }
1877
1878 #[tokio::test]
1879 async fn deadline_does_not_trip_before_signal() {
1880 let server = async {
1884 tokio::time::sleep(Duration::from_millis(40)).await;
1885 Err(ServeError::Io(std::io::Error::other("server error")))
1886 };
1887 let signalled = std::future::pending::<()>();
1888 let result = serve_with_drain_deadline(server, signalled, Duration::from_millis(10)).await;
1889 assert!(result.is_err());
1890 }
1891}
1892
1893#[cfg(all(test, feature = "handlers"))]
1894mod tests {
1895 use super::*;
1896 use boatramp_core::cose::{LocalSigner, TokenAlg};
1897 use boatramp_core::project::ProjectRef;
1898
1899 #[test]
1900 fn query_string_parses_and_url_decodes() {
1901 let q = parse_query_string("lang=fr&city=S%C3%A3o+Paulo&flag&dup=1&dup=2");
1902 assert_eq!(q.get("lang").map(String::as_str), Some("fr"));
1903 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")); }
1907
1908 #[test]
1909 fn cookie_header_parses_pairs() {
1910 let c = parse_cookie_header("beta=1; sid = abc ; empty=");
1911 assert_eq!(c.get("beta").map(String::as_str), Some("1"));
1912 assert_eq!(c.get("sid").map(String::as_str), Some("abc"));
1913 assert_eq!(c.get("empty").map(String::as_str), Some(""));
1914 }
1915
1916 #[test]
1917 fn apply_vary_merges_without_duplicates() {
1918 let base = (StatusCode::OK, "x").into_response();
1919 let r = apply_vary(base, &["accept-language".into()]);
1920 assert_eq!(r.headers().get(header::VARY).unwrap(), "accept-language");
1921 let r = apply_vary(r, &["cookie".into(), "accept-language".into()]);
1923 let v = r.headers().get(header::VARY).unwrap().to_str().unwrap();
1924 assert!(v.contains("accept-language") && v.contains("cookie"));
1925 assert_eq!(v.matches("accept-language").count(), 1);
1926 let plain = apply_vary((StatusCode::OK, "y").into_response(), &[]);
1928 assert!(plain.headers().get(header::VARY).is_none());
1929 }
1930
1931 #[tokio::test]
1935 async fn join_token_endpoint_mints_a_verifiable_bearer_token() {
1936 let keys: Arc<dyn Signer> = Arc::new(LocalSigner::generate(TokenAlg::Es256));
1937 let public = keys.public_key();
1938
1939 let resp = create_join_token(
1941 Extension(Issuer(Some(keys.clone()))),
1942 Json(CreateJoinTokenRequest {
1943 ttl_secs: Some(600),
1944 }),
1945 )
1946 .await;
1947 assert_eq!(resp.status(), StatusCode::CREATED);
1948 let body = axum::body::to_bytes(resp.into_body(), usize::MAX)
1949 .await
1950 .unwrap();
1951 let parsed: serde_json::Value = serde_json::from_slice(&body).unwrap();
1952 let token = parsed["token"].as_str().unwrap();
1953 let jti = cose::verify_join(token, &public, now_unix()).unwrap();
1954 assert!(!jti.is_empty());
1955
1956 let no_issuer = create_join_token(
1958 Extension(Issuer(None)),
1959 Json(CreateJoinTokenRequest { ttl_secs: None }),
1960 )
1961 .await;
1962 assert_eq!(no_issuer.status(), StatusCode::NOT_IMPLEMENTED);
1963 }
1964
1965 #[tokio::test]
1970 async fn function_write_path_deploy_rollback_alias_remove() {
1971 use boatramp_core::function::Lifecycle;
1972 use boatramp_core::kv::MemoryKv;
1973 use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
1974
1975 struct FakeStorage {
1978 present: bool,
1979 }
1980 #[async_trait::async_trait]
1981 impl Storage for FakeStorage {
1982 async fn get(&self, _: &str) -> Result<GetObject, StorageError> {
1983 Err(StorageError::NotFound(String::new()))
1984 }
1985 async fn get_range(
1986 &self,
1987 _: &str,
1988 _: u64,
1989 _: Option<u64>,
1990 ) -> Result<GetObject, StorageError> {
1991 Err(StorageError::NotFound(String::new()))
1992 }
1993 async fn put(
1994 &self,
1995 _: &str,
1996 _: ByteStream,
1997 _: PutMeta,
1998 ) -> Result<ObjectMeta, StorageError> {
1999 Err(StorageError::unsupported("fake"))
2000 }
2001 async fn head(&self, key: &str) -> Result<ObjectMeta, StorageError> {
2002 if self.present {
2003 Ok(ObjectMeta {
2004 key: key.to_string(),
2005 ..Default::default()
2006 })
2007 } else {
2008 Err(StorageError::NotFound(key.to_string()))
2009 }
2010 }
2011 async fn delete(&self, _: &str) -> Result<(), StorageError> {
2012 Ok(())
2013 }
2014 async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
2015 Ok(Vec::new())
2016 }
2017 }
2018
2019 async fn body_json(resp: Response) -> (StatusCode, serde_json::Value) {
2020 let status = resp.status();
2021 let bytes = axum::body::to_bytes(resp.into_body(), usize::MAX)
2022 .await
2023 .unwrap();
2024 let value = if bytes.is_empty() {
2025 serde_json::Value::Null
2026 } else {
2027 serde_json::from_slice(&bytes).unwrap()
2028 };
2029 (status, value)
2030 }
2031
2032 let deploy = DeployStore::new(
2033 Arc::new(FakeStorage { present: true }),
2034 Arc::new(MemoryKv::new()),
2035 );
2036 let v1 = "a".repeat(64);
2037 let v2 = "b".repeat(64);
2038
2039 let (st, body) = body_json(
2041 deploy_function(
2042 State(deploy.clone()),
2043 axum::extract::Extension(crate::ProjectContext::default()),
2044 axum::extract::Extension(Arc::new(crate::HandlerRuntime::disabled())),
2045 axum::extract::Query(DeployFunctionQuery::default()),
2046 Path("greeter".to_string()),
2047 Json(FunctionUpsert {
2048 component: v1.clone(),
2049 config: Default::default(),
2050 lifecycle: Lifecycle::Independent,
2051 }),
2052 )
2053 .await,
2054 )
2055 .await;
2056 assert_eq!(st, StatusCode::OK);
2057 assert_eq!(body["active"], v1);
2058
2059 let (_, body) = body_json(
2061 deploy_function(
2062 State(deploy.clone()),
2063 axum::extract::Extension(crate::ProjectContext::default()),
2064 axum::extract::Extension(Arc::new(crate::HandlerRuntime::disabled())),
2065 axum::extract::Query(DeployFunctionQuery::default()),
2066 Path("greeter".to_string()),
2067 Json(FunctionUpsert {
2068 component: v2.clone(),
2069 config: Default::default(),
2070 lifecycle: Lifecycle::Independent,
2071 }),
2072 )
2073 .await,
2074 )
2075 .await;
2076 assert_eq!(body["active"], v2);
2077 assert_eq!(body["versions"].as_array().unwrap().len(), 2);
2078
2079 let (st, body) = body_json(
2081 rollback_function(
2082 State(deploy.clone()),
2083 axum::extract::Extension(crate::ProjectContext::default()),
2084 Path("greeter".to_string()),
2085 Json(RollbackBody { to: v1.clone() }),
2086 )
2087 .await,
2088 )
2089 .await;
2090 assert_eq!(st, StatusCode::OK);
2091 assert_eq!(body["active"], v1);
2092
2093 let resp = rollback_function(
2095 State(deploy.clone()),
2096 axum::extract::Extension(crate::ProjectContext::default()),
2097 Path("greeter".to_string()),
2098 Json(RollbackBody { to: "c".repeat(64) }),
2099 )
2100 .await;
2101 assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
2102
2103 let (st, body) = body_json(
2105 alias_function(
2106 State(deploy.clone()),
2107 axum::extract::Extension(crate::ProjectContext::default()),
2108 Path(("greeter".to_string(), "prod".to_string())),
2109 Json(AliasBody {
2110 version: v2.clone(),
2111 }),
2112 )
2113 .await,
2114 )
2115 .await;
2116 assert_eq!(st, StatusCode::OK);
2117 assert_eq!(body["aliases"]["prod"], v2);
2118
2119 let (st, _) = body_json(
2121 remove_function(
2122 State(deploy.clone()),
2123 axum::extract::Extension(crate::ProjectContext::default()),
2124 Path("greeter".to_string()),
2125 )
2126 .await,
2127 )
2128 .await;
2129 assert_eq!(st, StatusCode::NO_CONTENT);
2130 assert!(deploy
2131 .get_function(ProjectRef::DEFAULT, "greeter")
2132 .await
2133 .unwrap()
2134 .is_none());
2135
2136 let empty = DeployStore::new(
2138 Arc::new(FakeStorage { present: false }),
2139 Arc::new(MemoryKv::new()),
2140 );
2141 let resp = deploy_function(
2142 State(empty),
2143 axum::extract::Extension(crate::ProjectContext::default()),
2144 axum::extract::Extension(Arc::new(crate::HandlerRuntime::disabled())),
2145 axum::extract::Query(DeployFunctionQuery::default()),
2146 Path("orphan".to_string()),
2147 Json(FunctionUpsert {
2148 component: v1.clone(),
2149 config: Default::default(),
2150 lifecycle: Lifecycle::default(),
2151 }),
2152 )
2153 .await;
2154 assert_eq!(resp.status(), StatusCode::BAD_REQUEST);
2155 }
2156
2157 struct StubControl {
2161 admits: std::sync::Mutex<Vec<(String, String)>>,
2162 respond: StubJoin,
2163 }
2164 #[derive(Clone, Copy)]
2165 enum StubJoin {
2166 Admit,
2167 Spent,
2168 Invalid,
2169 Revoked,
2170 }
2171
2172 #[async_trait::async_trait]
2173 impl MeshControl for StubControl {
2174 async fn admit(
2175 &self,
2176 mesh_pubkey_hex: &str,
2177 jti: &str,
2178 _proof: &[u8],
2179 _proof_iat: u64,
2180 _now: u64,
2181 _advertise_addr: Option<&str>,
2182 ) -> Result<JoinOutcome, String> {
2183 self.admits
2184 .lock()
2185 .unwrap()
2186 .push((mesh_pubkey_hex.to_string(), jti.to_string()));
2187 Ok(match self.respond {
2188 StubJoin::Admit => JoinOutcome::Admitted {
2189 members: vec!["signed-member".to_string()],
2190 addrs: std::collections::BTreeMap::from([(7u64, "https://x:7000".to_string())]),
2191 },
2192 StubJoin::Spent => JoinOutcome::TokenSpent,
2193 StubJoin::Invalid => JoinOutcome::ProofInvalid,
2194 StubJoin::Revoked => JoinOutcome::Revoked,
2195 })
2196 }
2197 async fn rotate_key(&self) -> Result<String, String> {
2198 Ok("cafe".to_string())
2199 }
2200 async fn revoke(&self, _node: u64) -> Result<(), String> {
2201 Ok(())
2202 }
2203 async fn members(&self) -> Result<Vec<MeshMember>, String> {
2204 Ok(Vec::new())
2205 }
2206 async fn promote(&self, _node: u64) -> Result<(), String> {
2207 Ok(())
2208 }
2209 }
2210
2211 #[tokio::test]
2215 async fn cluster_join_dispatches_and_maps_outcomes() {
2216 let keys: Arc<dyn Signer> = Arc::new(LocalSigner::generate(TokenAlg::Es256));
2217 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
2218 let auth = Auth::with_key(keys.public_key(), kv);
2219 let token = cose::mint_join(600, now_unix(), &*keys).await.unwrap();
2220 let req = |proof: &str| JoinRequest {
2221 token: token.clone(),
2222 mesh_pubkey: "302a300506032b6570032100feed".into(),
2223 possession_proof: proof.to_string(),
2224 proof_iat: now_unix(),
2225 advertise_addr: Some("https://joiner:7000".into()),
2226 };
2227
2228 let admitter = Arc::new(StubControl {
2230 admits: std::sync::Mutex::new(Vec::new()),
2231 respond: StubJoin::Admit,
2232 });
2233 let resp = cluster_join(
2234 Extension(auth.clone()),
2235 Extension(MeshControlHandle(Some(admitter.clone()))),
2236 Json(req("aa01")),
2237 )
2238 .await;
2239 assert_eq!(resp.status(), StatusCode::OK);
2240 assert_eq!(admitter.admits.lock().unwrap().len(), 1);
2241
2242 let spent = Arc::new(StubControl {
2244 admits: std::sync::Mutex::new(Vec::new()),
2245 respond: StubJoin::Spent,
2246 });
2247 assert_eq!(
2248 cluster_join(
2249 Extension(auth.clone()),
2250 Extension(MeshControlHandle(Some(spent))),
2251 Json(req("aa01")),
2252 )
2253 .await
2254 .status(),
2255 StatusCode::CONFLICT
2256 );
2257 let invalid = Arc::new(StubControl {
2258 admits: std::sync::Mutex::new(Vec::new()),
2259 respond: StubJoin::Invalid,
2260 });
2261 assert_eq!(
2262 cluster_join(
2263 Extension(auth.clone()),
2264 Extension(MeshControlHandle(Some(invalid))),
2265 Json(req("aa01")),
2266 )
2267 .await
2268 .status(),
2269 StatusCode::FORBIDDEN
2270 );
2271 let revoked = Arc::new(StubControl {
2273 admits: std::sync::Mutex::new(Vec::new()),
2274 respond: StubJoin::Revoked,
2275 });
2276 assert_eq!(
2277 cluster_join(
2278 Extension(auth.clone()),
2279 Extension(MeshControlHandle(Some(revoked))),
2280 Json(req("aa01")),
2281 )
2282 .await
2283 .status(),
2284 StatusCode::FORBIDDEN
2285 );
2286
2287 let ok = Arc::new(StubControl {
2289 admits: std::sync::Mutex::new(Vec::new()),
2290 respond: StubJoin::Admit,
2291 });
2292 assert_eq!(
2293 cluster_join(
2294 Extension(auth.clone()),
2295 Extension(MeshControlHandle(Some(ok))),
2296 Json(req("not-hex")),
2297 )
2298 .await
2299 .status(),
2300 StatusCode::BAD_REQUEST
2301 );
2302
2303 let none = cluster_join(
2305 Extension(auth),
2306 Extension(MeshControlHandle(None)),
2307 Json(req("aa01")),
2308 )
2309 .await;
2310 assert_eq!(none.status(), StatusCode::NOT_IMPLEMENTED);
2311 }
2312
2313 #[tokio::test]
2317 async fn bootstrap_mints_the_first_token_once() {
2318 use axum::http::{header::AUTHORIZATION, HeaderMap, HeaderValue};
2319 let keys: Arc<dyn Signer> = Arc::new(LocalSigner::generate(TokenAlg::Es256));
2320 let public = keys.public_key();
2321 let deploy = DeployStore::new(
2322 Arc::new(MemStorage::default()),
2323 Arc::new(MemoryKv::new()) as Arc<dyn KvStore>,
2324 );
2325 let secret = "s3cr3t-bootstrap-value";
2326 let gate = BootstrapGate::new(Some(secret));
2327 let issuer = Issuer(Some(keys.clone()));
2328 let bearer = |s: &str| {
2329 let mut h = HeaderMap::new();
2330 h.insert(
2331 AUTHORIZATION,
2332 HeaderValue::from_str(&format!("Bearer {s}")).unwrap(),
2333 );
2334 h
2335 };
2336 let req = || BootstrapRequest {
2337 roles: vec!["admin".to_string()],
2338 ttl_secs: None,
2339 };
2340
2341 let bad = bootstrap_token(
2343 State(deploy.clone()),
2344 Extension(issuer.clone()),
2345 Extension(gate.clone()),
2346 bearer("wrong"),
2347 Json(req()),
2348 )
2349 .await;
2350 assert_eq!(bad.status(), StatusCode::UNAUTHORIZED);
2351
2352 let ok = bootstrap_token(
2354 State(deploy.clone()),
2355 Extension(issuer.clone()),
2356 Extension(gate.clone()),
2357 bearer(secret),
2358 Json(req()),
2359 )
2360 .await;
2361 assert_eq!(ok.status(), StatusCode::CREATED);
2362 let body = axum::body::to_bytes(ok.into_body(), usize::MAX)
2363 .await
2364 .unwrap();
2365 let json: serde_json::Value = serde_json::from_slice(&body).unwrap();
2366 let token = json["token"].as_str().unwrap();
2367 let id = json["id"].as_str().unwrap();
2368 let verified = cose::verify(token, &public, now_unix()).unwrap();
2369 assert!(verified.roles.iter().any(|r| r.name == "admin"));
2370 assert!(deploy
2371 .list_token_meta()
2372 .await
2373 .unwrap()
2374 .iter()
2375 .any(|m| m.revocation_id == id));
2376
2377 let reuse = bootstrap_token(
2379 State(deploy.clone()),
2380 Extension(issuer.clone()),
2381 Extension(gate),
2382 bearer(secret),
2383 Json(req()),
2384 )
2385 .await;
2386 assert_eq!(reuse.status(), StatusCode::CONFLICT);
2387
2388 let disabled = bootstrap_token(
2390 State(deploy),
2391 Extension(issuer),
2392 Extension(BootstrapGate(None)),
2393 bearer(secret),
2394 Json(req()),
2395 )
2396 .await;
2397 assert_eq!(disabled.status(), StatusCode::NOT_IMPLEMENTED);
2398 }
2399
2400 #[tokio::test]
2403 async fn cluster_rotate_key_returns_the_new_pubkey_or_501() {
2404 let control = Arc::new(StubControl {
2405 admits: std::sync::Mutex::new(Vec::new()),
2406 respond: StubJoin::Admit,
2407 });
2408 let resp = cluster_rotate_key(Extension(MeshControlHandle(Some(control)))).await;
2409 assert_eq!(resp.status(), StatusCode::OK);
2410 let body = axum::body::to_bytes(resp.into_body(), usize::MAX)
2411 .await
2412 .unwrap();
2413 let parsed: serde_json::Value = serde_json::from_slice(&body).unwrap();
2414 assert_eq!(parsed["pubkey"].as_str(), Some("cafe"));
2415
2416 let none = cluster_rotate_key(Extension(MeshControlHandle(None))).await;
2417 assert_eq!(none.status(), StatusCode::NOT_IMPLEMENTED);
2418 }
2419
2420 #[test]
2421 fn gateway_addr_gate_refuses_metadata_and_private_per_posture() {
2422 use boatramp_core::security::SecurityProfile;
2423 let strict = SecurityProfile::MultiTenant.preset();
2424 let loose = SecurityProfile::SingleTenant.preset(); let public: IpAddr = "93.184.216.34".parse().unwrap(); let private: IpAddr = "10.1.2.3".parse().unwrap();
2428 let loopback: IpAddr = "127.0.0.1".parse().unwrap();
2429 let metadata: IpAddr = IpAddr::V4(CLOUD_METADATA_IPV4);
2430
2431 assert!(gateway_addr_allowed(public, &strict));
2433 assert!(!gateway_addr_allowed(private, &strict));
2434 assert!(!gateway_addr_allowed(loopback, &strict));
2435 assert!(!gateway_addr_allowed(metadata, &strict));
2436
2437 assert!(gateway_addr_allowed(public, &loose));
2440 assert!(gateway_addr_allowed(private, &loose));
2441 assert!(gateway_addr_allowed(loopback, &loose));
2442 assert!(!gateway_addr_allowed(metadata, &loose));
2443 }
2444
2445 #[tokio::test]
2446 async fn resolve_env_merges_static_and_host_secrets() {
2447 use boatramp_core::config::HandlersSiteConfig;
2448
2449 std::env::set_var("BOATRAMP_TEST_RESOLVE_SECRET", "topsecret");
2451
2452 let deploy_env = std::collections::BTreeMap::from([
2453 ("GREETING".to_string(), "hi".to_string()),
2454 ("OVERRIDE_ME".to_string(), "static".to_string()),
2455 ]);
2456 let site_handlers = HandlersSiteConfig {
2457 enabled: true,
2458 secrets: std::collections::BTreeMap::from([
2459 (
2461 "SECRET_TOKEN".to_string(),
2462 "BOATRAMP_TEST_RESOLVE_SECRET".to_string(),
2463 ),
2464 (
2465 "OVERRIDE_ME".to_string(),
2466 "BOATRAMP_TEST_RESOLVE_SECRET".to_string(),
2467 ),
2468 (
2469 "MISSING".to_string(),
2470 "BOATRAMP_TEST_NOT_SET_VAR".to_string(),
2471 ),
2472 ]),
2473 ..Default::default()
2474 };
2475 let env = resolve_env(
2478 "blog",
2479 boatramp_core::project::ProjectRef::DEFAULT,
2480 &deploy_env,
2481 &site_handlers,
2482 true,
2483 None,
2484 )
2485 .await
2486 .expect("resolves");
2487
2488 assert!(env.contains(&("GREETING".to_string(), "hi".to_string())));
2492 assert!(env.contains(&("SECRET_TOKEN".to_string(), "topsecret".to_string())));
2493 assert!(env.contains(&("OVERRIDE_ME".to_string(), "topsecret".to_string())));
2494 assert!(!env.iter().any(|(k, _)| k == "MISSING"));
2495
2496 std::env::remove_var("BOATRAMP_TEST_RESOLVE_SECRET");
2497 }
2498
2499 #[tokio::test]
2500 async fn multi_tenant_posture_refuses_a_host_env_handler_secret() {
2501 use boatramp_core::config::HandlersSiteConfig;
2502
2503 std::env::set_var("BOATRAMP_TEST_OTHER_TENANT_SECRET", "leak-me");
2508 let deploy_env = std::collections::BTreeMap::new();
2509 let bare = HandlersSiteConfig {
2510 enabled: true,
2511 secrets: std::collections::BTreeMap::from([(
2512 "STOLEN".to_string(),
2513 "BOATRAMP_TEST_OTHER_TENANT_SECRET".to_string(),
2514 )]),
2515 ..Default::default()
2516 };
2517 let err = resolve_env(
2518 "evil",
2519 boatramp_core::project::ProjectRef::DEFAULT,
2520 &deploy_env,
2521 &bare,
2522 false,
2523 None,
2524 )
2525 .await
2526 .expect_err("multi-tenant must refuse a bare host-env ref");
2527 assert!(
2528 err.contains("STOLEN"),
2529 "error names the offending guest var: {err}"
2530 );
2531 assert!(
2532 err.contains("multi-tenant"),
2533 "error steers the tenant: {err}"
2534 );
2535 assert!(
2536 !err.contains("leak-me"),
2537 "the host value must never appear (never read): {err}"
2538 );
2539
2540 let explicit = HandlersSiteConfig {
2542 enabled: true,
2543 secrets: std::collections::BTreeMap::from([(
2544 "STOLEN".to_string(),
2545 "env:BOATRAMP_TEST_OTHER_TENANT_SECRET".to_string(),
2546 )]),
2547 ..Default::default()
2548 };
2549 assert!(resolve_env(
2550 "evil",
2551 boatramp_core::project::ProjectRef::DEFAULT,
2552 &deploy_env,
2553 &explicit,
2554 false,
2555 None,
2556 )
2557 .await
2558 .is_err());
2559
2560 let reserved = HandlersSiteConfig {
2562 enabled: true,
2563 secrets: std::collections::BTreeMap::from([(
2564 "TOKEN".to_string(),
2565 "vault:kv/data/app#token".to_string(),
2566 )]),
2567 ..Default::default()
2568 };
2569 let err = resolve_env(
2570 "evil",
2571 boatramp_core::project::ProjectRef::DEFAULT,
2572 &deploy_env,
2573 &reserved,
2574 true,
2575 None,
2576 )
2577 .await
2578 .expect_err("a reserved scheme is not yet supported, even under single-tenant");
2579 assert!(err.contains("not yet supported"), "{err}");
2580
2581 let arbitrary = HandlersSiteConfig {
2585 enabled: true,
2586 secrets: std::collections::BTreeMap::from([(
2587 "KEY".to_string(),
2588 "aws:sm/prod/apikey".to_string(),
2589 )]),
2590 ..Default::default()
2591 };
2592 let err = resolve_env(
2593 "evil",
2594 boatramp_core::project::ProjectRef::DEFAULT,
2595 &deploy_env,
2596 &arbitrary,
2597 true,
2598 None,
2599 )
2600 .await
2601 .expect_err("any unknown scheme is reserved, even under single-tenant");
2602 assert!(
2603 err.contains("not yet supported") && err.contains("aws"),
2604 "provider-neutral reservation names the scheme: {err}"
2605 );
2606
2607 std::env::remove_var("BOATRAMP_TEST_OTHER_TENANT_SECRET");
2608 }
2609
2610 #[tokio::test]
2611 async fn function_resolve_secret_env_reads_host_and_matches_handler_semantics() {
2612 std::env::set_var("BOATRAMP_TEST_FN_SECRET", "fnsecret");
2617
2618 let static_env = std::collections::BTreeMap::from([
2619 ("STAGE".to_string(), "prod".to_string()),
2620 ("OVERRIDE_ME".to_string(), "static".to_string()),
2621 ]);
2622 let secrets = std::collections::BTreeMap::from([
2623 ("DB_URL".to_string(), "BOATRAMP_TEST_FN_SECRET".to_string()),
2625 (
2627 "OVERRIDE_ME".to_string(),
2628 "BOATRAMP_TEST_FN_SECRET".to_string(),
2629 ),
2630 (
2632 "MISSING".to_string(),
2633 "BOATRAMP_TEST_FN_NOT_SET".to_string(),
2634 ),
2635 ]);
2636 let env = resolve_secret_env(
2638 "fn/api",
2639 boatramp_core::project::ProjectRef::DEFAULT,
2640 &static_env,
2641 &secrets,
2642 true,
2643 None,
2644 )
2645 .await
2646 .expect("resolves");
2647
2648 assert!(env.contains(&("STAGE".to_string(), "prod".to_string())));
2649 assert!(env.contains(&("DB_URL".to_string(), "fnsecret".to_string())));
2651 assert!(env.contains(&("OVERRIDE_ME".to_string(), "fnsecret".to_string())));
2653 assert!(!env.iter().any(|(k, _)| k == "MISSING"));
2655
2656 std::env::remove_var("BOATRAMP_TEST_FN_SECRET");
2657 }
2658
2659 #[tokio::test]
2660 async fn multi_tenant_posture_refuses_a_host_env_function_secret() {
2661 std::env::set_var("BOATRAMP_TEST_FN_LEAK", "leak-me");
2666 let static_env = std::collections::BTreeMap::new();
2667 let secrets = std::collections::BTreeMap::from([(
2668 "DB_URL".to_string(),
2669 "BOATRAMP_TEST_FN_LEAK".to_string(),
2670 )]);
2671
2672 let err = resolve_secret_env(
2674 "fn/api",
2675 boatramp_core::project::ProjectRef::DEFAULT,
2676 &static_env,
2677 &secrets,
2678 false,
2679 None,
2680 )
2681 .await
2682 .expect_err("multi-tenant must refuse a function host-env ref");
2683 assert!(
2684 err.contains("DB_URL"),
2685 "error names the offending guest var: {err}"
2686 );
2687 assert!(
2688 !err.contains("leak-me"),
2689 "host value must never appear: {err}"
2690 );
2691
2692 let env = resolve_secret_env(
2694 "fn/api",
2695 boatramp_core::project::ProjectRef::DEFAULT,
2696 &static_env,
2697 &secrets,
2698 true,
2699 None,
2700 )
2701 .await
2702 .expect("resolves");
2703 assert!(env.contains(&("DB_URL".to_string(), "leak-me".to_string())));
2704
2705 std::env::remove_var("BOATRAMP_TEST_FN_LEAK");
2706 }
2707
2708 #[tokio::test]
2709 async fn boatramp_scheme_resolves_from_the_project_scoped_store() {
2710 use boatramp_core::project::ProjectRef;
2711 use boatramp_core::secret_store::SecretStore;
2712 use std::sync::Arc;
2713
2714 struct XorEnvelope;
2716 #[async_trait::async_trait]
2717 impl boatramp_core::envelope::KeyEnvelope for XorEnvelope {
2718 async fn wrap(
2719 &self,
2720 p: &[u8],
2721 ) -> Result<Vec<u8>, boatramp_core::envelope::EnvelopeError> {
2722 Ok(p.iter().map(|b| b ^ 0x5a).collect())
2723 }
2724 async fn unwrap(
2725 &self,
2726 c: &[u8],
2727 ) -> Result<Vec<u8>, boatramp_core::envelope::EnvelopeError> {
2728 Ok(c.iter().map(|b| b ^ 0x5a).collect())
2729 }
2730 }
2731
2732 let store = SecretStore::new(
2733 Arc::new(boatramp_core::kv::MemoryKv::new()),
2734 Arc::new(XorEnvelope),
2735 );
2736 store
2737 .set(ProjectRef::new("acme"), "api-key", b"s3cr3t")
2738 .await
2739 .unwrap();
2740
2741 let static_env = std::collections::BTreeMap::new();
2742 let secrets = std::collections::BTreeMap::from([
2743 ("API_KEY".to_string(), "boatramp:api-key".to_string()),
2744 ("MISSING".to_string(), "boatramp:not-set".to_string()),
2745 ]);
2746
2747 let env = resolve_secret_env(
2750 "site",
2751 ProjectRef::new("acme"),
2752 &static_env,
2753 &secrets,
2754 false,
2755 Some(&store),
2756 )
2757 .await
2758 .expect("boatramp refs resolve without the host-env gate");
2759 assert!(env.contains(&("API_KEY".to_string(), "s3cr3t".to_string())));
2760 assert!(!env.iter().any(|(k, _)| k == "MISSING"));
2762
2763 let other_secrets = std::collections::BTreeMap::from([(
2765 "API_KEY".to_string(),
2766 "boatramp:api-key".to_string(),
2767 )]);
2768 let other = resolve_secret_env(
2769 "site",
2770 ProjectRef::new("globex"),
2771 &static_env,
2772 &other_secrets,
2773 false,
2774 Some(&store),
2775 )
2776 .await
2777 .expect("resolves (a foreign project's secret is simply absent → skipped)");
2778 assert!(
2779 !other.iter().any(|(k, _)| k == "API_KEY"),
2780 "a tenant must not read another project's secret"
2781 );
2782
2783 let err = resolve_secret_env(
2785 "site",
2786 ProjectRef::new("acme"),
2787 &static_env,
2788 &other_secrets,
2789 false,
2790 None,
2791 )
2792 .await
2793 .expect_err("no store configured must fail closed");
2794 assert!(err.contains("no internal secret store"), "{err}");
2795 }
2796
2797 fn req() -> Request {
2798 Request::builder()
2799 .uri("/")
2800 .header(header::HOST, "example.com")
2801 .body(Body::empty())
2802 .unwrap()
2803 }
2804
2805 #[test]
2806 fn forwarded_headers_set_standard_triple() {
2807 let mut request = req();
2808 set_forwarded_headers(&mut request, "203.0.113.7".parse().unwrap());
2809 let h = request.headers();
2810 assert_eq!(h.get("x-forwarded-for").unwrap(), "203.0.113.7");
2811 assert_eq!(h.get("x-forwarded-host").unwrap(), "example.com");
2812 assert_eq!(h.get("x-forwarded-proto").unwrap(), "http");
2813 }
2814
2815 #[test]
2816 fn forwarded_for_overwrites_spoofed_value() {
2817 let mut request = Request::builder()
2820 .uri("/")
2821 .header(header::HOST, "example.com")
2822 .header("x-forwarded-for", "10.0.0.1, 1.2.3.4")
2823 .body(Body::empty())
2824 .unwrap();
2825 set_forwarded_headers(&mut request, "203.0.113.7".parse().unwrap());
2826 let values: Vec<_> = request
2827 .headers()
2828 .get_all("x-forwarded-for")
2829 .iter()
2830 .collect();
2831 assert_eq!(values.len(), 1);
2832 assert_eq!(values[0], "203.0.113.7");
2833 }
2834
2835 #[test]
2836 fn forwarded_proto_preserves_upstream_tls() {
2837 let mut request = Request::builder()
2839 .uri("/")
2840 .header(header::HOST, "example.com")
2841 .header("x-forwarded-proto", "https")
2842 .body(Body::empty())
2843 .unwrap();
2844 set_forwarded_headers(&mut request, "203.0.113.7".parse().unwrap());
2845 assert_eq!(request.headers().get("x-forwarded-proto").unwrap(), "https");
2846 }
2847
2848 #[test]
2849 fn forwarded_host_absent_when_no_host_header() {
2850 let mut request = Request::builder().uri("/").body(Body::empty()).unwrap();
2851 set_forwarded_headers(&mut request, "203.0.113.7".parse().unwrap());
2852 assert!(request.headers().get("x-forwarded-host").is_none());
2853 assert_eq!(
2854 request.headers().get("x-forwarded-for").unwrap(),
2855 "203.0.113.7"
2856 );
2857 }
2858
2859 use boatramp_core::kv::{KvStore, MemoryKv};
2862 use boatramp_core::messaging::{LogMessaging, Messaging};
2863 use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, StorageError};
2864
2865 const EVENT_CONSUMER: &[u8] =
2866 include_bytes!("../../boatramp-handlers/tests/fixtures/event-consumer.wasm");
2867
2868 #[derive(Default)]
2869 struct MemStorage {
2870 objects: std::sync::Mutex<std::collections::HashMap<String, Vec<u8>>>,
2871 }
2872
2873 #[async_trait::async_trait]
2874 impl boatramp_core::Storage for MemStorage {
2875 async fn get(&self, key: &str) -> Result<GetObject, StorageError> {
2876 let bytes = self
2877 .objects
2878 .lock()
2879 .unwrap()
2880 .get(key)
2881 .cloned()
2882 .ok_or_else(|| StorageError::NotFound(key.to_string()))?;
2883 let body: ByteStream =
2884 futures::stream::once(async move { Ok(bytes::Bytes::from(bytes)) }).boxed();
2885 Ok(GetObject {
2886 meta: ObjectMeta {
2887 key: key.to_string(),
2888 ..Default::default()
2889 },
2890 body,
2891 })
2892 }
2893 async fn get_range(
2894 &self,
2895 key: &str,
2896 _: u64,
2897 _: Option<u64>,
2898 ) -> Result<GetObject, StorageError> {
2899 self.get(key).await
2900 }
2901 async fn put(
2902 &self,
2903 key: &str,
2904 mut body: ByteStream,
2905 _: PutMeta,
2906 ) -> Result<ObjectMeta, StorageError> {
2907 use futures::StreamExt;
2908 let mut buf = Vec::new();
2909 while let Some(chunk) = body.next().await {
2910 buf.extend_from_slice(&chunk?);
2911 }
2912 self.objects.lock().unwrap().insert(key.to_string(), buf);
2913 Ok(ObjectMeta {
2914 key: key.to_string(),
2915 ..Default::default()
2916 })
2917 }
2918 async fn head(&self, key: &str) -> Result<ObjectMeta, StorageError> {
2919 self.objects
2920 .lock()
2921 .unwrap()
2922 .get(key)
2923 .map(|_| ObjectMeta {
2924 key: key.to_string(),
2925 ..Default::default()
2926 })
2927 .ok_or_else(|| StorageError::NotFound(key.to_string()))
2928 }
2929 async fn delete(&self, key: &str) -> Result<(), StorageError> {
2930 self.objects.lock().unwrap().remove(key);
2931 Ok(())
2932 }
2933 async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
2934 Ok(Vec::new())
2935 }
2936 }
2937
2938 fn observed_state_in(
2941 project: &str,
2942 workload: &str,
2943 host: &str,
2944 healthy: bool,
2945 phase: boatramp_core::compute::ReplicaPhase,
2946 ) -> boatramp_core::compute::ObservedInstance {
2947 use boatramp_core::compute::{Endpoint, InstanceHandle, ReplicaPhase, Scheme, Snapshot};
2948 boatramp_core::compute::ObservedInstance {
2949 handle: InstanceHandle {
2950 project: project.into(),
2951 workload: workload.into(),
2952 replica: 0,
2953 backend_ref: "ref-0".into(),
2954 },
2955 node: 1,
2956 backend: "vmm".into(),
2957 endpoint: Endpoint {
2958 scheme: Scheme::Http,
2959 host: host.into(),
2960 port: 80,
2961 },
2962 region: None,
2963 healthy,
2964 started_at: None,
2965 phase,
2966 snapshot: matches!(phase, ReplicaPhase::Zero).then(|| Snapshot {
2967 project: project.into(),
2968 workload: workload.into(),
2969 replica: 0,
2970 data_ref: "snap-0".into(),
2971 }),
2972 }
2973 }
2974
2975 fn observed_state(
2977 workload: &str,
2978 healthy: bool,
2979 phase: boatramp_core::compute::ReplicaPhase,
2980 ) -> boatramp_core::compute::ObservedInstance {
2981 observed_state_in("default", workload, "10.0.0.2", healthy, phase)
2982 }
2983
2984 #[tokio::test]
2985 async fn has_parked_replica_detects_a_zeroed_replica() {
2986 use boatramp_core::compute::ReplicaPhase;
2987 let storage = Arc::new(MemStorage::default());
2988 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
2989 let deploy = DeployStore::new(storage, kv);
2990
2991 assert!(!has_parked_replica(&deploy, "default", "w").await);
2993 deploy
2995 .set_replica_state(
2996 ProjectRef::DEFAULT,
2997 &observed_state("w", true, ReplicaPhase::Running),
2998 )
2999 .await
3000 .unwrap();
3001 assert!(!has_parked_replica(&deploy, "default", "w").await);
3002 deploy
3004 .set_replica_state(
3005 ProjectRef::DEFAULT,
3006 &observed_state("w", false, ReplicaPhase::Zero),
3007 )
3008 .await
3009 .unwrap();
3010 assert!(has_parked_replica(&deploy, "default", "w").await);
3011 }
3012
3013 #[tokio::test]
3014 async fn await_warm_returns_immediately_when_healthy_and_times_out_otherwise() {
3015 use boatramp_core::compute::ReplicaPhase;
3016 let storage = Arc::new(MemStorage::default());
3017 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3018 let deploy = DeployStore::new(storage, kv);
3019
3020 let empty = await_warm(
3022 &deploy,
3023 "default",
3024 "w",
3025 std::time::Duration::from_millis(150),
3026 )
3027 .await;
3028 assert!(empty.is_empty());
3029
3030 deploy
3032 .set_replica_state(
3033 ProjectRef::DEFAULT,
3034 &observed_state("w", true, ReplicaPhase::Running),
3035 )
3036 .await
3037 .unwrap();
3038 let warm = await_warm(&deploy, "default", "w", std::time::Duration::from_secs(5)).await;
3039 assert_eq!(warm, vec!["http://10.0.0.2:80".to_string()]);
3040 }
3041
3042 #[tokio::test]
3048 async fn compute_endpoints_are_project_scoped() {
3049 use boatramp_core::compute::ReplicaPhase;
3050 let storage = Arc::new(MemStorage::default());
3051 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3052 let deploy = DeployStore::new(storage, kv);
3053
3054 deploy
3056 .set_replica_state(
3057 ProjectRef::new("acme"),
3058 &observed_state_in("acme", "web", "10.0.0.5", true, ReplicaPhase::Running),
3059 )
3060 .await
3061 .unwrap();
3062 deploy
3063 .set_replica_state(
3064 ProjectRef::DEFAULT,
3065 &observed_state_in("default", "web", "10.0.0.9", true, ReplicaPhase::Running),
3066 )
3067 .await
3068 .unwrap();
3069
3070 assert_eq!(
3072 compute_endpoints(&deploy, "acme", "web").await,
3073 vec!["http://10.0.0.5:80".to_string()],
3074 "acme's web resolves against acme, not default"
3075 );
3076 assert_eq!(
3078 compute_endpoints(&deploy, "default", "web").await,
3079 vec!["http://10.0.0.9:80".to_string()]
3080 );
3081 assert!(compute_endpoints(&deploy, "beta", "web").await.is_empty());
3084 }
3085
3086 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3090 async fn dispatcher_delivers_at_least_once_then_dead_letters() {
3091 use boatramp_handlers::{Bindings, HandlerEngine, Limits};
3092 let storage = Arc::new(MemStorage::default());
3093 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3094 let mq = LogMessaging::new(storage, kv.clone());
3095 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3096 let hash = boatramp_core::deploy::sha256_hex(EVENT_CONSUMER);
3097 let bindings = Bindings::new("blog").with_keyvalue("blog", kv.clone());
3098 let topic = "blog/orders/created";
3099
3100 for _ in 0..3 {
3102 mq.publish(topic, b"ok").await.unwrap();
3103 }
3104 loop {
3105 let acked = dispatch_consumer_batch(
3106 &engine,
3107 &mq,
3108 &metrics::Metrics::default(),
3109 "blog",
3110 topic,
3111 "blog/",
3112 "",
3113 boatramp_core::messaging::StartPosition::Latest,
3114 &hash,
3115 EVENT_CONSUMER,
3116 &bindings,
3117 None,
3119 Limits::default(),
3120 Duration::from_secs(30),
3121 5,
3122 10,
3123 )
3124 .await;
3125 if acked == 0 {
3126 break;
3127 }
3128 }
3129 assert_eq!(
3130 kv.get("hkv/blog/delivered/orders/created").await.unwrap(),
3131 Some(b"3".to_vec())
3132 );
3133
3134 mq.publish(topic, b"fail").await.unwrap();
3137 for _ in 0..5 {
3138 dispatch_consumer_batch(
3139 &engine,
3140 &mq,
3141 &metrics::Metrics::default(),
3142 "blog",
3143 topic,
3144 "blog/",
3145 "",
3146 boatramp_core::messaging::StartPosition::Latest,
3147 &hash,
3148 EVENT_CONSUMER,
3149 &bindings,
3150 None,
3152 Limits::default(),
3153 Duration::ZERO,
3154 2,
3155 10,
3156 )
3157 .await;
3158 }
3159 assert_eq!(mq.dead_letter_count(topic).await.unwrap(), 1);
3160 assert_eq!(
3162 kv.get("hkv/blog/delivered/orders/created").await.unwrap(),
3163 Some(b"3".to_vec())
3164 );
3165 }
3166
3167 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3172 async fn consumer_groups_fan_out_through_the_dispatcher() {
3173 use boatramp_handlers::{Bindings, HandlerEngine, Limits};
3174 let storage = Arc::new(MemStorage::default());
3175 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3176 let mq = LogMessaging::new(storage, kv.clone());
3177 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3178 let hash = boatramp_core::deploy::sha256_hex(EVENT_CONSUMER);
3179 let bindings = Bindings::new("blog").with_keyvalue("blog", kv.clone());
3180 let topic = "blog/orders/created";
3181 let start = boatramp_core::messaging::StartPosition::Latest;
3182
3183 for g in ["billing", "audit"] {
3186 let n = dispatch_consumer_batch(
3187 &engine,
3188 &mq,
3189 &metrics::Metrics::default(),
3190 "blog",
3191 topic,
3192 "blog/",
3193 g,
3194 start,
3195 &hash,
3196 EVENT_CONSUMER,
3197 &bindings,
3198 None,
3200 Limits::default(),
3201 Duration::from_secs(30),
3202 5,
3203 10,
3204 )
3205 .await;
3206 assert_eq!(n, 0, "no events yet for group {g}");
3207 }
3208 mq.publish(topic, b"ok").await.unwrap();
3209
3210 for g in ["billing", "audit"] {
3212 let n = dispatch_consumer_batch(
3213 &engine,
3214 &mq,
3215 &metrics::Metrics::default(),
3216 "blog",
3217 topic,
3218 "blog/",
3219 g,
3220 start,
3221 &hash,
3222 EVENT_CONSUMER,
3223 &bindings,
3224 None,
3226 Limits::default(),
3227 Duration::from_secs(30),
3228 5,
3229 10,
3230 )
3231 .await;
3232 assert_eq!(n, 1, "group {g} should receive the message");
3233 }
3234 assert_eq!(
3236 kv.get("hkv/blog/delivered/orders/created").await.unwrap(),
3237 Some(b"2".to_vec())
3238 );
3239 }
3240
3241 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3245 async fn scheduler_runs_current_consumers_not_previews() {
3246 use boatramp_core::config::{ConsumerConfig, DeployConfig, HandlersSiteConfig, SiteConfig};
3247 use boatramp_core::deploy::{DeployStore, FileEntry, Manifest};
3248 use boatramp_handlers::{HandlerEngine, Limits};
3249 use futures::StreamExt;
3250
3251 let storage = Arc::new(MemStorage::default());
3252 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3253 let deploy = DeployStore::new(storage.clone(), kv.clone());
3254 let messaging: Arc<dyn Messaging> =
3255 Arc::new(LogMessaging::new(storage.clone(), kv.clone()));
3256
3257 let hash = boatramp_core::deploy::sha256_hex(EVENT_CONSUMER);
3259 let stream: ByteStream =
3260 futures::stream::once(async move { Ok(bytes::Bytes::from_static(EVENT_CONSUMER)) })
3261 .boxed();
3262 deploy.put_blob(&hash, stream).await.unwrap();
3263 let mut files = std::collections::BTreeMap::new();
3264 files.insert(
3265 "consumer.wasm".to_string(),
3266 FileEntry {
3267 hash: hash.clone(),
3268 size: EVENT_CONSUMER.len() as u64,
3269 content_type: None,
3270 variants: std::collections::BTreeMap::new(),
3271 },
3272 );
3273 let manifest = Manifest {
3274 files,
3275 config: DeployConfig {
3276 consumers: vec![ConsumerConfig {
3277 tenancy: None,
3278 token_claims: None,
3279 topic: "orders/created".into(),
3280 component: "consumer.wasm".into(),
3281 imports: vec!["wasi:keyvalue".into()],
3282 group: String::new(),
3283 start: Default::default(),
3284 }],
3285 ..Default::default()
3286 },
3287 ..Default::default()
3288 };
3289 let id = deploy.put_manifest(&manifest).await.unwrap();
3290 deploy
3291 .activate(ProjectRef::DEFAULT, "blog", &id)
3292 .await
3293 .unwrap();
3294 deploy
3295 .set_site_config(
3296 ProjectRef::DEFAULT,
3297 "blog",
3298 &SiteConfig {
3299 handlers: Some(HandlersSiteConfig {
3300 enabled: true,
3301 allow_imports: vec!["wasi:keyvalue".into()],
3302 ..Default::default()
3303 }),
3304 ..Default::default()
3305 },
3306 )
3307 .await
3308 .unwrap();
3309
3310 messaging
3312 .publish("blog/orders/created", b"live")
3313 .await
3314 .unwrap();
3315 messaging
3316 .publish("blog/_preview/abc/orders/created", b"preview")
3317 .await
3318 .unwrap();
3319
3320 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3321 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, Some(messaging));
3322 let inner = rt.inner.clone().unwrap();
3323 let mut cache = std::collections::HashMap::new();
3324 let mut crons = std::collections::HashMap::new();
3325 let mut sweep = std::collections::HashMap::new();
3326 let now = CronNow {
3327 minute: 0,
3328 hour: 0,
3329 dom: 1,
3330 month: 1,
3331 dow: 0,
3332 minute_stamp: 0,
3333 };
3334 for _ in 0..3 {
3335 run_scheduler_tick(&inner, &deploy, &mut cache, &mut crons, &mut sweep, now)
3336 .await
3337 .unwrap();
3338 }
3339
3340 assert_eq!(
3342 kv.get("hkv/blog/delivered/orders/created").await.unwrap(),
3343 Some(b"1".to_vec())
3344 );
3345 assert_eq!(
3348 kv.get("hkv/blog/_preview/abc/delivered/orders/created")
3349 .await
3350 .unwrap(),
3351 None
3352 );
3353 }
3354
3355 const KV_COUNTER: &[u8] =
3360 include_bytes!("../../boatramp-handlers/tests/fixtures/kv-counter.wasm");
3361
3362 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3367 async fn function_invoker_runs_target_buffers_and_meters() {
3368 use boatramp_core::deploy::DeployStore;
3369 use boatramp_core::function::{Function, FunctionVersion, Lifecycle, Owner};
3370 use boatramp_handlers::{HandlerEngine, InvokeError, InvokeRequest, Invoker, Limits};
3371 use futures::StreamExt;
3372
3373 const HTTP_200: &[u8] =
3376 include_bytes!("../../boatramp-handlers/tests/fixtures/http-200.wasm");
3377
3378 let storage = Arc::new(MemStorage::default());
3379 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3380 let deploy = DeployStore::new(storage.clone(), kv.clone());
3381
3382 let hash = boatramp_core::deploy::sha256_hex(HTTP_200);
3383 let stream: ByteStream =
3384 futures::stream::once(async move { Ok(bytes::Bytes::from_static(HTTP_200)) }).boxed();
3385 deploy.put_blob(&hash, stream).await.unwrap();
3386 let function = Function {
3387 name: "target".into(),
3388 owner: Owner::Project("default".into()),
3389 versions: vec![FunctionVersion {
3390 id: "v1".into(),
3391 component: hash.clone(),
3392 created: 0,
3393 lifecycle: Lifecycle::Independent,
3394 }],
3395 active: "v1".into(),
3396 aliases: Default::default(),
3397 config: Default::default(),
3398 };
3399 deploy
3400 .put_function(ProjectRef::DEFAULT, &function)
3401 .await
3402 .unwrap();
3403
3404 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3405 let rt = HandlerRuntime::new(engine, kv, storage, None, None);
3406 rt.set_invoker(deploy.clone());
3407 let invoker = rt.inner.as_ref().unwrap().invoker.get().unwrap().clone();
3408
3409 let request = || InvokeRequest {
3410 method: "GET".into(),
3411 path: "/".into(),
3412 headers: vec![],
3413 body: vec![],
3414 };
3415
3416 let response = invoker.invoke("target", request(), 1).await.unwrap();
3418 assert_eq!(response.status, 200);
3419
3420 let metering = deploy
3422 .get_metering(ProjectRef::DEFAULT, "target")
3423 .await
3424 .unwrap()
3425 .unwrap();
3426 assert_eq!(metering.invocations, 1);
3427
3428 let err = invoker.invoke("ghost", request(), 1).await.unwrap_err();
3430 assert!(matches!(err, InvokeError::NotFound));
3431 }
3432
3433 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3439 async fn federation_runner_enforces_the_safelist_before_planning() {
3440 use boatramp_core::deploy::DeployStore;
3441 use boatramp_core::project::ProjectRef;
3442 use boatramp_handlers::{GraphqlRequest, HandlerEngine, Limits, SupergraphRunError};
3443
3444 let storage = Arc::new(MemStorage::default());
3445 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3446 let deploy = DeployStore::new(storage.clone(), kv.clone());
3447 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3448 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
3449 rt.set_invoker(deploy.clone());
3450 let runner = rt
3451 .inner
3452 .as_ref()
3453 .unwrap()
3454 .federation_runner
3455 .get()
3456 .unwrap()
3457 .scoped(ProjectRef::new("default"), Vec::new());
3458
3459 let req = |query: &str| GraphqlRequest {
3460 query: Some(query.to_string()),
3461 persisted_hash: None,
3462 variables: "{}".to_string(),
3463 operation_name: None,
3464 authorization: None,
3465 };
3466
3467 assert!(matches!(
3469 runner.run(req("{ me { id } }"), 1).await,
3470 Err(SupergraphRunError::NotSafelisted)
3471 ));
3472
3473 let query = "{ me { id } }";
3476 let hash = crate::graphql_apq::sha256_hex(query);
3477 kv.put(&format!("hapq/default/{hash}"), query.as_bytes().to_vec())
3478 .await
3479 .unwrap();
3480 assert!(matches!(
3481 runner.run(req(query), 1).await,
3482 Err(SupergraphRunError::PlanFailed(_))
3483 ));
3484
3485 let persisted = GraphqlRequest {
3487 query: None,
3488 persisted_hash: Some("deadbeef".to_string()),
3489 variables: "{}".to_string(),
3490 operation_name: None,
3491 authorization: None,
3492 };
3493 assert!(matches!(
3494 runner.run(persisted, 1).await,
3495 Err(SupergraphRunError::NotSafelisted)
3496 ));
3497 }
3498
3499 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3508 #[ignore = "run on the host toolchain (real libsql static-musl segfault); wired in the CI graphql-propagation gate"]
3509 async fn graphql_run_propagates_caller_principal_to_a_scoped_subfetch() {
3510 use boatramp_core::deploy::{sha256_hex, DeployStore};
3511 use boatramp_core::function::{
3512 Function, FunctionConfig, FunctionVersion, Lifecycle, Owner,
3513 };
3514 use boatramp_core::project::ProjectRef;
3515 use boatramp_core::sql::{SqlBackends, SqlValue};
3516 use boatramp_core::tenancy::{AccessMode, ScopeAxis, Tenancy, TenantSource};
3517 use boatramp_handlers::{GraphqlRequest, HandlerEngine, Limits, ScopeFact};
3518
3519 const PROBE: &[u8] = include_bytes!("../tests/fixtures/graphql-scope-probe.wasm");
3523
3524 let storage = Arc::new(MemStorage::default());
3525 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3526 let deploy = DeployStore::new(storage.clone(), kv.clone());
3527
3528 let sql_dir = std::env::temp_dir().join(format!("br-gqlprop-{}", std::process::id()));
3531 let _ = std::fs::remove_dir_all(&sql_dir);
3532 let backends = boatramp_storage::LibsqlSqlBackends::local(&sql_dir);
3533 let db = backends
3534 .database("default", "fn/scopeprobe", "")
3535 .await
3536 .unwrap();
3537 {
3538 let mut tx = db.begin().await.unwrap();
3539 tx.execute(
3540 "CREATE TABLE items (id TEXT PRIMARY KEY, tenant_id TEXT)",
3541 &[],
3542 )
3543 .await
3544 .unwrap();
3545 for (id, tenant) in [("a1", "tenant_A"), ("b1", "tenant_B"), ("b2", "tenant_B")] {
3546 tx.execute(
3547 "INSERT INTO items (id, tenant_id) VALUES (?1, ?2)",
3548 &[SqlValue::Text(id.into()), SqlValue::Text(tenant.into())],
3549 )
3550 .await
3551 .unwrap();
3552 }
3553 tx.commit().await.unwrap();
3554 }
3555 let sql: Arc<dyn SqlBackends> = Arc::new(backends);
3556
3557 let hash = sha256_hex(PROBE);
3560 let stream: ByteStream =
3561 futures::stream::once(async move { Ok(bytes::Bytes::from_static(PROBE)) }).boxed();
3562 deploy.put_blob(&hash, stream).await.unwrap();
3563 let function = Function {
3564 name: "scopeprobe".into(),
3565 owner: Owner::Project("default".into()),
3566 versions: vec![FunctionVersion {
3567 id: "v1".into(),
3568 component: hash.clone(),
3569 created: 0,
3570 lifecycle: Lifecycle::Independent,
3571 }],
3572 active: "v1".into(),
3573 aliases: Default::default(),
3574 config: FunctionConfig {
3575 imports: vec!["sql".into()],
3576 tenancy: Some(Tenancy::Scoped {
3577 column: "tenant_id".into(),
3578 sources: vec![TenantSource::None],
3579 read: AccessMode::Own,
3580 write: AccessMode::None,
3581 }),
3582 ..Default::default()
3583 },
3584 };
3585 deploy
3586 .put_function(ProjectRef::DEFAULT, &function)
3587 .await
3588 .unwrap();
3589
3590 crate::graphql_registry::publish(
3592 kv.as_ref(),
3593 "default",
3594 "scopeprobe",
3595 "type Query { items: [Item!]! }\ntype Item @key(fields: \"id\") { id: ID! }",
3596 )
3597 .await
3598 .unwrap();
3599 let query = "{ items { id } }";
3600 let op_hash = crate::graphql_apq::sha256_hex(query);
3601 kv.put(
3602 &format!("hapq/default/{op_hash}"),
3603 query.as_bytes().to_vec(),
3604 )
3605 .await
3606 .unwrap();
3607
3608 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3609 let rt = HandlerRuntime::new(engine, kv.clone(), storage, Some(sql), None);
3610 rt.set_invoker(deploy.clone());
3611 let fed = rt
3612 .inner
3613 .as_ref()
3614 .unwrap()
3615 .federation_runner
3616 .get()
3617 .unwrap()
3618 .clone();
3619 let req = || GraphqlRequest {
3620 query: Some(query.to_string()),
3621 persisted_hash: None,
3622 variables: "{}".to_string(),
3623 operation_name: None,
3624 authorization: None,
3625 };
3626
3627 let runner_b = fed.scoped(
3629 ProjectRef::new("default"),
3630 vec![ScopeFact {
3631 axis: ScopeAxis::Tenant,
3632 value: SqlValue::Text("tenant_B".into()),
3633 }],
3634 );
3635 let body_b = String::from_utf8_lossy(&runner_b.run(req(), 0).await.unwrap()).into_owned();
3636 assert!(
3637 body_b.contains("\"b1\"") && body_b.contains("\"b2\"") && !body_b.contains("\"a1\""),
3638 "graphql::run propagated principal B → the subgraph read ONLY tenant B's rows: {body_b}"
3639 );
3640
3641 let runner_empty = fed.scoped(ProjectRef::new("default"), Vec::new());
3644 let body_none =
3645 String::from_utf8_lossy(&runner_empty.run(req(), 0).await.unwrap()).into_owned();
3646 assert!(
3647 !body_none.contains("\"a1\"")
3648 && !body_none.contains("\"b1\"")
3649 && !body_none.contains("\"b2\""),
3650 "with NO propagated principal the subgraph fails closed — no rows leak: {body_none}"
3651 );
3652
3653 let _ = std::fs::remove_dir_all(&sql_dir);
3654 println!(
3655 "GRAPHQL-RUN PRINCIPAL PROPAGATION OK: graphql::run carried the caller's resolved principal \
3656 (tenant B) to a federated subgraph sub-fetch, which resolved its own-tenancy from the \
3657 inherited principal and returned ONLY tenant B's rows over a real libsql engine; with no \
3658 principal the same sub-fetch failed closed (no rows) — the async lane can drive a \
3659 tenant-scoped supergraph read/write with no request bearer, symmetric to emit::invoke"
3660 );
3661 }
3662
3663 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3672 #[ignore = "run on the host toolchain (real libsql static-musl segfault); wired in the CI gateway-target gate"]
3673 async fn gateway_forces_target_confinement_onto_a_wasm_subgraph() {
3674 use boatramp_core::deploy::{sha256_hex, DeployStore};
3675 use boatramp_core::function::{
3676 Function, FunctionConfig, FunctionVersion, Lifecycle, Owner,
3677 };
3678 use boatramp_core::project::ProjectRef;
3679 use boatramp_core::sql::{SqlBackends, SqlValue};
3680 use boatramp_core::tenancy::{
3681 AccessMode, PublicPredicate, PublicSubset, PublicTerm, TableScope, Tenancy,
3682 TenancySchema, TenantSource,
3683 };
3684 use boatramp_handlers::{HandlerEngine, Limits};
3685
3686 const PROBE: &[u8] = include_bytes!("../tests/fixtures/graphql-scope-probe.wasm");
3692
3693 let storage = Arc::new(MemStorage::default());
3694 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3695 let deploy = DeployStore::new(storage.clone(), kv.clone());
3696
3697 let sql_dir = std::env::temp_dir().join(format!("br-gwtarget-{}", std::process::id()));
3700 let _ = std::fs::remove_dir_all(&sql_dir);
3701 let backends = boatramp_storage::LibsqlSqlBackends::local(&sql_dir);
3702 let db = backends
3703 .database("default", "fn/scopeprobe", "")
3704 .await
3705 .unwrap();
3706 {
3707 let mut tx = db.begin().await.unwrap();
3708 tx.execute(
3709 "CREATE TABLE items (id TEXT PRIMARY KEY, tenant_id TEXT, published INTEGER)",
3710 &[],
3711 )
3712 .await
3713 .unwrap();
3714 for (id, tenant, published) in [
3715 ("a_pub", "tenant_A", 1),
3716 ("b_pub", "tenant_B", 1),
3717 ("b_priv", "tenant_B", 0),
3718 ] {
3719 tx.execute(
3720 "INSERT INTO items (id, tenant_id, published) VALUES (?1, ?2, ?3)",
3721 &[
3722 SqlValue::Text(id.into()),
3723 SqlValue::Text(tenant.into()),
3724 SqlValue::Integer(published),
3725 ],
3726 )
3727 .await
3728 .unwrap();
3729 }
3730 tx.commit().await.unwrap();
3731 }
3732 let sql: Arc<dyn SqlBackends> = Arc::new(backends);
3733
3734 let hash = sha256_hex(PROBE);
3737 let stream: ByteStream =
3738 futures::stream::once(async move { Ok(bytes::Bytes::from_static(PROBE)) }).boxed();
3739 deploy.put_blob(&hash, stream).await.unwrap();
3740 let function = Function {
3741 name: "scopeprobe".into(),
3742 owner: Owner::Project("default".into()),
3743 versions: vec![FunctionVersion {
3744 id: "v1".into(),
3745 component: hash.clone(),
3746 created: 0,
3747 lifecycle: Lifecycle::Independent,
3748 }],
3749 active: "v1".into(),
3750 aliases: Default::default(),
3751 config: FunctionConfig {
3752 imports: vec!["sql".into()],
3753 tenancy: Some(Tenancy::Scoped {
3754 column: "tenant_id".into(),
3755 sources: vec![TenantSource::None],
3756 read: AccessMode::Own,
3757 write: AccessMode::None,
3758 }),
3759 ..Default::default()
3760 },
3761 };
3762 deploy
3763 .put_function(ProjectRef::DEFAULT, &function)
3764 .await
3765 .unwrap();
3766
3767 let sdl = "type Query { items: [Item!]! @tenant(scope: target, via: [domain], public: \"items\") }\n\
3770 type Item @key(fields: \"id\") { id: ID! }";
3771 let sg = crate::graphql_federation::compose(&[("scopeprobe".into(), sdl.into())]).unwrap();
3772 let plan = crate::graphql_plan::plan("{ items { id } }", &sg).unwrap();
3773
3774 let mut schema = TenancySchema {
3775 default_tenant_key: "tenant_id".into(),
3776 tables: std::collections::BTreeMap::from([("items".into(), TableScope::Tenant)]),
3777 ..Default::default()
3778 };
3779 schema.public_subsets.insert(
3780 "items".into(),
3781 PublicSubset {
3782 predicate: PublicPredicate {
3783 terms: vec![PublicTerm::Cmp {
3784 column: "published".into(),
3785 op: boatramp_core::tenancy::PublicCmp::Eq,
3786 value: boatramp_core::tenancy::PublicLiteral::Int(1),
3787 }],
3788 },
3789 world_public: true,
3790 listable: true,
3791 },
3792 );
3793
3794 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3795 let rt = HandlerRuntime::new(engine, kv.clone(), storage, Some(sql), None);
3796 rt.set_invoker(deploy.clone());
3797 let inner = rt.inner.as_ref().unwrap();
3798 let invoker = inner.invoker.get().unwrap().clone();
3799
3800 let make_router = |domain: Option<&str>| {
3801 crate::graphql_gateway::BackendRouter::new(
3802 invoker.scoped(ProjectRef::new("default"), Vec::new()),
3803 "default".to_string(),
3804 inner.sql.clone(),
3805 std::collections::BTreeMap::new(),
3806 None,
3807 )
3808 .with_target_inputs(Some(crate::graphql_gateway::TargetInputs {
3809 schema: Arc::new(schema.clone()),
3810 domain_context: domain.map(str::to_string),
3811 target_handle: None,
3812 capability_anchor: None,
3813 }))
3814 };
3815
3816 let router_b = make_router(Some("tenant_B"));
3818 let out_b = crate::graphql_gateway::execute(&plan, &router_b, &serde_json::json!({})).await;
3819 let body_b = out_b.to_string();
3820 assert!(
3821 body_b.contains("\"b_pub\"")
3822 && !body_b.contains("\"b_priv\"")
3823 && !body_b.contains("\"a_pub\""),
3824 "gateway forced target B onto the wasm subgraph → ONLY B's PUBLIC row (b_pub), never B's \
3825 private row nor A's: {body_b}"
3826 );
3827
3828 let router_none = make_router(None);
3830 let out_none =
3831 crate::graphql_gateway::execute(&plan, &router_none, &serde_json::json!({})).await;
3832 let body_none = out_none.to_string();
3833 assert!(
3834 !body_none.contains("\"a_pub\"")
3835 && !body_none.contains("\"b_pub\"")
3836 && !body_none.contains("\"b_priv\""),
3837 "with no resolved target the wasm-subgraph target fetch fails closed — no rows: {body_none}"
3838 );
3839
3840 let _ = std::fs::remove_dir_all(&sql_dir);
3841 println!(
3842 "GATEWAY WASM-TARGET OK: the /graphql gateway resolved target tenant B from the routed \
3843 domain and FORCED a HostTenancy::target binding onto the wasm subgraph invocation \
3844 (invoke_target), confining the guest's own SQL to tenant=B AND published=1 over a real \
3845 libsql engine — returned ONLY B's public row (b_pub), never B's private row (b_priv) nor \
3846 tenant A's (a_pub); with no resolved target the fetch failed closed"
3847 );
3848 }
3849
3850 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
3860 #[ignore = "run on the host toolchain (real libsql static-musl segfault); wired in the CI gateway-own-propagation gate"]
3861 async fn gateway_propagates_caller_own_principal_to_a_wasm_subgraph() {
3862 use boatramp_core::deploy::{sha256_hex, DeployStore};
3863 use boatramp_core::function::{
3864 Function, FunctionConfig, FunctionVersion, Lifecycle, Owner,
3865 };
3866 use boatramp_core::project::ProjectRef;
3867 use boatramp_core::sql::{SqlBackends, SqlValue};
3868 use boatramp_core::tenancy::{AccessMode, ScopeAxis, Tenancy, TenantSource};
3869 use boatramp_handlers::{HandlerEngine, Limits, ScopeFact};
3870
3871 const PROBE: &[u8] = include_bytes!("../tests/fixtures/graphql-scope-probe.wasm");
3875
3876 let storage = Arc::new(MemStorage::default());
3877 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
3878 let deploy = DeployStore::new(storage.clone(), kv.clone());
3879
3880 let sql_dir = std::env::temp_dir().join(format!("br-gwown-{}", std::process::id()));
3881 let _ = std::fs::remove_dir_all(&sql_dir);
3882 let backends = boatramp_storage::LibsqlSqlBackends::local(&sql_dir);
3883 let db = backends
3884 .database("default", "fn/scopeprobe", "")
3885 .await
3886 .unwrap();
3887 {
3888 let mut tx = db.begin().await.unwrap();
3889 tx.execute(
3890 "CREATE TABLE items (id TEXT PRIMARY KEY, tenant_id TEXT)",
3891 &[],
3892 )
3893 .await
3894 .unwrap();
3895 for (id, tenant) in [("a1", "tenant_A"), ("b1", "tenant_B"), ("b2", "tenant_B")] {
3896 tx.execute(
3897 "INSERT INTO items (id, tenant_id) VALUES (?1, ?2)",
3898 &[SqlValue::Text(id.into()), SqlValue::Text(tenant.into())],
3899 )
3900 .await
3901 .unwrap();
3902 }
3903 tx.commit().await.unwrap();
3904 }
3905 let sql: Arc<dyn SqlBackends> = Arc::new(backends);
3906
3907 let hash = sha256_hex(PROBE);
3910 let stream: ByteStream =
3911 futures::stream::once(async move { Ok(bytes::Bytes::from_static(PROBE)) }).boxed();
3912 deploy.put_blob(&hash, stream).await.unwrap();
3913 let function = Function {
3914 name: "scopeprobe".into(),
3915 owner: Owner::Project("default".into()),
3916 versions: vec![FunctionVersion {
3917 id: "v1".into(),
3918 component: hash.clone(),
3919 created: 0,
3920 lifecycle: Lifecycle::Independent,
3921 }],
3922 active: "v1".into(),
3923 aliases: Default::default(),
3924 config: FunctionConfig {
3925 imports: vec!["sql".into()],
3926 tenancy: Some(Tenancy::Scoped {
3927 column: "tenant_id".into(),
3928 sources: vec![TenantSource::None],
3929 read: AccessMode::Own,
3930 write: AccessMode::None,
3931 }),
3932 ..Default::default()
3933 },
3934 };
3935 deploy
3936 .put_function(ProjectRef::DEFAULT, &function)
3937 .await
3938 .unwrap();
3939
3940 let sdl = "type Query { items: [Item!]! }\ntype Item @key(fields: \"id\") { id: ID! }";
3943 let sg = crate::graphql_federation::compose(&[("scopeprobe".into(), sdl.into())]).unwrap();
3944 let plan = crate::graphql_plan::plan("{ items { id } }", &sg).unwrap();
3945
3946 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
3947 let rt = HandlerRuntime::new(engine, kv.clone(), storage, Some(sql), None);
3948 rt.set_invoker(deploy.clone());
3949 let inner = rt.inner.as_ref().unwrap();
3950 let invoker = inner.invoker.get().unwrap().clone();
3951
3952 let make_router = |facts: Vec<ScopeFact>| {
3955 crate::graphql_gateway::BackendRouter::new(
3956 invoker.scoped(ProjectRef::new("default"), facts),
3957 "default".to_string(),
3958 inner.sql.clone(),
3959 std::collections::BTreeMap::new(),
3960 None,
3961 )
3962 };
3963
3964 let router_b = make_router(vec![ScopeFact {
3966 axis: ScopeAxis::Tenant,
3967 value: SqlValue::Text("tenant_B".into()),
3968 }]);
3969 let body_b = crate::graphql_gateway::execute(&plan, &router_b, &serde_json::json!({}))
3970 .await
3971 .to_string();
3972 assert!(
3973 body_b.contains("\"b1\"") && body_b.contains("\"b2\"") && !body_b.contains("\"a1\""),
3974 "gateway propagated principal B → the wasm subgraph read ONLY tenant B's rows: {body_b}"
3975 );
3976
3977 let router_anon = make_router(Vec::new());
3979 let body_anon =
3980 crate::graphql_gateway::execute(&plan, &router_anon, &serde_json::json!({}))
3981 .await
3982 .to_string();
3983 assert!(
3984 !body_anon.contains("\"a1\"")
3985 && !body_anon.contains("\"b1\"")
3986 && !body_anon.contains("\"b2\""),
3987 "with an empty principal the wasm-subgraph own fetch fails closed — no rows: {body_anon}"
3988 );
3989
3990 let _ = std::fs::remove_dir_all(&sql_dir);
3991 println!(
3992 "GATEWAY OWN-PROPAGATION OK: the /graphql gateway propagated the caller's resolved OWN \
3993 principal (tenant B) to a wasm subgraph fetch, which scoped its own read to the \
3994 inherited principal and returned ONLY tenant B's rows over a real libsql engine; with \
3995 an empty principal the same own fetch failed closed (no rows) — anon is never widened"
3996 );
3997 }
3998
3999 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4005 #[ignore = "run on the host toolchain (real libsql static-musl segfault); wired in the CI gateway-target gate"]
4006 async fn gateway_target_or_null_includes_the_shared_base() {
4007 use boatramp_core::deploy::{sha256_hex, DeployStore};
4008 use boatramp_core::function::{
4009 Function, FunctionConfig, FunctionVersion, Lifecycle, Owner,
4010 };
4011 use boatramp_core::project::ProjectRef;
4012 use boatramp_core::sql::{SqlBackends, SqlValue};
4013 use boatramp_core::tenancy::{
4014 AccessMode, PublicPredicate, PublicSubset, PublicTerm, TableScope, Tenancy,
4015 TenancySchema, TenantSource,
4016 };
4017 use boatramp_handlers::{HandlerEngine, Limits};
4018
4019 const PROBE: &[u8] = include_bytes!("../tests/fixtures/graphql-scope-probe.wasm");
4020
4021 let storage = Arc::new(MemStorage::default());
4022 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
4023 let deploy = DeployStore::new(storage.clone(), kv.clone());
4024
4025 let sql_dir = std::env::temp_dir().join(format!("br-gwton-{}", std::process::id()));
4026 let _ = std::fs::remove_dir_all(&sql_dir);
4027 let backends = boatramp_storage::LibsqlSqlBackends::local(&sql_dir);
4028 let db = backends
4029 .database("default", "fn/scopeprobe", "")
4030 .await
4031 .unwrap();
4032 {
4033 let mut tx = db.begin().await.unwrap();
4034 tx.execute(
4035 "CREATE TABLE items (id TEXT PRIMARY KEY, tenant_id TEXT, published INTEGER)",
4036 &[],
4037 )
4038 .await
4039 .unwrap();
4040 for (id, tenant, published) in [
4044 ("base_pub", None, 1),
4045 ("base_priv", None, 0),
4046 ("b_pub", Some("tenant_B"), 1),
4047 ("b_priv", Some("tenant_B"), 0),
4048 ("a_pub", Some("tenant_A"), 1),
4049 ] {
4050 tx.execute(
4051 "INSERT INTO items (id, tenant_id, published) VALUES (?1, ?2, ?3)",
4052 &[
4053 SqlValue::Text(id.into()),
4054 tenant
4055 .map(|t| SqlValue::Text(t.into()))
4056 .unwrap_or(SqlValue::Null),
4057 SqlValue::Integer(published),
4058 ],
4059 )
4060 .await
4061 .unwrap();
4062 }
4063 tx.commit().await.unwrap();
4064 }
4065 let sql: Arc<dyn SqlBackends> = Arc::new(backends);
4066
4067 let hash = sha256_hex(PROBE);
4068 let stream: ByteStream =
4069 futures::stream::once(async move { Ok(bytes::Bytes::from_static(PROBE)) }).boxed();
4070 deploy.put_blob(&hash, stream).await.unwrap();
4071 let function = Function {
4072 name: "scopeprobe".into(),
4073 owner: Owner::Project("default".into()),
4074 versions: vec![FunctionVersion {
4075 id: "v1".into(),
4076 component: hash.clone(),
4077 created: 0,
4078 lifecycle: Lifecycle::Independent,
4079 }],
4080 active: "v1".into(),
4081 aliases: Default::default(),
4082 config: FunctionConfig {
4083 imports: vec!["sql".into()],
4084 tenancy: Some(Tenancy::Scoped {
4085 column: "tenant_id".into(),
4086 sources: vec![TenantSource::None],
4087 read: AccessMode::Own,
4088 write: AccessMode::None,
4089 }),
4090 ..Default::default()
4091 },
4092 };
4093 deploy
4094 .put_function(ProjectRef::DEFAULT, &function)
4095 .await
4096 .unwrap();
4097
4098 let sdl = "type Query { items: [Item!]! @tenant(scope: target_or_null, via: [domain], public: \"items\") }\n\
4100 type Item @key(fields: \"id\") { id: ID! }";
4101 let sg = crate::graphql_federation::compose(&[("scopeprobe".into(), sdl.into())]).unwrap();
4102 let plan = crate::graphql_plan::plan("{ items { id } }", &sg).unwrap();
4103
4104 let mut schema = TenancySchema {
4105 default_tenant_key: "tenant_id".into(),
4106 tables: std::collections::BTreeMap::from([("items".into(), TableScope::Tenant)]),
4107 ..Default::default()
4108 };
4109 schema.public_subsets.insert(
4110 "items".into(),
4111 PublicSubset {
4112 predicate: PublicPredicate {
4113 terms: vec![PublicTerm::Cmp {
4114 column: "published".into(),
4115 op: boatramp_core::tenancy::PublicCmp::Eq,
4116 value: boatramp_core::tenancy::PublicLiteral::Int(1),
4117 }],
4118 },
4119 world_public: true,
4120 listable: true,
4121 },
4122 );
4123
4124 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4125 let rt = HandlerRuntime::new(engine, kv.clone(), storage, Some(sql), None);
4126 rt.set_invoker(deploy.clone());
4127 let inner = rt.inner.as_ref().unwrap();
4128 let router = crate::graphql_gateway::BackendRouter::new(
4129 inner
4130 .invoker
4131 .get()
4132 .unwrap()
4133 .clone()
4134 .scoped(ProjectRef::new("default"), Vec::new()),
4135 "default".to_string(),
4136 inner.sql.clone(),
4137 std::collections::BTreeMap::new(),
4138 None,
4139 )
4140 .with_target_inputs(Some(crate::graphql_gateway::TargetInputs {
4141 schema: Arc::new(schema),
4142 domain_context: Some("tenant_B".into()),
4143 target_handle: None,
4144 capability_anchor: None,
4145 }));
4146
4147 let out = crate::graphql_gateway::execute(&plan, &router, &serde_json::json!({})).await;
4148 let body = out.to_string();
4149 assert!(
4150 body.contains("\"base_pub\"") && body.contains("\"b_pub\""),
4151 "target_or_null returned B's public row AND the shared base floor: {body}"
4152 );
4153 assert!(
4154 !body.contains("\"base_priv\"")
4155 && !body.contains("\"b_priv\"")
4156 && !body.contains("\"a_pub\""),
4157 "the public subset confines BOTH B and base (no base_priv, no b_priv), and no other \
4158 tenant's rows (no a_pub): {body}"
4159 );
4160
4161 let _ = std::fs::remove_dir_all(&sql_dir);
4162 println!(
4163 "GATEWAY TARGET-OR-NULL OK: scope:target_or_null on a wasm subgraph read tenant B's \
4164 public row (b_pub) PLUS the shared NULL-tenant base floor (base_pub), each confined to \
4165 published=1 — never B's private row (b_priv), never a non-public base row (base_priv), \
4166 never another tenant's row (a_pub); the base-only funnel keeps its inheritance floor"
4167 );
4168 }
4169
4170 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4176 async fn scheduler_drains_a_non_default_projects_invocation_in_its_own_tenant() {
4177 use crate::scheduler::{run_scheduler_tick, CronNow};
4178 use boatramp_core::deploy::DeployStore;
4179 use boatramp_core::function::{
4180 Function, FunctionVersion, Invocation, InvocationStatus, InvokeMode, Lifecycle, Owner,
4181 };
4182 use boatramp_handlers::{HandlerEngine, Limits};
4183 use futures::StreamExt;
4184
4185 const HTTP_200: &[u8] =
4186 include_bytes!("../../boatramp-handlers/tests/fixtures/http-200.wasm");
4187
4188 let storage = Arc::new(MemStorage::default());
4189 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
4190 let deploy = DeployStore::new(storage.clone(), kv.clone());
4191
4192 let hash = boatramp_core::deploy::sha256_hex(HTTP_200);
4193 let stream: ByteStream =
4194 futures::stream::once(async move { Ok(bytes::Bytes::from_static(HTTP_200)) }).boxed();
4195 deploy.put_blob(&hash, stream).await.unwrap();
4196
4197 let acme = ProjectRef::new("acme");
4199 let function = Function {
4200 name: "worker".into(),
4201 owner: Owner::Project("acme".into()),
4202 versions: vec![FunctionVersion {
4203 id: "v1".into(),
4204 component: hash.clone(),
4205 created: 0,
4206 lifecycle: Lifecycle::Independent,
4207 }],
4208 active: "v1".into(),
4209 aliases: Default::default(),
4210 config: Default::default(),
4211 };
4212 deploy.put_function(acme, &function).await.unwrap();
4213 let inv = Invocation {
4214 id: "inv1".into(),
4215 function: "worker".into(),
4216 version: "v1".into(),
4217 mode: InvokeMode::Async,
4218 status: InvocationStatus::Queued,
4219 idempotency_key: None,
4220 attempts: 0,
4221 lease_expires: None,
4222 request_b64: None,
4223 request_content_type: None,
4224 result: None,
4225 created: 0,
4226 updated: 0,
4227 };
4228 deploy.put_invocation(acme, &inv).await.unwrap();
4229
4230 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4231 let rt = HandlerRuntime::new(engine, kv.clone(), storage.clone(), None, None);
4232 let inner = rt.inner.as_ref().unwrap();
4233
4234 let mut wasm_cache = std::collections::HashMap::new();
4237 let mut cron_state = std::collections::HashMap::new();
4238 let mut sweep = std::collections::HashMap::new();
4239 let now = CronNow {
4240 minute: 0,
4241 hour: 0,
4242 dom: 1,
4243 month: 1,
4244 dow: 0,
4245 minute_stamp: 0,
4246 };
4247 run_scheduler_tick(
4248 inner,
4249 &deploy,
4250 &mut wasm_cache,
4251 &mut cron_state,
4252 &mut sweep,
4253 now,
4254 )
4255 .await
4256 .unwrap();
4257
4258 let settled = poll_invocation_settled(&deploy, acme, "worker", "inv1").await;
4261 assert_eq!(settled.status, InvocationStatus::Succeeded);
4263 let metering = deploy.get_metering(acme, "worker").await.unwrap().unwrap();
4265 assert_eq!(metering.invocations, 1);
4266 assert!(deploy
4268 .get_invocation(ProjectRef::DEFAULT, "worker", "inv1")
4269 .await
4270 .unwrap()
4271 .is_none());
4272 assert!(deploy
4273 .get_metering(ProjectRef::DEFAULT, "worker")
4274 .await
4275 .unwrap()
4276 .is_none());
4277 }
4278
4279 #[cfg(feature = "handlers")]
4283 async fn poll_invocation_settled(
4284 deploy: &boatramp_core::deploy::DeployStore,
4285 project: ProjectRef<'_>,
4286 function: &str,
4287 id: &str,
4288 ) -> boatramp_core::function::Invocation {
4289 use boatramp_core::function::InvocationStatus;
4290 for _ in 0..200 {
4291 if let Some(inv) = deploy.get_invocation(project, function, id).await.unwrap() {
4292 if matches!(
4293 inv.status,
4294 InvocationStatus::Succeeded | InvocationStatus::Failed
4295 ) {
4296 return inv;
4297 }
4298 }
4299 tokio::time::sleep(std::time::Duration::from_millis(50)).await;
4300 }
4301 panic!("invocation {function}/{id} never settled");
4302 }
4303
4304 #[cfg(feature = "handlers")]
4309 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4310 async fn drain_reclaims_an_expired_lease_and_skips_a_live_one() {
4311 use crate::scheduler::{run_scheduler_tick, CronNow};
4312 use boatramp_core::deploy::DeployStore;
4313 use boatramp_core::function::{
4314 Function, FunctionVersion, Invocation, InvocationStatus, InvokeMode, Lifecycle, Owner,
4315 };
4316 use boatramp_handlers::{HandlerEngine, Limits};
4317 use futures::StreamExt;
4318
4319 const HTTP_200: &[u8] =
4320 include_bytes!("../../boatramp-handlers/tests/fixtures/http-200.wasm");
4321
4322 let storage = Arc::new(MemStorage::default());
4323 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
4324 let deploy = DeployStore::new(storage.clone(), kv.clone());
4325 let hash = boatramp_core::deploy::sha256_hex(HTTP_200);
4326 let stream: ByteStream =
4327 futures::stream::once(async move { Ok(bytes::Bytes::from_static(HTTP_200)) }).boxed();
4328 deploy.put_blob(&hash, stream).await.unwrap();
4329
4330 let function = Function {
4331 name: "worker".into(),
4332 owner: Owner::Project("default".into()),
4333 versions: vec![FunctionVersion {
4334 id: "v1".into(),
4335 component: hash.clone(),
4336 created: 0,
4337 lifecycle: Lifecycle::Independent,
4338 }],
4339 active: "v1".into(),
4340 aliases: Default::default(),
4341 config: Default::default(),
4342 };
4343 deploy
4344 .put_function(ProjectRef::DEFAULT, &function)
4345 .await
4346 .unwrap();
4347
4348 let base = Invocation {
4351 id: String::new(),
4352 function: "worker".into(),
4353 version: "v1".into(),
4354 mode: InvokeMode::Async,
4355 status: InvocationStatus::Running,
4356 idempotency_key: None,
4357 attempts: 1,
4358 lease_expires: None,
4359 request_b64: None,
4360 request_content_type: None,
4361 result: None,
4362 created: 0,
4363 updated: 0,
4364 };
4365 let orphan = Invocation {
4366 id: "orphan".into(),
4367 lease_expires: Some(1),
4368 ..base.clone()
4369 };
4370 deploy
4371 .put_invocation(ProjectRef::DEFAULT, &orphan)
4372 .await
4373 .unwrap();
4374 let live = Invocation {
4375 id: "live".into(),
4376 lease_expires: Some(u64::MAX),
4377 ..base.clone()
4378 };
4379 deploy
4380 .put_invocation(ProjectRef::DEFAULT, &live)
4381 .await
4382 .unwrap();
4383
4384 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4385 let rt = HandlerRuntime::new(engine, kv.clone(), storage.clone(), None, None);
4386 let inner = rt.inner.as_ref().unwrap();
4387
4388 let now = CronNow {
4389 minute: 0,
4390 hour: 0,
4391 dom: 1,
4392 month: 1,
4393 dow: 0,
4394 minute_stamp: 0,
4395 };
4396 let mut wasm_cache = std::collections::HashMap::new();
4397 let mut cron_state = std::collections::HashMap::new();
4398 let mut sweep = std::collections::HashMap::new();
4399 run_scheduler_tick(
4400 inner,
4401 &deploy,
4402 &mut wasm_cache,
4403 &mut cron_state,
4404 &mut sweep,
4405 now,
4406 )
4407 .await
4408 .unwrap();
4409
4410 let settled =
4412 poll_invocation_settled(&deploy, ProjectRef::DEFAULT, "worker", "orphan").await;
4413 assert_eq!(settled.status, InvocationStatus::Succeeded);
4414 assert_eq!(settled.attempts, 2, "a reclaim counts as another attempt");
4415 assert_eq!(
4416 settled.lease_expires, None,
4417 "a settled invocation drops its lease"
4418 );
4419 let live_after = deploy
4421 .get_invocation(ProjectRef::DEFAULT, "worker", "live")
4422 .await
4423 .unwrap()
4424 .unwrap();
4425 assert_eq!(live_after.status, InvocationStatus::Running);
4426 assert_eq!(live_after.attempts, 1, "a live lease is never reclaimed");
4427 assert_eq!(live_after.lease_expires, Some(u64::MAX));
4428 }
4429
4430 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4445 async fn guest_kv_is_isolated_between_same_named_functions_in_two_projects() {
4446 use boatramp_core::deploy::DeployStore;
4447 use boatramp_core::function::{
4448 Function, FunctionConfig, FunctionVersion, Lifecycle, Owner,
4449 };
4450 use boatramp_handlers::{HandlerEngine, Limits};
4451 use futures::StreamExt;
4452
4453 let storage = Arc::new(MemStorage::default());
4454 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
4455 let deploy = DeployStore::new(storage.clone(), kv.clone());
4456
4457 let hash = boatramp_core::deploy::sha256_hex(KV_COUNTER);
4460 let stream: ByteStream =
4461 futures::stream::once(async move { Ok(bytes::Bytes::from_static(KV_COUNTER)) }).boxed();
4462 deploy.put_blob(&hash, stream).await.unwrap();
4463
4464 let store = Function {
4469 name: "store".into(),
4470 owner: Owner::Project("default".into()),
4471 versions: vec![FunctionVersion {
4472 id: "v1".into(),
4473 component: hash.clone(),
4474 created: 0,
4475 lifecycle: Lifecycle::Independent,
4476 }],
4477 active: "v1".into(),
4478 aliases: Default::default(),
4479 config: FunctionConfig {
4480 imports: vec!["wasi:keyvalue".into()],
4481 ..Default::default()
4482 },
4483 };
4484 let acme = ProjectRef::new("acme");
4485 let globex = ProjectRef::new("globex");
4486
4487 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4488 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
4489 let inner = rt.inner.as_ref().unwrap();
4490
4491 let request = || {
4492 axum::http::Request::builder()
4493 .method("GET")
4494 .uri("/")
4495 .body(axum::body::Body::empty())
4496 .unwrap()
4497 };
4498
4499 let component = store.resolve(&store.active).unwrap().to_owned();
4502 for project in [acme, globex, ProjectRef::DEFAULT] {
4503 let (response, _) = execute_function(
4504 inner,
4505 &deploy,
4506 project,
4507 &store,
4508 &component,
4509 request(),
4510 0,
4511 boatramp_handlers::Lane::Sync,
4512 crate::function_runtime::FnTenant::Request,
4513 )
4514 .await;
4515 assert!(response.status().is_success(), "invocation should succeed");
4516 }
4517
4518 assert_eq!(
4521 kv.get("hkv/acme/fn/store/hits").await.unwrap(),
4522 Some(b"1".to_vec()),
4523 "acme's write must be tenant-qualified"
4524 );
4525 assert_eq!(
4526 kv.get("hkv/globex/fn/store/hits").await.unwrap(),
4527 Some(b"1".to_vec()),
4528 "globex's write must be tenant-qualified"
4529 );
4530 assert_eq!(
4531 kv.get("hkv/fn/store/hits").await.unwrap(),
4532 Some(b"1".to_vec()),
4533 "the default project must keep the byte-identical pre-project key"
4534 );
4535 }
4538
4539 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4543 async fn scheduler_fires_crons_with_dedup_and_overlap_skip() {
4544 use boatramp_core::config::{
4545 CronConfig, DeployConfig, HandlerConfig, HandlersSiteConfig, Overlap, SiteConfig,
4546 };
4547 use boatramp_core::deploy::{DeployStore, FileEntry, Manifest};
4548 use boatramp_handlers::{HandlerEngine, Limits};
4549 use futures::StreamExt;
4550 use std::sync::atomic::Ordering;
4551
4552 let storage = Arc::new(MemStorage::default());
4553 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
4554 let deploy = DeployStore::new(storage.clone(), kv.clone());
4555
4556 let hash = boatramp_core::deploy::sha256_hex(KV_COUNTER);
4557 let stream: ByteStream =
4558 futures::stream::once(async move { Ok(bytes::Bytes::from_static(KV_COUNTER)) }).boxed();
4559 deploy.put_blob(&hash, stream).await.unwrap();
4560 let mut files = std::collections::BTreeMap::new();
4561 files.insert(
4562 "counter.wasm".to_string(),
4563 FileEntry {
4564 hash: hash.clone(),
4565 size: KV_COUNTER.len() as u64,
4566 content_type: None,
4567 variants: std::collections::BTreeMap::new(),
4568 },
4569 );
4570 let manifest = Manifest {
4571 files,
4572 config: DeployConfig {
4573 handlers: vec![HandlerConfig {
4574 tenancy: None,
4575 token_claims: None,
4576 route: "/".into(),
4577 methods: Vec::new(),
4578 component: "counter.wasm".into(),
4579 imports: vec!["wasi:keyvalue".into()],
4580 streaming: false,
4581 limits: None,
4582 env: std::collections::BTreeMap::new(),
4583 invoke_targets: Vec::new(),
4584 }],
4585 crons: vec![CronConfig {
4586 schedule: "* * * * *".into(),
4587 route: "/".into(),
4588 overlap: Overlap::Skip,
4589 }],
4590 ..Default::default()
4591 },
4592 ..Default::default()
4593 };
4594 let id = deploy.put_manifest(&manifest).await.unwrap();
4595 deploy
4596 .activate(ProjectRef::DEFAULT, "blog", &id)
4597 .await
4598 .unwrap();
4599 deploy
4600 .set_site_config(
4601 ProjectRef::DEFAULT,
4602 "blog",
4603 &SiteConfig {
4604 handlers: Some(HandlersSiteConfig {
4605 enabled: true,
4606 allow_imports: vec!["wasi:keyvalue".into()],
4607 ..Default::default()
4608 }),
4609 ..Default::default()
4610 },
4611 )
4612 .await
4613 .unwrap();
4614
4615 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4616 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
4617 let inner = rt.inner.clone().unwrap();
4618 let mut wasm = std::collections::HashMap::new();
4619 let mut crons = std::collections::HashMap::new();
4620 let mut sweep = std::collections::HashMap::new();
4621 let at = |stamp| CronNow {
4622 minute: 0,
4623 hour: 0,
4624 dom: 1,
4625 month: 1,
4626 dow: 0,
4627 minute_stamp: stamp,
4628 };
4629
4630 let (_, handles) =
4632 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, at(100))
4633 .await
4634 .unwrap();
4635 for h in handles {
4636 h.await.unwrap();
4637 }
4638 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), Some(b"1".to_vec()));
4639
4640 let (_, handles) =
4642 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, at(100))
4643 .await
4644 .unwrap();
4645 assert!(handles.is_empty());
4646 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), Some(b"1".to_vec()));
4647
4648 let (_, handles) =
4650 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, at(101))
4651 .await
4652 .unwrap();
4653 for h in handles {
4654 h.await.unwrap();
4655 }
4656 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), Some(b"2".to_vec()));
4657
4658 crons
4662 .get("default|blog|cron|0")
4663 .unwrap()
4664 .running
4665 .store(true, Ordering::Release);
4666 let (_, handles) =
4667 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, at(102))
4668 .await
4669 .unwrap();
4670 assert!(handles.is_empty());
4671 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), Some(b"2".to_vec()));
4672 }
4673
4674 #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
4678 async fn cron_leader_gate_suppresses_crons_off_leader() {
4679 use boatramp_core::config::{
4680 CronConfig, DeployConfig, HandlerConfig, HandlersSiteConfig, Overlap, SiteConfig,
4681 };
4682 use boatramp_core::deploy::{DeployStore, FileEntry, Manifest};
4683 use boatramp_handlers::{HandlerEngine, Limits};
4684 use futures::StreamExt;
4685
4686 let storage = Arc::new(MemStorage::default());
4687 let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
4688 let deploy = DeployStore::new(storage.clone(), kv.clone());
4689
4690 let hash = boatramp_core::deploy::sha256_hex(KV_COUNTER);
4691 let stream: ByteStream =
4692 futures::stream::once(async move { Ok(bytes::Bytes::from_static(KV_COUNTER)) }).boxed();
4693 deploy.put_blob(&hash, stream).await.unwrap();
4694 let mut files = std::collections::BTreeMap::new();
4695 files.insert(
4696 "counter.wasm".to_string(),
4697 FileEntry {
4698 hash: hash.clone(),
4699 size: KV_COUNTER.len() as u64,
4700 content_type: None,
4701 variants: std::collections::BTreeMap::new(),
4702 },
4703 );
4704 let manifest = Manifest {
4705 files,
4706 config: DeployConfig {
4707 handlers: vec![HandlerConfig {
4708 tenancy: None,
4709 token_claims: None,
4710 route: "/".into(),
4711 methods: Vec::new(),
4712 component: "counter.wasm".into(),
4713 imports: vec!["wasi:keyvalue".into()],
4714 streaming: false,
4715 limits: None,
4716 env: std::collections::BTreeMap::new(),
4717 invoke_targets: Vec::new(),
4718 }],
4719 crons: vec![CronConfig {
4720 schedule: "* * * * *".into(),
4721 route: "/".into(),
4722 overlap: Overlap::Skip,
4723 }],
4724 ..Default::default()
4725 },
4726 ..Default::default()
4727 };
4728 let id = deploy.put_manifest(&manifest).await.unwrap();
4729 deploy
4730 .activate(ProjectRef::DEFAULT, "blog", &id)
4731 .await
4732 .unwrap();
4733 deploy
4734 .set_site_config(
4735 ProjectRef::DEFAULT,
4736 "blog",
4737 &SiteConfig {
4738 handlers: Some(HandlersSiteConfig {
4739 enabled: true,
4740 allow_imports: vec!["wasi:keyvalue".into()],
4741 ..Default::default()
4742 }),
4743 ..Default::default()
4744 },
4745 )
4746 .await
4747 .unwrap();
4748
4749 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4750 let rt = HandlerRuntime::new(engine, kv.clone(), storage, None, None);
4751 rt.set_cron_leader_gate(Arc::new(|| false));
4753 let inner = rt.inner.clone().unwrap();
4754 let mut wasm = std::collections::HashMap::new();
4755 let mut crons = std::collections::HashMap::new();
4756 let mut sweep = std::collections::HashMap::new();
4757 let now = CronNow {
4758 minute: 0,
4759 hour: 0,
4760 dom: 1,
4761 month: 1,
4762 dow: 0,
4763 minute_stamp: 100,
4764 };
4765
4766 let (_, handles) =
4767 run_scheduler_tick(&inner, &deploy, &mut wasm, &mut crons, &mut sweep, now)
4768 .await
4769 .unwrap();
4770 assert!(handles.is_empty(), "a non-leader must not fire crons");
4772 assert_eq!(kv.get("hkv/blog/hits").await.unwrap(), None);
4773 }
4774
4775 #[tokio::test]
4779 async fn project_tenancy_knobs_override_wins_else_node_base() {
4780 use boatramp_core::security::ResolvedProjectTenancy;
4781 use boatramp_handlers::{HandlerEngine, Limits};
4782
4783 let kv: Arc<dyn boatramp_core::kv::KvStore> = Arc::new(boatramp_core::kv::MemoryKv::new());
4784 let storage: Arc<dyn boatramp_core::Storage> = Arc::new(MemStorage::default());
4785 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4786 let rt = HandlerRuntime::new(engine, kv, storage, None, None);
4787 rt.set_tenancy_posture(true, false);
4789 let mut overrides = std::collections::BTreeMap::new();
4791 overrides.insert(
4792 "preview".to_string(),
4793 ResolvedProjectTenancy {
4794 require_tenancy_declaration: true,
4795 allow_cross_tenant_db: true,
4796 capability_max_ttl_secs: Some(1800),
4797 },
4798 );
4799 rt.set_project_tenancy_overrides(overrides);
4800 let inner = rt.inner.as_ref().unwrap();
4801
4802 let p = inner.project_tenancy_knobs("preview");
4804 assert!(p.require_tenancy_declaration);
4805 assert!(p.allow_cross_tenant_db);
4806 assert_eq!(p.capability_max_ttl_secs, Some(1800));
4807 let b = inner.project_tenancy_knobs("prod");
4809 assert!(b.require_tenancy_declaration);
4810 assert!(!b.allow_cross_tenant_db);
4811 assert_eq!(b.capability_max_ttl_secs, None);
4812 }
4813
4814 #[tokio::test]
4821 async fn build_bindings_dispatches_named_sql_databases_with_least_privilege() {
4822 use boatramp_core::config::HandlersSiteConfig;
4823 use boatramp_core::project::ProjectRef;
4824 use boatramp_handlers::{HandlerEngine, Limits};
4825
4826 let kv: Arc<dyn boatramp_core::kv::KvStore> = Arc::new(boatramp_core::kv::MemoryKv::new());
4827 let storage: Arc<dyn boatramp_core::Storage> = Arc::new(MemStorage::default());
4828 let sql_dir =
4830 std::env::temp_dir().join(format!("boatramp-named-sql-{}", std::process::id()));
4831 let _ = std::fs::remove_dir_all(&sql_dir);
4832 let sql: Arc<dyn boatramp_core::sql::SqlBackends> =
4833 Arc::new(boatramp_storage::LibsqlSqlBackends::local(&sql_dir));
4834
4835 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4836 let rt = HandlerRuntime::new(engine, kv, storage, Some(sql), None);
4837 let inner = rt.inner.as_ref().unwrap();
4838
4839 let site = HandlersSiteConfig {
4843 enabled: true,
4844 allow_imports: vec!["sql".into(), "sql:product".into(), "sql:privileged".into()],
4845 tenancy: Some(boatramp_core::tenancy::Tenancy::Disabled),
4846 ..Default::default()
4847 };
4848 let env = std::collections::BTreeMap::new();
4849 let build = |imports: &[&str]| {
4850 let imports: Vec<String> = imports.iter().copied().map(String::from).collect();
4851 let site = &site;
4852 let env = &env;
4853 async move {
4854 crate::handler_dispatch::build_bindings(
4855 inner,
4856 ProjectRef::new("default"),
4857 "shop",
4858 "shop",
4859 None,
4860 &imports,
4861 site,
4862 env,
4863 &[],
4864 0,
4865 None,
4866 None,
4867 None,
4868 None,
4869 None,
4870 None,
4871 None,
4872 None,
4873 )
4874 .await
4875 .expect("no secrets → resolves")
4876 .sql_database_names()
4877 }
4878 };
4879
4880 assert_eq!(build(&["sql", "sql:product"]).await, vec!["", "product"]);
4883 assert_eq!(
4885 build(&["sql", "sql:*"]).await,
4886 vec!["", "privileged", "product"]
4887 );
4888 assert!(build(&["sql:secret"]).await.is_empty());
4890 assert_eq!(build(&["sql:product"]).await, vec!["product"]);
4892
4893 let _ = std::fs::remove_dir_all(&sql_dir);
4894 }
4895
4896 #[tokio::test]
4911 async fn consumer_signed_context_dispatch_resolves_per_message_tenant_on_a_real_engine() {
4912 use crate::handler_dispatch::ConsumerRebuild;
4913 use boatramp_core::config::HandlersSiteConfig;
4914 use boatramp_core::cose::{mint_context, Signer};
4915 use boatramp_core::orm::{Expr, Select, SelectItem};
4916 use boatramp_core::sql::{Dialect, SqlBackends, SqlValue};
4917 use boatramp_core::tenancy::{AccessMode, Tenancy, TenantSource};
4918 use boatramp_handlers::{HandlerEngine, Limits, TenantAxis, TenantDenied};
4919
4920 let sql_dir =
4922 std::env::temp_dir().join(format!("boatramp-consumer-sctx-{}", std::process::id()));
4923 let _ = std::fs::remove_dir_all(&sql_dir);
4924 let backends = boatramp_storage::LibsqlSqlBackends::local(&sql_dir);
4925 let db = backends.database("default", "worker", "").await.unwrap();
4926 {
4927 let mut tx = db.begin().await.unwrap();
4928 tx.execute(
4929 "CREATE TABLE notes (id TEXT PRIMARY KEY, tenant_id TEXT, body TEXT)",
4930 &[],
4931 )
4932 .await
4933 .unwrap();
4934 for (id, tenant, body) in [("n1", "acme", "acme-note"), ("n2", "globex", "globex-note")]
4935 {
4936 tx.execute(
4937 "INSERT INTO notes (id, tenant_id, body) VALUES (?1, ?2, ?3)",
4938 &[
4939 SqlValue::Text(id.into()),
4940 SqlValue::Text(tenant.into()),
4941 SqlValue::Text(body.into()),
4942 ],
4943 )
4944 .await
4945 .unwrap();
4946 }
4947 tx.commit().await.unwrap();
4948 }
4949
4950 let kv: Arc<dyn boatramp_core::kv::KvStore> = Arc::new(boatramp_core::kv::MemoryKv::new());
4951 let storage: Arc<dyn boatramp_core::Storage> = Arc::new(MemStorage::default());
4952 let sql: Arc<dyn boatramp_core::sql::SqlBackends> =
4953 Arc::new(boatramp_storage::LibsqlSqlBackends::local(&sql_dir));
4954 let engine = HandlerEngine::new(Limits::default(), 16).unwrap();
4955 let rt = HandlerRuntime::new(engine, kv, storage, Some(sql), None);
4956 rt.set_tenancy_posture(true, false);
4959 let signer: Arc<dyn Signer> = Arc::new(LocalSigner::generate(TokenAlg::Es256));
4961 rt.set_session_signer(signer.clone());
4962 let stranger: Arc<dyn Signer> = Arc::new(LocalSigner::generate(TokenAlg::Es256));
4963 let inner = rt.inner.as_ref().unwrap();
4964
4965 let tenancy = Tenancy::Scoped {
4967 column: "tenant_id".into(),
4968 sources: vec![TenantSource::SignedContext],
4969 read: AccessMode::Own,
4970 write: AccessMode::Own,
4971 };
4972 let site = HandlersSiteConfig {
4973 enabled: true,
4974 allow_imports: vec!["sql".into()],
4975 ..Default::default()
4977 };
4978 let imports = vec!["sql".to_string()];
4979 let rebuild = ConsumerRebuild {
4980 inner,
4981 project: ProjectRef::new("default"),
4982 site: "worker",
4983 scope: "worker",
4984 imports: &imports,
4985 site_handlers: &site,
4986 tenancy: Some(&tenancy),
4987 token_claims: None,
4988 };
4989
4990 async fn read_bodies(
4992 db: &dyn boatramp_core::sql::SqlBackend,
4993 scope: &boatramp_core::orm::Scope,
4994 ) -> Vec<String> {
4995 let mut q = Select {
4996 columns: vec![SelectItem {
4997 expr: Expr::col("body"),
4998 alias: None,
4999 }],
5000 ..Select::from("notes")
5001 };
5002 q.force_scope(scope).unwrap();
5003 let (sql, params) = q.compile(Dialect::Sqlite).unwrap();
5004 let mut tx = db.begin().await.unwrap();
5005 let rows = tx.query(&sql, ¶ms).await.expect("query runs");
5006 tx.commit().await.unwrap();
5007 let mut out: Vec<String> = rows
5008 .rows
5009 .into_iter()
5010 .flatten()
5011 .filter_map(|v| match v {
5012 SqlValue::Text(s) => Some(s),
5013 _ => None,
5014 })
5015 .collect();
5016 out.sort();
5017 out
5018 }
5019
5020 let now = boatramp_core::time::now_unix();
5021
5022 let env = mint_context("acme", 3600, now, signer.as_ref())
5024 .await
5025 .unwrap();
5026 let sealed = rebuild.bindings_for(Some(&env)).await.unwrap();
5027 let ht = sealed
5028 .resolved_tenancy()
5029 .expect("a signed-context consumer resolves a tenancy");
5030 assert_eq!(
5031 ht.value(),
5032 Some(&SqlValue::Text("acme".into())),
5033 "the sealed envelope's tenant is the consumer's own principal"
5034 );
5035 assert!(
5037 ht.facts()
5038 .iter()
5039 .any(|f| f.value == SqlValue::Text("acme".into())),
5040 "the resolved fact carries acme (this is the graphql::run caller principal too)"
5041 );
5042 let read_scope = ht.orm_scope(TenantAxis::Read).unwrap().unwrap();
5043 assert_eq!(
5044 read_bodies(db.as_ref(), &read_scope).await,
5045 vec!["acme-note".to_string()],
5046 "the resolved tenant scopes a real engine to acme's row ONLY (never globex)"
5047 );
5048 assert!(
5050 ht.orm_scope(TenantAxis::Write).unwrap().is_some(),
5051 "the write axis resolves the sealed tenant"
5052 );
5053
5054 let plain = rebuild.bindings_for(None).await.unwrap();
5056 let ht = plain
5057 .resolved_tenancy()
5058 .expect("the tenancy decision is present (but factless)");
5059 assert!(
5060 ht.value().is_none(),
5061 "an unsealed message resolves no own tenant"
5062 );
5063 assert!(
5064 matches!(ht.orm_scope(TenantAxis::Read), Err(TenantDenied::NoSource)),
5065 "a scoped op with no resolved tenant fails closed (never runs unscoped)"
5066 );
5067
5068 let forged = mint_context("globex", 3600, now, stranger.as_ref())
5070 .await
5071 .unwrap();
5072 let forged_b = rebuild.bindings_for(Some(&forged)).await.unwrap();
5073 let ht = forged_b
5074 .resolved_tenancy()
5075 .expect("the tenancy decision is present (but factless)");
5076 assert!(
5077 ht.value().is_none(),
5078 "a stranger-signed envelope resolves no own tenant (never masquerades as globex)"
5079 );
5080 assert!(
5081 matches!(ht.orm_scope(TenantAxis::Read), Err(TenantDenied::NoSource)),
5082 "a forged envelope fails the scoped op closed"
5083 );
5084
5085 let expired = mint_context("acme", 3600, now.saturating_sub(7200), signer.as_ref())
5088 .await
5089 .unwrap();
5090 let expired_b = rebuild.bindings_for(Some(&expired)).await.unwrap();
5091 let ht = expired_b
5092 .resolved_tenancy()
5093 .expect("the tenancy decision is present (but factless)");
5094 assert!(
5095 ht.value().is_none()
5096 && matches!(ht.orm_scope(TenantAxis::Read), Err(TenantDenied::NoSource)),
5097 "an expired envelope resolves no own tenant and fails the scoped op closed"
5098 );
5099
5100 let env_b = mint_context("globex", 3600, now, signer.as_ref())
5105 .await
5106 .unwrap();
5107 let a = rebuild.bindings_for(Some(&env)).await.unwrap();
5108 let b = rebuild.bindings_for(Some(&env_b)).await.unwrap();
5109 let ht_a = a.resolved_tenancy().unwrap();
5110 let ht_b = b.resolved_tenancy().unwrap();
5111 assert_eq!(ht_a.value(), Some(&SqlValue::Text("acme".into())));
5112 assert_eq!(ht_b.value(), Some(&SqlValue::Text("globex".into())));
5113 assert_eq!(
5114 read_bodies(
5115 db.as_ref(),
5116 &ht_a.orm_scope(TenantAxis::Read).unwrap().unwrap()
5117 )
5118 .await,
5119 vec!["acme-note".to_string()],
5120 "message A stays scoped to acme"
5121 );
5122 assert_eq!(
5123 read_bodies(
5124 db.as_ref(),
5125 &ht_b.orm_scope(TenantAxis::Read).unwrap().unwrap()
5126 )
5127 .await,
5128 vec!["globex-note".to_string()],
5129 "message B (interleaved) scopes to globex ONLY — no bleed from A's binding"
5130 );
5131
5132 println!(
5133 "CONSUMER SIGNED-CONTEXT DISPATCH OK: the site-consumers async lane resolves each \
5134 claimed message's host-sealed signed_context per message; a valid fleet-signed \
5135 envelope scopes a real libsql engine to the originator's tenant ONLY, while an \
5136 unsealed or forged message resolves no principal and fails an own op closed (never \
5137 cross-tenant, never unscoped). The resolved principal is the same value that \
5138 propagates onto a graphql::run sub-fetch."
5139 );
5140
5141 let _ = std::fs::remove_dir_all(&sql_dir);
5142 }
5143}