1use crate::time::now_unix;
19use std::collections::{BTreeMap, BTreeSet};
20use std::sync::{Arc, Mutex};
21
22use futures::StreamExt;
23use sha2::{Digest, Sha256};
24
25use crate::config::SiteConfig;
26use crate::domain_verify::{DomainVerification, VerificationMethod};
27use crate::error::DeployError;
28use crate::kv::{KvStore, WriteOp};
29use crate::project::{DomainOwner, ProjectRef};
30use crate::site::SiteName;
31use crate::{ByteStream, GetObject, PutMeta, Storage, StorageError};
32
33pub use boatramp_types::file::{FileEntry, Variant};
40pub use boatramp_types::manifest::{sha256_hex, Manifest};
41
42pub use boatramp_types::deploy::{
48 BlobMismatch, BlobReadError, DeployMeta, DeployMetaInput, DeploymentList, GcReport,
49 HistoryEntry, ScrubReport,
50};
51
52#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
56pub struct GcOptions {
57 pub grace_secs: u64,
61 pub keep_last: Option<usize>,
65 pub keep_age_secs: Option<u64>,
68}
69
70const MAX_HISTORY: usize = 100;
72
73fn canon_host(host: &str) -> String {
80 crate::host::Host::new(host).routing_key()
81}
82
83fn backfill_replica_project(state: &mut crate::compute::ObservedInstance, project: &str) {
89 if state.handle.project.is_empty() {
90 state.handle.project = project.to_string();
91 }
92 if let Some(snap) = state.snapshot.as_mut() {
93 if snap.project.is_empty() {
94 snap.project = project.to_string();
95 }
96 }
97}
98
99fn is_blob_key(key: &str) -> bool {
102 match key.split_once('/') {
103 Some((shard, hash)) => {
104 shard.len() == 2
105 && hash.len() == 64
106 && shard.bytes().all(|b| b.is_ascii_hexdigit())
107 && hash.bytes().all(|b| b.is_ascii_hexdigit())
108 }
109 None => false,
110 }
111}
112
113#[derive(Clone)]
122pub struct DeployStore(Arc<DeployStoreInner>);
123
124impl std::ops::Deref for DeployStore {
125 type Target = DeployStoreInner;
126 fn deref(&self) -> &DeployStoreInner {
127 &self.0
128 }
129}
130
131pub struct DeployStoreInner {
133 storage: Arc<dyn Storage>,
134 kv: Arc<dyn KvStore>,
135 domain_claim_lock: Arc<futures::lock::Mutex<()>>,
142 site_config_cache: Arc<std::sync::RwLock<std::collections::HashMap<String, Arc<SiteConfig>>>>,
150 blob_body_cache: Arc<std::sync::RwLock<BlobBodyCache>>,
157 domain_cache: Arc<std::sync::RwLock<DomainResolveCache>>,
165 domain_epoch: Arc<std::sync::atomic::AtomicU64>,
171}
172
173#[derive(Default)]
176struct DomainResolveCache {
177 epoch: u64,
178 map: std::collections::HashMap<String, Option<DomainOwner>>,
179}
180
181pub enum BlobBody {
184 Cached(bytes::Bytes),
186 Mapped(bytes::Bytes),
190 Stream(GetObject),
192}
193
194#[derive(Default)]
199struct BlobBodyCache {
200 map: std::collections::HashMap<String, bytes::Bytes>,
201 bytes: usize,
202}
203
204const SMALL_BLOB_CACHE_MAX: u64 = 256 * 1024;
207const BLOB_BODY_CACHE_MAX_BYTES: usize = 64 * 1024 * 1024;
209
210pub(crate) mod keys {
211 use crate::project::ProjectRef;
228
229 pub fn manifest(id: &str) -> String {
231 format!("manifests/{id}")
232 }
233
234 pub fn meta(id: &str) -> String {
236 format!("meta/{id}")
237 }
238
239 pub fn current(project: ProjectRef<'_>, site: &str) -> String {
241 format!("project/{project}/current/{site}")
242 }
243
244 pub fn current_prefix(project: ProjectRef<'_>) -> String {
246 format!("project/{project}/current/")
247 }
248
249 pub fn alias(project: ProjectRef<'_>, site: &str, name: &str) -> String {
251 format!("project/{project}/alias/{site}/{name}")
252 }
253
254 pub fn alias_prefix(project: ProjectRef<'_>, site: &str) -> String {
256 format!("project/{project}/alias/{site}/")
257 }
258
259 pub fn alias_project_prefix(project: ProjectRef<'_>) -> String {
261 format!("project/{project}/alias/")
262 }
263
264 pub fn secret(project: ProjectRef<'_>, name: &str) -> String {
270 format!("project/{project}/secret/{name}")
271 }
272
273 pub fn secret_prefix(project: ProjectRef<'_>) -> String {
275 format!("project/{project}/secret/")
276 }
277
278 pub fn email_profile(project: ProjectRef<'_>, name: &str) -> String {
285 format!("project/{project}/email/{name}")
286 }
287
288 pub fn email_profile_prefix(project: ProjectRef<'_>) -> String {
290 format!("project/{project}/email/")
291 }
292
293 pub fn project_config(project: ProjectRef<'_>, key: &str) -> String {
297 format!("project/{project}/config/{key}")
298 }
299
300 pub fn graphql_safelist_prefix(project: ProjectRef<'_>) -> String {
306 format!("hapq/{project}/")
307 }
308
309 pub fn graphql_registry_prefix(project: ProjectRef<'_>) -> String {
314 format!("graphql/{project}/")
315 }
316
317 pub fn graphql_subgraph_prefix(project: ProjectRef<'_>) -> String {
322 format!("graphql/{project}/subgraph/")
323 }
324
325 pub fn blob(hash: &str) -> String {
327 if hash.len() >= 2 {
328 format!("{}/{}", &hash[..2], hash)
329 } else {
330 hash.to_string()
331 }
332 }
333
334 pub fn site_pointer(project: ProjectRef<'_>, site: &str) -> String {
338 format!("project/{project}/site/{site}")
339 }
340
341 pub fn site_prefix(project: ProjectRef<'_>) -> String {
343 format!("project/{project}/site/")
344 }
345
346 pub fn site_config_blob(hash: &str) -> String {
351 format!("siteconfig/{hash}")
352 }
353
354 pub fn domain(host: &str) -> String {
358 format!("domain/{}", super::canon_host(host))
359 }
360
361 pub fn wildcard(suffix: &str) -> String {
364 format!("wildcard/{}", super::canon_host(suffix))
365 }
366
367 pub fn domain_verification(project: ProjectRef<'_>, site: &str, host: &str) -> String {
370 format!(
371 "project/{project}/domainverify/{site}/{}",
372 crate::domain_verify::normalize_host(host)
373 )
374 }
375
376 pub fn domain_verification_prefix(project: ProjectRef<'_>, site: &str) -> String {
378 format!("project/{project}/domainverify/{site}/")
379 }
380
381 pub fn domain_verification_project_prefix(project: ProjectRef<'_>) -> String {
383 format!("project/{project}/domainverify/")
384 }
385
386 pub fn http_challenge_index(host: &str, token: &str) -> String {
392 format!(
393 "httpchallenge/{}/{token}",
394 crate::domain_verify::normalize_host(host)
395 )
396 }
397
398 pub fn daemon_config_blob(hash: &str) -> String {
401 format!("daemonconfig/{hash}")
402 }
403
404 pub fn history(project: ProjectRef<'_>, site: &str) -> String {
406 format!("project/{project}/history/{site}")
407 }
408
409 pub fn history_prefix(project: ProjectRef<'_>) -> String {
411 format!("project/{project}/history/")
412 }
413
414 pub const PROJECT_ROOT: &str = "project/";
418
419 pub fn project_of_key(key: &str) -> Option<&str> {
422 key.strip_prefix(PROJECT_ROOT)
423 .and_then(|rest| rest.split('/').next())
424 .filter(|p| !p.is_empty())
425 }
426}
427
428#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
434pub struct ProjectTeardownPlan {
435 pub project: String,
437 pub sites: Vec<String>,
440 pub functions: Vec<String>,
442 pub compute: Vec<ComputeTeardown>,
445 pub secrets: Vec<String>,
447 pub safelist: usize,
449 pub subgraphs: Vec<String>,
451 pub other_families: std::collections::BTreeMap<String, usize>,
455}
456
457impl ProjectTeardownPlan {
458 pub fn is_empty(&self) -> bool {
461 self.sites.is_empty()
462 && self.functions.is_empty()
463 && self.compute.is_empty()
464 && self.secrets.is_empty()
465 && self.safelist == 0
466 && self.subgraphs.is_empty()
467 && self.other_families.is_empty()
468 }
469
470 pub fn all_volumes(&self) -> Vec<String> {
474 let mut vols: std::collections::BTreeSet<String> = Default::default();
475 for c in &self.compute {
476 for v in &c.volumes {
477 vols.insert(v.clone());
478 }
479 }
480 vols.into_iter().collect()
481 }
482}
483
484#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
487pub struct ComputeTeardown {
488 pub name: String,
490 pub volumes: Vec<String>,
492}
493
494pub async fn load_project_tenancy(
502 kv: &dyn KvStore,
503 project: ProjectRef<'_>,
504) -> Result<Option<crate::tenancy::TenancySchema>, DeployError> {
505 match kv.get(&keys::project_config(project, "tenancy")).await? {
506 Some(bytes) => Ok(Some(
507 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
508 )),
509 None => Ok(None),
510 }
511}
512
513impl DeployStore {
514 pub fn new(storage: Arc<dyn Storage>, kv: Arc<dyn KvStore>) -> Self {
516 Self(Arc::new(DeployStoreInner {
517 storage,
518 kv,
519 domain_claim_lock: Arc::new(futures::lock::Mutex::new(())),
520 site_config_cache: Arc::new(std::sync::RwLock::new(std::collections::HashMap::new())),
521 blob_body_cache: Arc::new(std::sync::RwLock::new(BlobBodyCache::default())),
522 domain_cache: Arc::new(std::sync::RwLock::new(DomainResolveCache::default())),
523 domain_epoch: Arc::new(std::sync::atomic::AtomicU64::new(0)),
524 }))
525 }
526
527 pub fn kv(&self) -> &Arc<dyn KvStore> {
530 &self.kv
531 }
532
533 pub async fn ready(&self) -> Result<(), DeployError> {
537 self.kv.get("__readyz_probe__").await?;
538 Ok(())
539 }
540
541 pub async fn put_manifest(&self, manifest: &Manifest) -> Result<String, DeployError> {
543 self.put_manifest_with(manifest, DeployMetaInput::default())
544 .await
545 }
546
547 pub async fn put_manifest_with(
555 &self,
556 manifest: &Manifest,
557 input: DeployMetaInput,
558 ) -> Result<String, DeployError> {
559 let id = manifest.id()?;
560
561 let existing = self.get_meta(&id).await?;
565 let created_at = existing
566 .as_ref()
567 .map(|m| m.created_at)
568 .unwrap_or_else(now_unix);
569 let meta = DeployMeta {
570 version: crate::SCHEMA_VERSION,
571 created_at,
572 file_count: manifest.files.len() as u64,
573 total_size: manifest.files.values().map(|entry| entry.size).sum(),
574 source: input
575 .source
576 .or_else(|| existing.as_ref().and_then(|m| m.source.clone())),
577 branch: input
578 .branch
579 .or_else(|| existing.as_ref().and_then(|m| m.branch.clone())),
580 author: input
581 .author
582 .or_else(|| existing.as_ref().and_then(|m| m.author.clone())),
583 message: input
584 .message
585 .or_else(|| existing.as_ref().and_then(|m| m.message.clone())),
586 tag: input
587 .tag
588 .or_else(|| existing.as_ref().and_then(|m| m.tag.clone())),
589 tags: if input.tags.is_empty() {
592 existing.map(|m| m.tags).unwrap_or_default()
593 } else {
594 input.tags
595 },
596 };
597 self.kv
598 .write_batch(vec![
599 WriteOp::Put(keys::manifest(&id), manifest.to_bytes()?),
600 WriteOp::Put(keys::meta(&id), serde_json::to_vec(&meta)?),
601 ])
602 .await?;
603 Ok(id)
604 }
605
606 pub async fn get_meta(&self, id: &str) -> Result<Option<DeployMeta>, DeployError> {
608 match self.kv.get(&keys::meta(id)).await? {
609 Some(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
610 None => Ok(None),
611 }
612 }
613
614 pub async fn get_manifest(&self, id: &str) -> Result<Option<Manifest>, DeployError> {
616 match self.kv.get(&keys::manifest(id)).await? {
617 Some(bytes) => Ok(Some(Manifest::from_bytes(&bytes)?)),
618 None => Ok(None),
619 }
620 }
621
622 pub async fn resolve_manifest_id(&self, prefix: &str) -> Result<Option<String>, DeployError> {
629 if self.kv.get(&keys::manifest(prefix)).await?.is_some() {
631 return Ok(Some(prefix.to_string()));
632 }
633 let key_prefix = keys::manifest(prefix);
634 let strip = "manifests/".len();
635 let keys = self.kv.list_prefix(&key_prefix).await?;
636 let mut ids = keys.iter().map(|key| &key[strip..]);
637 match (ids.next(), ids.next()) {
638 (Some(only), None) => Ok(Some(only.to_string())),
639 (Some(_), Some(_)) => Err(DeployError::Ambiguous(prefix.to_string())),
640 _ => Ok(None),
641 }
642 }
643
644 pub async fn has_blob(&self, hash: &str) -> Result<bool, DeployError> {
646 match self.storage.head(&keys::blob(hash)).await {
647 Ok(_) => Ok(true),
648 Err(StorageError::NotFound(_)) => Ok(false),
649 Err(err) => Err(err.into()),
650 }
651 }
652
653 pub async fn missing_blobs(&self, manifest: &Manifest) -> Result<Vec<String>, DeployError> {
655 let mut missing = Vec::new();
656 for hash in manifest.blob_hashes() {
657 if !self.has_blob(&hash).await? {
658 missing.push(hash);
659 }
660 }
661 Ok(missing)
662 }
663
664 pub async fn put_blob(&self, hash: &str, body: ByteStream) -> Result<(), DeployError> {
669 let hasher = Arc::new(Mutex::new(Sha256::new()));
670 let tap = hasher.clone();
671 let verified: ByteStream = body
672 .map(move |chunk| {
673 if let Ok(bytes) = &chunk {
674 tap.lock().unwrap().update(bytes);
675 }
676 chunk
677 })
678 .boxed();
679
680 let key = keys::blob(hash);
681 self.storage.put(&key, verified, PutMeta::default()).await?;
682
683 let actual = hex::encode(hasher.lock().unwrap().clone().finalize());
684 if actual != hash {
685 let _ = self.storage.delete(&key).await;
686 return Err(DeployError::HashMismatch {
687 expected: hash.to_string(),
688 actual,
689 });
690 }
691 Ok(())
692 }
693
694 pub async fn open_blob(&self, hash: &str) -> Result<GetObject, DeployError> {
696 Ok(self.storage.get(&keys::blob(hash)).await?)
697 }
698
699 pub fn blob_file(&self, hash: &str) -> Option<std::fs::File> {
704 self.storage.local_file(&keys::blob(hash))
705 }
706
707 pub async fn open_blob_cached(&self, hash: &str, size: u64) -> Result<BlobBody, DeployError> {
714 use futures::TryStreamExt;
715 if size > SMALL_BLOB_CACHE_MAX {
716 if let Some(bytes) = self.storage.mapped(&keys::blob(hash)) {
721 return Ok(BlobBody::Mapped(bytes));
722 }
723 return Ok(BlobBody::Stream(self.open_blob(hash).await?));
724 }
725 if let Some(bytes) = self.blob_body_cache.read().unwrap().map.get(hash).cloned() {
726 return Ok(BlobBody::Cached(bytes));
727 }
728 let object = self.storage.get(&keys::blob(hash)).await?;
731 let mut buf = bytes::BytesMut::with_capacity(size as usize);
732 let mut body = object.body;
733 while let Some(chunk) = body.try_next().await? {
734 buf.extend_from_slice(&chunk);
735 }
736 let bytes = buf.freeze();
737 {
738 let mut cache = self.blob_body_cache.write().unwrap();
739 if cache.bytes.saturating_add(bytes.len()) > BLOB_BODY_CACHE_MAX_BYTES {
741 cache.map.clear();
742 cache.bytes = 0;
743 }
744 if cache.map.insert(hash.to_string(), bytes.clone()).is_none() {
745 cache.bytes = cache.bytes.saturating_add(bytes.len());
746 }
747 }
748 Ok(BlobBody::Cached(bytes))
749 }
750
751 pub async fn open_blob_range(
753 &self,
754 hash: &str,
755 offset: u64,
756 len: Option<u64>,
757 ) -> Result<GetObject, DeployError> {
758 Ok(self
759 .storage
760 .get_range(&keys::blob(hash), offset, len)
761 .await?)
762 }
763
764 pub async fn get_site_config(
767 &self,
768 project: ProjectRef<'_>,
769 site: &str,
770 ) -> Result<Option<SiteConfig>, DeployError> {
771 let Some(hash) = self.kv.get(&keys::site_pointer(project, site)).await? else {
772 return Ok(None);
773 };
774 let hash = String::from_utf8_lossy(&hash).into_owned();
775 match self.kv.get(&keys::site_config_blob(&hash)).await? {
776 Some(bytes) => Ok(Some(SiteConfig::from_json(&bytes)?)),
777 None => Ok(None),
779 }
780 }
781
782 pub async fn get_project_tenancy(
788 &self,
789 project: ProjectRef<'_>,
790 ) -> Result<Option<crate::tenancy::TenancySchema>, DeployError> {
791 match self
792 .kv
793 .get(&keys::project_config(project, "tenancy"))
794 .await?
795 {
796 Some(bytes) => Ok(Some(
797 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
798 )),
799 None => Ok(None),
800 }
801 }
802
803 pub async fn set_project_tenancy(
809 &self,
810 project: ProjectRef<'_>,
811 schema: &crate::tenancy::TenancySchema,
812 ) -> Result<(), DeployError> {
813 schema.validate().map_err(DeployError::Invalid)?;
814 let bytes = serde_json::to_vec(schema).map_err(|e| DeployError::Serde(e.to_string()))?;
815 self.kv
816 .put(&keys::project_config(project, "tenancy"), bytes)
817 .await?;
818 Ok(())
819 }
820
821 pub async fn clear_project_tenancy(&self, project: ProjectRef<'_>) -> Result<(), DeployError> {
825 self.kv
826 .delete(&keys::project_config(project, "tenancy"))
827 .await?;
828 Ok(())
829 }
830
831 pub async fn get_site_config_cached(
838 &self,
839 project: ProjectRef<'_>,
840 site: &str,
841 ) -> Result<Option<Arc<SiteConfig>>, DeployError> {
842 let Some(hash) = self.kv.get(&keys::site_pointer(project, site)).await? else {
843 return Ok(None);
844 };
845 let hash = String::from_utf8_lossy(&hash).into_owned();
846 if let Some(cfg) = self.site_config_cache.read().unwrap().get(&hash).cloned() {
847 return Ok(Some(cfg));
848 }
849 let Some(bytes) = self.kv.get(&keys::site_config_blob(&hash)).await? else {
850 return Ok(None); };
852 let cfg = Arc::new(SiteConfig::from_json(&bytes)?);
853 {
854 let mut cache = self.site_config_cache.write().unwrap();
855 if cache.len() >= 512 {
859 cache.clear();
860 }
861 cache.insert(hash, Arc::clone(&cfg));
862 }
863 Ok(Some(cfg))
864 }
865
866 pub async fn set_site_config(
886 &self,
887 project: ProjectRef<'_>,
888 site: &str,
889 config: &SiteConfig,
890 ) -> Result<(), DeployError> {
891 let _claim = self.domain_claim_lock.lock().await;
892 self.set_site_config_locked(project, site, config).await
893 }
894
895 async fn set_site_config_locked(
900 &self,
901 project: ProjectRef<'_>,
902 site: &str,
903 config: &SiteConfig,
904 ) -> Result<(), DeployError> {
905 let owner = DomainOwner::new(project.as_str(), site);
906 for host in config.domains.exact_hosts() {
910 self.ensure_host_claimable(&keys::domain(host), host, &owner)
911 .await?;
912 }
913 for wildcard in &config.domains.wildcards {
914 if let Some(suffix) = wildcard.strip_prefix("*.") {
915 self.ensure_host_claimable(&keys::wildcard(suffix), wildcard, &owner)
916 .await?;
917 }
918 }
919
920 let body = config.to_json()?;
921 let hash = sha256_hex(&body);
922
923 let mut ops = Vec::new();
924 if let Some(old) = self.get_site_config(project, site).await? {
925 for host in old.domains.exact_hosts() {
926 ops.push(WriteOp::Delete(keys::domain(host)));
927 }
928 for wildcard in &old.domains.wildcards {
929 if let Some(suffix) = wildcard.strip_prefix("*.") {
930 ops.push(WriteOp::Delete(keys::wildcard(suffix)));
931 }
932 }
933 }
934
935 ops.push(WriteOp::Put(keys::site_config_blob(&hash), body));
937 ops.push(WriteOp::Put(
938 keys::site_pointer(project, site),
939 hash.into_bytes(),
940 ));
941
942 let contexts = &config.domains.contexts;
948 let primary_ctx = config
949 .domains
950 .primary
951 .as_deref()
952 .and_then(|p| contexts.get(p));
953 let owner_for = |key: &str| -> Vec<u8> {
954 let ctx = contexts.get(key).or(primary_ctx);
955 match ctx {
956 Some(tag) => owner.clone().with_context(tag.clone()).to_bytes(),
957 None => owner.to_bytes(),
958 }
959 };
960 for host in config.domains.exact_hosts() {
961 ops.push(WriteOp::Put(keys::domain(host), owner_for(host)));
962 }
963 for wildcard in &config.domains.wildcards {
964 if let Some(suffix) = wildcard.strip_prefix("*.") {
965 let bytes = match contexts.get(wildcard) {
967 Some(tag) => owner.clone().with_context(tag.clone()).to_bytes(),
968 None => owner.to_bytes(),
969 };
970 ops.push(WriteOp::Put(keys::wildcard(suffix), bytes));
971 }
972 }
973 self.kv.write_batch(ops).await?;
974 self.bump_domain_epoch();
977 Ok(())
978 }
979
980 async fn ensure_host_claimable(
986 &self,
987 key: &str,
988 label: &str,
989 owner: &DomainOwner,
990 ) -> Result<(), DeployError> {
991 if let Some(bytes) = self.kv.get(key).await? {
992 let held = DomainOwner::from_bytes(&bytes);
993 if &held != owner {
994 return Err(DeployError::Conflict(format!(
995 "{label} is already attached to site `{}` in project `{}`",
996 held.site, held.project
997 )));
998 }
999 }
1000 Ok(())
1001 }
1002
1003 pub async fn resolve_site_by_host(
1008 &self,
1009 host: &str,
1010 ) -> Result<Option<DomainOwner>, DeployError> {
1011 let host = canon_host(host);
1014 let host = host.as_str();
1015 let epoch = self.domain_epoch.load(std::sync::atomic::Ordering::Acquire);
1019 {
1020 let cache = self.domain_cache.read().unwrap();
1021 if cache.epoch == epoch {
1022 if let Some(owner) = cache.map.get(host) {
1023 return Ok(owner.clone());
1024 }
1025 }
1026 }
1027 let resolved = self.resolve_site_by_host_uncached(host).await?;
1028 {
1032 let mut cache = self.domain_cache.write().unwrap();
1033 let current = self.domain_epoch.load(std::sync::atomic::Ordering::Acquire);
1034 if cache.epoch != current {
1035 cache.epoch = current;
1036 cache.map.clear();
1037 }
1038 if current == epoch {
1039 if cache.map.len() >= 4096 {
1043 cache.map.clear();
1044 }
1045 cache.map.insert(host.to_string(), resolved.clone());
1046 }
1047 }
1048 Ok(resolved)
1049 }
1050
1051 async fn resolve_site_by_host_uncached(
1055 &self,
1056 host: &str,
1057 ) -> Result<Option<DomainOwner>, DeployError> {
1058 if let Some(bytes) = self.kv.get(&keys::domain(host)).await? {
1059 return Ok(Some(DomainOwner::from_bytes(&bytes)));
1060 }
1061 let mut rest = host;
1062 while let Some((_, parent)) = rest.split_once('.') {
1063 if let Some(bytes) = self.kv.get(&keys::wildcard(parent)).await? {
1064 return Ok(Some(DomainOwner::from_bytes(&bytes)));
1065 }
1066 rest = parent;
1067 }
1068 Ok(None)
1069 }
1070
1071 fn bump_domain_epoch(&self) {
1074 self.domain_epoch
1075 .fetch_add(1, std::sync::atomic::Ordering::Release);
1076 }
1077
1078 pub async fn all_sites(&self, project: ProjectRef<'_>) -> Result<Vec<String>, DeployError> {
1084 let mut sites = BTreeSet::new();
1085 for prefix in [
1086 keys::current_prefix(project),
1087 keys::site_prefix(project),
1088 keys::history_prefix(project),
1089 ] {
1090 for key in self.kv.list_prefix(&prefix).await? {
1091 if let Some(name) = key.strip_prefix(&prefix) {
1092 if !name.is_empty() {
1093 sites.insert(name.to_string());
1094 }
1095 }
1096 }
1097 }
1098 Ok(sites.into_iter().collect())
1099 }
1100
1101 pub async fn all_sites_all(&self) -> Result<Vec<(String, String)>, DeployError> {
1105 let mut out = Vec::new();
1106 for project in self.discover_projects().await? {
1107 for site in self.all_sites(ProjectRef::new(&project)).await? {
1108 out.push((project.clone(), site));
1109 }
1110 }
1111 Ok(out)
1112 }
1113
1114 pub async fn discover_projects(&self) -> Result<Vec<String>, DeployError> {
1119 let mut projects = BTreeSet::new();
1120 for key in self.kv.list_prefix(keys::PROJECT_ROOT).await? {
1121 if let Some(p) = keys::project_of_key(&key) {
1122 projects.insert(p.to_string());
1123 }
1124 }
1125 Ok(projects.into_iter().collect())
1126 }
1127
1128 pub async fn put_function(
1136 &self,
1137 project: ProjectRef<'_>,
1138 f: &crate::function::Function,
1139 ) -> Result<(), DeployError> {
1140 let bytes = serde_json::to_vec(f).map_err(|e| DeployError::Serde(e.to_string()))?;
1141 self.kv
1142 .put(
1143 &crate::function::keys::meta(project.as_str(), &f.name),
1144 bytes,
1145 )
1146 .await?;
1147 Ok(())
1148 }
1149
1150 pub async fn get_function(
1152 &self,
1153 project: ProjectRef<'_>,
1154 name: &str,
1155 ) -> Result<Option<crate::function::Function>, DeployError> {
1156 match self
1157 .kv
1158 .get(&crate::function::keys::meta(project.as_str(), name))
1159 .await?
1160 {
1161 Some(bytes) => Ok(Some(
1162 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
1163 )),
1164 None => Ok(None),
1165 }
1166 }
1167
1168 pub async fn list_stored_functions(
1174 &self,
1175 project: ProjectRef<'_>,
1176 ) -> Result<Vec<crate::function::Function>, DeployError> {
1177 let prefix = crate::function::keys::functions_prefix(project.as_str());
1178 let mut out = Vec::new();
1179 for key in self.kv.list_prefix(&prefix).await? {
1180 if key[prefix.len()..].contains('/') {
1181 continue;
1182 }
1183 if let Some(bytes) = self.kv.get(&key).await? {
1184 if let Ok(f) = serde_json::from_slice(&bytes) {
1185 out.push(f);
1186 }
1187 }
1188 }
1189 Ok(out)
1190 }
1191
1192 pub async fn delete_function(
1195 &self,
1196 project: ProjectRef<'_>,
1197 name: &str,
1198 ) -> Result<bool, DeployError> {
1199 let meta = crate::function::keys::meta(project.as_str(), name);
1200 let existed = self.kv.get(&meta).await?.is_some();
1201 for key in self.kv.list_prefix(&format!("{meta}/")).await? {
1208 self.kv.delete(&key).await?;
1209 }
1210 self.kv.delete(&meta).await?;
1211 self.kv
1212 .delete(&crate::function::keys::metering(project.as_str(), name))
1213 .await?;
1214 Ok(existed)
1215 }
1216
1217 pub async fn put_trigger(
1221 &self,
1222 project: ProjectRef<'_>,
1223 function: &str,
1224 trigger: &crate::function::FunctionTrigger,
1225 ) -> Result<(), DeployError> {
1226 let bytes = serde_json::to_vec(trigger).map_err(|e| DeployError::Serde(e.to_string()))?;
1227 self.kv
1228 .put(
1229 &crate::function::keys::trigger(project.as_str(), function, &trigger.id),
1230 bytes,
1231 )
1232 .await?;
1233 Ok(())
1234 }
1235
1236 pub async fn get_trigger(
1238 &self,
1239 project: ProjectRef<'_>,
1240 function: &str,
1241 id: &str,
1242 ) -> Result<Option<crate::function::FunctionTrigger>, DeployError> {
1243 match self
1244 .kv
1245 .get(&crate::function::keys::trigger(
1246 project.as_str(),
1247 function,
1248 id,
1249 ))
1250 .await?
1251 {
1252 Some(bytes) => Ok(Some(
1253 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
1254 )),
1255 None => Ok(None),
1256 }
1257 }
1258
1259 pub async fn list_triggers(
1261 &self,
1262 project: ProjectRef<'_>,
1263 function: &str,
1264 ) -> Result<Vec<crate::function::FunctionTrigger>, DeployError> {
1265 let prefix = crate::function::keys::triggers_prefix(project.as_str(), function);
1266 let mut out = Vec::new();
1267 for key in self.kv.list_prefix(&prefix).await? {
1268 if let Some(bytes) = self.kv.get(&key).await? {
1269 if let Ok(t) = serde_json::from_slice(&bytes) {
1270 out.push(t);
1271 }
1272 }
1273 }
1274 Ok(out)
1275 }
1276
1277 pub async fn delete_trigger(
1279 &self,
1280 project: ProjectRef<'_>,
1281 function: &str,
1282 id: &str,
1283 ) -> Result<bool, DeployError> {
1284 let key = crate::function::keys::trigger(project.as_str(), function, id);
1285 let existed = self.kv.get(&key).await?.is_some();
1286 self.kv.delete(&key).await?;
1287 Ok(existed)
1288 }
1289
1290 pub async fn put_invocation(
1294 &self,
1295 project: ProjectRef<'_>,
1296 inv: &crate::function::Invocation,
1297 ) -> Result<(), DeployError> {
1298 let bytes = serde_json::to_vec(inv).map_err(|e| DeployError::Serde(e.to_string()))?;
1299 self.kv
1300 .put(
1301 &crate::function::keys::invocation(project.as_str(), &inv.function, &inv.id),
1302 bytes,
1303 )
1304 .await?;
1305 Ok(())
1306 }
1307
1308 pub async fn get_invocation(
1310 &self,
1311 project: ProjectRef<'_>,
1312 function: &str,
1313 id: &str,
1314 ) -> Result<Option<crate::function::Invocation>, DeployError> {
1315 match self
1316 .kv
1317 .get(&crate::function::keys::invocation(
1318 project.as_str(),
1319 function,
1320 id,
1321 ))
1322 .await?
1323 {
1324 Some(bytes) => Ok(Some(
1325 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
1326 )),
1327 None => Ok(None),
1328 }
1329 }
1330
1331 pub async fn list_invocations(
1333 &self,
1334 project: ProjectRef<'_>,
1335 function: &str,
1336 ) -> Result<Vec<crate::function::Invocation>, DeployError> {
1337 let prefix = crate::function::keys::invocations_prefix(project.as_str(), function);
1338 let mut out = Vec::new();
1339 for key in self.kv.list_prefix(&prefix).await? {
1340 if let Some(bytes) = self.kv.get(&key).await? {
1341 if let Ok(inv) = serde_json::from_slice(&bytes) {
1342 out.push(inv);
1343 }
1344 }
1345 }
1346 Ok(out)
1347 }
1348
1349 pub async fn put_idempotency(
1352 &self,
1353 project: ProjectRef<'_>,
1354 function: &str,
1355 key: &str,
1356 invocation_id: &str,
1357 ) -> Result<(), DeployError> {
1358 self.kv
1359 .put(
1360 &crate::function::keys::idempotency(project.as_str(), function, key),
1361 invocation_id.as_bytes().to_vec(),
1362 )
1363 .await?;
1364 Ok(())
1365 }
1366
1367 pub async fn get_idempotency(
1369 &self,
1370 project: ProjectRef<'_>,
1371 function: &str,
1372 key: &str,
1373 ) -> Result<Option<String>, DeployError> {
1374 match self
1375 .kv
1376 .get(&crate::function::keys::idempotency(
1377 project.as_str(),
1378 function,
1379 key,
1380 ))
1381 .await?
1382 {
1383 Some(bytes) => Ok(Some(String::from_utf8_lossy(&bytes).into_owned())),
1384 None => Ok(None),
1385 }
1386 }
1387
1388 pub async fn get_metering(
1392 &self,
1393 project: ProjectRef<'_>,
1394 function: &str,
1395 ) -> Result<Option<crate::function::Metering>, DeployError> {
1396 match self
1397 .kv
1398 .get(&crate::function::keys::metering(project.as_str(), function))
1399 .await?
1400 {
1401 Some(bytes) => Ok(Some(
1402 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
1403 )),
1404 None => Ok(None),
1405 }
1406 }
1407
1408 pub async fn put_metering(
1410 &self,
1411 project: ProjectRef<'_>,
1412 metering: &crate::function::Metering,
1413 ) -> Result<(), DeployError> {
1414 let bytes = serde_json::to_vec(metering).map_err(|e| DeployError::Serde(e.to_string()))?;
1415 self.kv
1416 .put(
1417 &crate::function::keys::metering(project.as_str(), &metering.function),
1418 bytes,
1419 )
1420 .await?;
1421 Ok(())
1422 }
1423
1424 pub async fn list_metering(
1427 &self,
1428 project: ProjectRef<'_>,
1429 ) -> Result<Vec<crate::function::Metering>, DeployError> {
1430 let prefix = crate::function::keys::metering_prefix(project.as_str());
1431 let mut out = Vec::new();
1432 for key in self.kv.list_prefix(&prefix).await? {
1433 if let Some(bytes) = self.kv.get(&key).await? {
1434 if let Ok(m) = serde_json::from_slice(&bytes) {
1435 out.push(m);
1436 }
1437 }
1438 }
1439 Ok(out)
1440 }
1441
1442 pub async fn put_managed_notification(
1447 &self,
1448 project: ProjectRef<'_>,
1449 record: &crate::blob_notify::ManagedNotification,
1450 ) -> Result<(), DeployError> {
1451 let bytes = record
1452 .to_json()
1453 .map_err(|e| DeployError::Serde(e.to_string()))?;
1454 self.kv
1455 .put(
1456 &crate::blob_notify::blobnotify_key(
1457 project.as_str(),
1458 &record.function,
1459 &record.prefix,
1460 ),
1461 bytes,
1462 )
1463 .await?;
1464 Ok(())
1465 }
1466
1467 pub async fn get_managed_notification(
1469 &self,
1470 project: ProjectRef<'_>,
1471 function: &str,
1472 prefix: &str,
1473 ) -> Result<Option<crate::blob_notify::ManagedNotification>, DeployError> {
1474 match self
1475 .kv
1476 .get(&crate::blob_notify::blobnotify_key(
1477 project.as_str(),
1478 function,
1479 prefix,
1480 ))
1481 .await?
1482 {
1483 Some(bytes) => Ok(Some(
1484 crate::blob_notify::ManagedNotification::from_json(&bytes)
1485 .map_err(|e| DeployError::Serde(e.to_string()))?,
1486 )),
1487 None => Ok(None),
1488 }
1489 }
1490
1491 pub async fn list_managed_notifications(
1493 &self,
1494 project: ProjectRef<'_>,
1495 function: &str,
1496 ) -> Result<Vec<crate::blob_notify::ManagedNotification>, DeployError> {
1497 let prefix = crate::blob_notify::blobnotify_function_prefix(project.as_str(), function);
1498 let mut out = Vec::new();
1499 for key in self.kv.list_prefix(&prefix).await? {
1500 if let Some(bytes) = self.kv.get(&key).await? {
1501 if let Ok(record) = crate::blob_notify::ManagedNotification::from_json(&bytes) {
1502 out.push(record);
1503 }
1504 }
1505 }
1506 Ok(out)
1507 }
1508
1509 pub async fn remove_managed_notification(
1511 &self,
1512 project: ProjectRef<'_>,
1513 function: &str,
1514 prefix: &str,
1515 ) -> Result<(), DeployError> {
1516 self.kv
1517 .delete(&crate::blob_notify::blobnotify_key(
1518 project.as_str(),
1519 function,
1520 prefix,
1521 ))
1522 .await?;
1523 Ok(())
1524 }
1525
1526 pub async fn put_workflow(
1530 &self,
1531 project: ProjectRef<'_>,
1532 workflow: &crate::workflow::Workflow,
1533 ) -> Result<(), DeployError> {
1534 let bytes = serde_json::to_vec(workflow).map_err(|e| DeployError::Serde(e.to_string()))?;
1535 self.kv
1536 .put(
1537 &crate::workflow::keys::definition(project.as_str(), &workflow.name),
1538 bytes,
1539 )
1540 .await?;
1541 Ok(())
1542 }
1543
1544 pub async fn get_workflow(
1546 &self,
1547 project: ProjectRef<'_>,
1548 name: &str,
1549 ) -> Result<Option<crate::workflow::Workflow>, DeployError> {
1550 match self
1551 .kv
1552 .get(&crate::workflow::keys::definition(project.as_str(), name))
1553 .await?
1554 {
1555 Some(bytes) => Ok(Some(
1556 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
1557 )),
1558 None => Ok(None),
1559 }
1560 }
1561
1562 pub async fn list_workflows(
1566 &self,
1567 project: ProjectRef<'_>,
1568 ) -> Result<Vec<crate::workflow::Workflow>, DeployError> {
1569 let prefix = crate::workflow::keys::definitions_prefix(project.as_str());
1570 let mut out = Vec::new();
1571 for key in self.kv.list_prefix(&prefix).await? {
1572 if key[prefix.len()..].contains('/') {
1573 continue;
1574 }
1575 if let Some(bytes) = self.kv.get(&key).await? {
1576 if let Ok(w) = serde_json::from_slice(&bytes) {
1577 out.push(w);
1578 }
1579 }
1580 }
1581 Ok(out)
1582 }
1583
1584 pub async fn delete_workflow(
1587 &self,
1588 project: ProjectRef<'_>,
1589 name: &str,
1590 ) -> Result<bool, DeployError> {
1591 let key = crate::workflow::keys::definition(project.as_str(), name);
1592 let existed = self.kv.get(&key).await?.is_some();
1593 self.kv.delete(&key).await?;
1594 Ok(existed)
1595 }
1596
1597 pub async fn put_workflow_run(
1599 &self,
1600 project: ProjectRef<'_>,
1601 run: &crate::workflow::WorkflowRun,
1602 ) -> Result<(), DeployError> {
1603 let bytes = serde_json::to_vec(run).map_err(|e| DeployError::Serde(e.to_string()))?;
1604 self.kv
1605 .put(
1606 &crate::workflow::keys::run(project.as_str(), &run.workflow, &run.id),
1607 bytes,
1608 )
1609 .await?;
1610 Ok(())
1611 }
1612
1613 pub async fn get_workflow_run(
1615 &self,
1616 project: ProjectRef<'_>,
1617 workflow: &str,
1618 id: &str,
1619 ) -> Result<Option<crate::workflow::WorkflowRun>, DeployError> {
1620 match self
1621 .kv
1622 .get(&crate::workflow::keys::run(project.as_str(), workflow, id))
1623 .await?
1624 {
1625 Some(bytes) => Ok(Some(
1626 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
1627 )),
1628 None => Ok(None),
1629 }
1630 }
1631
1632 pub async fn list_workflow_runs(
1634 &self,
1635 project: ProjectRef<'_>,
1636 workflow: &str,
1637 ) -> Result<Vec<crate::workflow::WorkflowRun>, DeployError> {
1638 let prefix = crate::workflow::keys::runs_prefix(project.as_str(), workflow);
1639 let mut out = Vec::new();
1640 for key in self.kv.list_prefix(&prefix).await? {
1641 if let Some(bytes) = self.kv.get(&key).await? {
1642 if let Ok(run) = serde_json::from_slice(&bytes) {
1643 out.push(run);
1644 }
1645 }
1646 }
1647 Ok(out)
1648 }
1649
1650 pub async fn get_domain_verification(
1652 &self,
1653 project: ProjectRef<'_>,
1654 site: &SiteName,
1655 host: &str,
1656 ) -> Result<Option<DomainVerification>, DeployError> {
1657 match self
1658 .kv
1659 .get(&keys::domain_verification(project, site.as_str(), host))
1660 .await?
1661 {
1662 Some(bytes) => Ok(Some(DomainVerification::from_json(&bytes)?)),
1663 None => Ok(None),
1664 }
1665 }
1666
1667 pub async fn list_domain_verifications(
1669 &self,
1670 project: ProjectRef<'_>,
1671 site: &SiteName,
1672 ) -> Result<Vec<DomainVerification>, DeployError> {
1673 let prefix = keys::domain_verification_prefix(project, site.as_str());
1674 let mut out = Vec::new();
1675 for key in self.kv.list_prefix(&prefix).await? {
1676 if let Some(bytes) = self.kv.get(&key).await? {
1677 out.push(DomainVerification::from_json(&bytes)?);
1678 }
1679 }
1680 out.sort_by(|a, b| a.host.cmp(&b.host));
1681 Ok(out)
1682 }
1683
1684 pub async fn list_all_domain_verifications(
1691 &self,
1692 ) -> Result<Vec<(String, String, DomainVerification)>, DeployError> {
1693 let mut out = Vec::new();
1694 for project in self.discover_projects().await? {
1695 let pref = ProjectRef::new(&project);
1696 let scan = keys::domain_verification_project_prefix(pref);
1697 for key in self.kv.list_prefix(&scan).await? {
1698 let Some((site, _host)) = key
1701 .strip_prefix(&scan)
1702 .and_then(|rest| rest.split_once('/'))
1703 else {
1704 continue;
1705 };
1706 if let Some(bytes) = self.kv.get(&key).await? {
1707 out.push((
1708 project.clone(),
1709 site.to_string(),
1710 DomainVerification::from_json(&bytes)?,
1711 ));
1712 }
1713 }
1714 }
1715 Ok(out)
1716 }
1717
1718 pub async fn find_pending_http_challenge(
1727 &self,
1728 host: &str,
1729 token: &str,
1730 now_unix: u64,
1731 ) -> Result<Option<DomainVerification>, DeployError> {
1732 let host = crate::domain_verify::normalize_host(host);
1733 let Some(owner_bytes) = self
1738 .kv
1739 .get(&keys::http_challenge_index(&host, token))
1740 .await?
1741 else {
1742 return Ok(None);
1743 };
1744 let owner = DomainOwner::from_bytes(&owner_bytes);
1745 let site = SiteName::new(owner.site);
1746 let Some(v) = self
1747 .get_domain_verification(ProjectRef::new(&owner.project), &site, &host)
1748 .await?
1749 else {
1750 return Ok(None);
1751 };
1752 if v.method == VerificationMethod::Http
1753 && v.host == host
1754 && v.matches(token)
1755 && !v.is_expired(now_unix)
1756 {
1757 Ok(Some(v))
1758 } else {
1759 Ok(None)
1760 }
1761 }
1762
1763 async fn put_domain_verification(
1764 &self,
1765 project: ProjectRef<'_>,
1766 site: &SiteName,
1767 verification: &DomainVerification,
1768 ) -> Result<(), DeployError> {
1769 let mut ops = vec![WriteOp::Put(
1770 keys::domain_verification(project, site.as_str(), &verification.host),
1771 verification.to_json()?,
1772 )];
1773 if verification.method == VerificationMethod::Http {
1779 ops.push(WriteOp::Put(
1780 keys::http_challenge_index(&verification.host, &verification.token),
1781 DomainOwner::new(project.as_str(), site.as_str()).to_bytes(),
1782 ));
1783 }
1784 self.kv.write_batch(ops).await?;
1785 Ok(())
1786 }
1787
1788 pub async fn get_managed_dns(
1792 &self,
1793 project: ProjectRef<'_>,
1794 site: &SiteName,
1795 host: &str,
1796 ) -> Result<Option<crate::dns_managed::ManagedDns>, DeployError> {
1797 match self
1798 .kv
1799 .get(&crate::dns_managed::dnsmanaged_key(
1800 project.as_str(),
1801 site.as_str(),
1802 host,
1803 ))
1804 .await?
1805 {
1806 Some(bytes) => Ok(Some(crate::dns_managed::ManagedDns::from_json(&bytes)?)),
1807 None => Ok(None),
1808 }
1809 }
1810
1811 pub async fn set_managed_dns(
1813 &self,
1814 project: ProjectRef<'_>,
1815 site: &SiteName,
1816 ledger: &crate::dns_managed::ManagedDns,
1817 ) -> Result<(), DeployError> {
1818 self.kv
1819 .put(
1820 &crate::dns_managed::dnsmanaged_key(project.as_str(), site.as_str(), &ledger.host),
1821 ledger.to_json()?,
1822 )
1823 .await?;
1824 Ok(())
1825 }
1826
1827 pub async fn remove_managed_dns(
1829 &self,
1830 project: ProjectRef<'_>,
1831 site: &SiteName,
1832 host: &str,
1833 ) -> Result<(), DeployError> {
1834 self.kv
1835 .delete(&crate::dns_managed::dnsmanaged_key(
1836 project.as_str(),
1837 site.as_str(),
1838 host,
1839 ))
1840 .await?;
1841 Ok(())
1842 }
1843
1844 pub async fn list_managed_dns(
1847 &self,
1848 project: ProjectRef<'_>,
1849 site: &SiteName,
1850 ) -> Result<Vec<crate::dns_managed::ManagedDns>, DeployError> {
1851 let prefix = crate::dns_managed::dnsmanaged_site_prefix(project.as_str(), site.as_str());
1852 let mut out = Vec::new();
1853 for key in self.kv.list_prefix(&prefix).await? {
1854 if let Some(bytes) = self.kv.get(&key).await? {
1855 out.push(crate::dns_managed::ManagedDns::from_json(&bytes)?);
1856 }
1857 }
1858 out.sort_by(|a, b| a.host.cmp(&b.host));
1859 Ok(out)
1860 }
1861
1862 pub async fn start_domain_verification(
1869 &self,
1870 project: ProjectRef<'_>,
1871 site: &SiteName,
1872 host: &str,
1873 method: VerificationMethod,
1874 now_unix: u64,
1875 ) -> Result<DomainVerification, DeployError> {
1876 if let Some(existing) = self.get_domain_verification(project, site, host).await? {
1877 if existing.verified || existing.method == method {
1878 return Ok(existing);
1879 }
1880 }
1881 const MAX_PENDING_VERIFICATIONS_PER_SITE: usize = 64;
1888 let pending = self
1889 .list_domain_verifications(project, site)
1890 .await?
1891 .into_iter()
1892 .filter(|v| !v.verified && !v.is_expired(now_unix))
1893 .count();
1894 if pending >= MAX_PENDING_VERIFICATIONS_PER_SITE {
1895 return Err(DeployError::Conflict(format!(
1896 "too many pending domain verifications for site {site} \
1897 (max {MAX_PENDING_VERIFICATIONS_PER_SITE}); verify or remove some first"
1898 )));
1899 }
1900 let verification = DomainVerification::new(host, method, now_unix);
1901 self.put_domain_verification(project, site, &verification)
1902 .await?;
1903 Ok(verification)
1904 }
1905
1906 pub async fn is_domain_verified(
1908 &self,
1909 project: ProjectRef<'_>,
1910 site: &SiteName,
1911 host: &str,
1912 ) -> Result<bool, DeployError> {
1913 Ok(self
1914 .get_domain_verification(project, site, host)
1915 .await?
1916 .is_some_and(|v| v.verified))
1917 }
1918
1919 pub async fn mark_domain_verified(
1922 &self,
1923 project: ProjectRef<'_>,
1924 site: &SiteName,
1925 host: &str,
1926 ) -> Result<DomainVerification, DeployError> {
1927 let mut verification = self
1928 .get_domain_verification(project, site, host)
1929 .await?
1930 .ok_or_else(|| {
1931 DeployError::NotFound(format!("no verification challenge for {host}"))
1932 })?;
1933 verification.verified = true;
1934 self.put_domain_verification(project, site, &verification)
1935 .await?;
1936 Ok(verification)
1937 }
1938
1939 pub async fn remove_domain_verification(
1942 &self,
1943 project: ProjectRef<'_>,
1944 site: &SiteName,
1945 host: &str,
1946 ) -> Result<bool, DeployError> {
1947 let Some(v) = self.get_domain_verification(project, site, host).await? else {
1948 return Ok(false);
1949 };
1950 let mut ops = vec![WriteOp::Delete(keys::domain_verification(
1951 project,
1952 site.as_str(),
1953 host,
1954 ))];
1955 if v.method == VerificationMethod::Http {
1956 ops.push(WriteOp::Delete(keys::http_challenge_index(
1957 &v.host, &v.token,
1958 )));
1959 }
1960 self.kv.write_batch(ops).await?;
1961 Ok(true)
1962 }
1963
1964 pub async fn attach_verified_domain(
1975 &self,
1976 project: ProjectRef<'_>,
1977 site: &SiteName,
1978 host: &str,
1979 ) -> Result<SiteConfig, DeployError> {
1980 if !self.is_domain_verified(project, site, host).await? {
1981 return Err(DeployError::NotFound(format!(
1982 "{host} is not verified for {site}; run domain verification first"
1983 )));
1984 }
1985 if host.trim_start().starts_with("*.")
1989 && self
1990 .get_domain_verification(project, site, host)
1991 .await?
1992 .map(|v| v.method)
1993 != Some(VerificationMethod::Dns)
1994 {
1995 return Err(DeployError::Conflict(format!(
1996 "wildcard {host} must be verified via DNS \
1997 (an HTTP token proves only the base host, not the subtree)"
1998 )));
1999 }
2000 let _claim = self.domain_claim_lock.lock().await;
2004 let mut config = self
2005 .get_site_config(project, site.as_str())
2006 .await?
2007 .unwrap_or_default();
2008 let domains = &mut config.domains;
2009 if let Some(suffix) = host.strip_prefix("*.") {
2010 let wildcard = format!("*.{}", suffix.trim_end_matches('.').to_ascii_lowercase());
2011 if !domains.wildcards.contains(&wildcard) {
2012 domains.wildcards.push(wildcard);
2013 }
2014 } else {
2015 let host = host.trim().trim_end_matches('.').to_ascii_lowercase();
2016 if domains.primary.is_none() {
2017 domains.primary = Some(host);
2018 } else if domains.primary.as_deref() != Some(host.as_str())
2019 && !domains.aliases.contains(&host)
2020 {
2021 domains.aliases.push(host);
2022 }
2023 }
2024 self.set_site_config_locked(project, site.as_str(), &config)
2025 .await?;
2026 Ok(config)
2027 }
2028
2029 pub async fn set_alias(
2035 &self,
2036 project: ProjectRef<'_>,
2037 site: &str,
2038 name: &str,
2039 id: &str,
2040 ) -> Result<(), DeployError> {
2041 let manifest = self
2042 .get_manifest(id)
2043 .await?
2044 .ok_or_else(|| DeployError::NotFound(format!("deployment {id}")))?;
2045 let missing = self.missing_blobs(&manifest).await?;
2046 if !missing.is_empty() {
2047 return Err(DeployError::Incomplete(missing));
2048 }
2049 self.kv
2050 .put(&keys::alias(project, site, name), id.as_bytes().to_vec())
2051 .await?;
2052 Ok(())
2053 }
2054
2055 pub async fn get_alias(
2057 &self,
2058 project: ProjectRef<'_>,
2059 site: &str,
2060 name: &str,
2061 ) -> Result<Option<String>, DeployError> {
2062 match self.kv.get(&keys::alias(project, site, name)).await? {
2063 Some(bytes) => Ok(Some(String::from_utf8_lossy(&bytes).into_owned())),
2064 None => Ok(None),
2065 }
2066 }
2067
2068 pub async fn remove_alias(
2070 &self,
2071 project: ProjectRef<'_>,
2072 site: &str,
2073 name: &str,
2074 ) -> Result<bool, DeployError> {
2075 let key = keys::alias(project, site, name);
2076 let existed = self.kv.get(&key).await?.is_some();
2077 if existed {
2078 self.kv.delete(&key).await?;
2079 }
2080 Ok(existed)
2081 }
2082
2083 pub async fn list_aliases(
2085 &self,
2086 project: ProjectRef<'_>,
2087 site: &str,
2088 ) -> Result<BTreeMap<String, String>, DeployError> {
2089 let prefix = keys::alias_prefix(project, site);
2090 let mut out = BTreeMap::new();
2091 for key in self.kv.list_prefix(&prefix).await? {
2092 if let Some(bytes) = self.kv.get(&key).await? {
2093 let name = key.strip_prefix(&prefix).unwrap_or(&key).to_string();
2094 out.insert(name, String::from_utf8_lossy(&bytes).into_owned());
2095 }
2096 }
2097 Ok(out)
2098 }
2099
2100 pub async fn put_token_meta(&self, meta: &crate::authz::TokenMeta) -> Result<(), DeployError> {
2105 self.kv
2106 .put(
2107 &crate::authz::token_meta_key(&meta.revocation_id),
2108 serde_json::to_vec(meta)?,
2109 )
2110 .await?;
2111 Ok(())
2112 }
2113
2114 pub async fn list_token_meta(&self) -> Result<Vec<crate::authz::TokenMeta>, DeployError> {
2116 let mut out = Vec::new();
2117 for key in self.kv.list_prefix(crate::authz::TOKEN_META_PREFIX).await? {
2118 if let Some(bytes) = self.kv.get(&key).await? {
2119 if let Ok(meta) = serde_json::from_slice::<crate::authz::TokenMeta>(&bytes) {
2120 out.push(meta);
2121 }
2122 }
2123 }
2124 Ok(out)
2125 }
2126
2127 pub async fn revoke_token(&self, id_or_prefix: &str) -> Result<bool, DeployError> {
2132 let ids: Vec<String> = self
2133 .list_token_meta()
2134 .await?
2135 .into_iter()
2136 .map(|m| m.revocation_id)
2137 .collect();
2138 let matches: Vec<&String> = ids
2139 .iter()
2140 .filter(|id| id.starts_with(id_or_prefix))
2141 .collect();
2142 if let [id] = matches.as_slice() {
2143 self.kv
2144 .put(&crate::authz::revoked_key(id), Vec::new())
2145 .await?;
2146 self.kv.delete(&crate::authz::token_meta_key(id)).await?;
2147 Ok(true)
2148 } else {
2149 Ok(false)
2150 }
2151 }
2152
2153 pub async fn bootstrap_consumed(&self, secret_hash: &str) -> Result<bool, DeployError> {
2158 Ok(self
2159 .kv
2160 .get(&crate::authz::bootstrap_key(secret_hash))
2161 .await?
2162 .is_some())
2163 }
2164
2165 pub async fn mark_bootstrap_consumed(&self, secret_hash: &str) -> Result<(), DeployError> {
2167 self.kv
2168 .put(&crate::authz::bootstrap_key(secret_hash), Vec::new())
2169 .await?;
2170 Ok(())
2171 }
2172
2173 pub async fn get_authz_policy(&self) -> Result<Option<crate::authz::AuthzPolicy>, DeployError> {
2176 match self.kv.get(crate::authz::POLICY_KEY).await? {
2177 Some(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
2178 None => Ok(None),
2179 }
2180 }
2181
2182 pub async fn set_authz_policy(
2186 &self,
2187 policy: &crate::authz::AuthzPolicy,
2188 ) -> Result<(), DeployError> {
2189 self.kv
2190 .put(crate::authz::POLICY_KEY, serde_json::to_vec(policy)?)
2191 .await?;
2192 Ok(())
2193 }
2194
2195 pub async fn add_root_anchor(&self, pubkey: &str) -> Result<(), DeployError> {
2200 self.kv
2201 .put(&crate::authz::root_anchor_key(pubkey), Vec::new())
2202 .await?;
2203 Ok(())
2204 }
2205
2206 pub async fn remove_root_anchor(&self, pubkey: &str) -> Result<(), DeployError> {
2208 self.kv
2209 .delete(&crate::authz::root_anchor_key(pubkey))
2210 .await?;
2211 Ok(())
2212 }
2213
2214 pub async fn list_root_anchors(&self) -> Result<Vec<String>, DeployError> {
2216 Ok(self
2217 .kv
2218 .list_prefix(crate::authz::ROOT_ANCHOR_PREFIX)
2219 .await?
2220 .iter()
2221 .filter_map(|k| {
2222 k.strip_prefix(crate::authz::ROOT_ANCHOR_PREFIX)
2223 .map(String::from)
2224 })
2225 .collect())
2226 }
2227
2228 const DAEMON_CURRENT_KEY: &'static str = "daemon/current";
2232 const DAEMON_HISTORY_KEY: &'static str = "daemon/history";
2234 const DAEMON_HISTORY_MAX: usize = 20;
2236
2237 pub async fn daemon_config_generation(&self) -> Result<Option<String>, DeployError> {
2240 Ok(self
2241 .kv
2242 .get(Self::DAEMON_CURRENT_KEY)
2243 .await?
2244 .map(|b| String::from_utf8_lossy(&b).into_owned()))
2245 }
2246
2247 pub async fn get_daemon_config(
2250 &self,
2251 ) -> Result<Option<crate::daemon_config::DaemonConfig>, DeployError> {
2252 let Some(hash) = self.daemon_config_generation().await? else {
2253 return Ok(None);
2254 };
2255 match self.kv.get(&keys::daemon_config_blob(&hash)).await? {
2256 Some(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
2257 None => Ok(None),
2259 }
2260 }
2261
2262 pub async fn daemon_config_history(&self) -> Result<Vec<String>, DeployError> {
2265 match self.kv.get(Self::DAEMON_HISTORY_KEY).await? {
2266 Some(bytes) => Ok(serde_json::from_slice(&bytes)?),
2267 None => Ok(Vec::new()),
2268 }
2269 }
2270
2271 pub async fn set_daemon_config(
2277 &self,
2278 config: &crate::daemon_config::DaemonConfig,
2279 ) -> Result<String, DeployError> {
2280 let body = serde_json::to_vec(config)?;
2281 let hash = sha256_hex(&body);
2282 let mut history = self.daemon_config_history().await?;
2283 if let Some(current) = self.daemon_config_generation().await? {
2284 if current != hash {
2285 history.push(current);
2286 if history.len() > Self::DAEMON_HISTORY_MAX {
2287 let overflow = history.len() - Self::DAEMON_HISTORY_MAX;
2288 history.drain(0..overflow);
2289 }
2290 }
2291 }
2292 let ops = vec![
2293 WriteOp::Put(keys::daemon_config_blob(&hash), body),
2294 WriteOp::Put(
2295 Self::DAEMON_HISTORY_KEY.to_string(),
2296 serde_json::to_vec(&history)?,
2297 ),
2298 WriteOp::Put(
2299 Self::DAEMON_CURRENT_KEY.to_string(),
2300 hash.clone().into_bytes(),
2301 ),
2302 ];
2303 self.kv.write_batch(ops).await?;
2304 Ok(hash)
2305 }
2306
2307 pub async fn rollback_daemon_config(&self) -> Result<Option<String>, DeployError> {
2312 let mut history = self.daemon_config_history().await?;
2313 let Some(prev) = history.pop() else {
2314 return Ok(None);
2315 };
2316 let ops = vec![
2317 WriteOp::Put(
2318 Self::DAEMON_HISTORY_KEY.to_string(),
2319 serde_json::to_vec(&history)?,
2320 ),
2321 WriteOp::Put(
2322 Self::DAEMON_CURRENT_KEY.to_string(),
2323 prev.clone().into_bytes(),
2324 ),
2325 ];
2326 self.kv.write_batch(ops).await?;
2327 Ok(Some(prev))
2328 }
2329
2330 pub async fn put_compute_spec(
2335 &self,
2336 spec: &crate::compute::ComputeSpec,
2337 ) -> Result<String, DeployError> {
2338 let id = spec.id();
2339 self.kv
2340 .put(&crate::compute::spec_key(&id), serde_json::to_vec(spec)?)
2341 .await?;
2342 Ok(id)
2343 }
2344
2345 pub async fn get_compute_spec(
2347 &self,
2348 hash: &str,
2349 ) -> Result<Option<crate::compute::ComputeSpec>, DeployError> {
2350 match self.kv.get(&crate::compute::spec_key(hash)).await? {
2351 Some(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
2352 None => Ok(None),
2353 }
2354 }
2355
2356 pub async fn set_compute_workload(
2360 &self,
2361 project: ProjectRef<'_>,
2362 workload: &crate::compute::ComputeWorkload,
2363 ) -> Result<(), DeployError> {
2364 self.kv
2365 .put(
2366 &crate::compute::workload_key(project.as_str(), &workload.name),
2367 serde_json::to_vec(workload)?,
2368 )
2369 .await?;
2370 Ok(())
2371 }
2372
2373 pub async fn get_compute_workload(
2375 &self,
2376 project: ProjectRef<'_>,
2377 name: &str,
2378 ) -> Result<Option<crate::compute::ComputeWorkload>, DeployError> {
2379 match self
2380 .kv
2381 .get(&crate::compute::workload_key(project.as_str(), name))
2382 .await?
2383 {
2384 Some(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)),
2385 None => Ok(None),
2386 }
2387 }
2388
2389 pub async fn list_compute_workloads(
2391 &self,
2392 project: ProjectRef<'_>,
2393 ) -> Result<Vec<crate::compute::ComputeWorkload>, DeployError> {
2394 let prefix = crate::compute::workloads_prefix(project.as_str());
2395 let mut out = Vec::new();
2396 for key in self.kv.list_prefix(&prefix).await? {
2397 if let Some(bytes) = self.kv.get(&key).await? {
2398 if let Ok(w) = serde_json::from_slice::<crate::compute::ComputeWorkload>(&bytes) {
2399 out.push(w);
2400 }
2401 }
2402 }
2403 Ok(out)
2404 }
2405
2406 pub async fn list_compute_workloads_all(
2409 &self,
2410 ) -> Result<Vec<(String, crate::compute::ComputeWorkload)>, DeployError> {
2411 let mut out = Vec::new();
2412 for project in self.discover_projects().await? {
2413 for w in self
2414 .list_compute_workloads(ProjectRef::new(&project))
2415 .await?
2416 {
2417 out.push((project.clone(), w));
2418 }
2419 }
2420 Ok(out)
2421 }
2422
2423 pub async fn delete_compute_workload(
2426 &self,
2427 project: ProjectRef<'_>,
2428 name: &str,
2429 ) -> Result<bool, DeployError> {
2430 let key = crate::compute::workload_key(project.as_str(), name);
2431 let existed = self.kv.get(&key).await?.is_some();
2432 if existed {
2433 self.kv.delete(&key).await?;
2434 }
2435 Ok(existed)
2436 }
2437
2438 pub async fn set_replica_state(
2442 &self,
2443 project: ProjectRef<'_>,
2444 state: &crate::compute::ObservedInstance,
2445 ) -> Result<(), DeployError> {
2446 self.kv
2447 .put(
2448 &crate::compute::replica_state_key(
2449 project.as_str(),
2450 &state.handle.workload,
2451 state.handle.replica,
2452 ),
2453 serde_json::to_vec(state)?,
2454 )
2455 .await?;
2456 Ok(())
2457 }
2458
2459 pub async fn list_replica_states(
2464 &self,
2465 project: ProjectRef<'_>,
2466 workload: &str,
2467 ) -> Result<Vec<crate::compute::ObservedInstance>, DeployError> {
2468 let mut out = Vec::new();
2469 for key in self
2470 .kv
2471 .list_prefix(&crate::compute::replica_state_prefix(
2472 project.as_str(),
2473 workload,
2474 ))
2475 .await?
2476 {
2477 if let Some(bytes) = self.kv.get(&key).await? {
2478 if let Ok(mut state) =
2479 serde_json::from_slice::<crate::compute::ObservedInstance>(&bytes)
2480 {
2481 backfill_replica_project(&mut state, project.as_str());
2482 out.push(state);
2483 }
2484 }
2485 }
2486 Ok(out)
2487 }
2488
2489 pub async fn list_all_replica_states(
2493 &self,
2494 ) -> Result<Vec<crate::compute::ObservedInstance>, DeployError> {
2495 let mut out = Vec::new();
2496 for project in self.discover_projects().await? {
2497 let prefix = crate::compute::replica_states_project_prefix(&project);
2498 for key in self.kv.list_prefix(&prefix).await? {
2499 if let Some(bytes) = self.kv.get(&key).await? {
2500 if let Ok(mut state) =
2501 serde_json::from_slice::<crate::compute::ObservedInstance>(&bytes)
2502 {
2503 backfill_replica_project(&mut state, &project);
2507 out.push(state);
2508 }
2509 }
2510 }
2511 }
2512 Ok(out)
2513 }
2514
2515 pub async fn delete_replica_state(
2517 &self,
2518 project: ProjectRef<'_>,
2519 workload: &str,
2520 replica: u32,
2521 ) -> Result<(), DeployError> {
2522 self.kv
2523 .delete(&crate::compute::replica_state_key(
2524 project.as_str(),
2525 workload,
2526 replica,
2527 ))
2528 .await?;
2529 Ok(())
2530 }
2531
2532 pub async fn put_project(&self, p: &crate::project::Project) -> Result<String, DeployError> {
2544 let hash = p.id();
2545 let body = serde_json::to_vec(p).map_err(|e| DeployError::Serde(e.to_string()))?;
2546 let pointer = crate::project::pointer_key(&p.name);
2547 let mut history: Vec<String> =
2550 match self.kv.get(&crate::project::history_key(&p.name)).await? {
2551 Some(bytes) => serde_json::from_slice(&bytes).unwrap_or_default(),
2552 None => Vec::new(),
2553 };
2554 if let Some(current) = self.kv.get(&pointer).await? {
2555 let current = String::from_utf8_lossy(¤t).into_owned();
2556 if current != hash {
2557 history.retain(|h| h != ¤t);
2558 history.insert(0, current);
2559 history.truncate(MAX_HISTORY);
2560 }
2561 }
2562 self.kv
2563 .write_batch(vec![
2564 WriteOp::Put(crate::project::spec_key(&hash), body),
2565 WriteOp::Put(pointer, hash.clone().into_bytes()),
2566 WriteOp::Put(
2567 crate::project::history_key(&p.name),
2568 serde_json::to_vec(&history).map_err(|e| DeployError::Serde(e.to_string()))?,
2569 ),
2570 ])
2571 .await?;
2572 Ok(hash)
2573 }
2574
2575 pub fn default_project_record() -> crate::project::Project {
2580 crate::project::Project {
2581 version: crate::SCHEMA_VERSION,
2582 name: crate::project::DEFAULT_PROJECT.to_string(),
2583 created_at: now_unix(),
2584 meta: crate::project::ProjectMeta::default(),
2585 config: crate::project::ProjectConfig::default(),
2586 secrets_ref: None,
2587 }
2588 }
2589
2590 pub async fn ensure_default_project(&self) -> Result<bool, DeployError> {
2598 let pointer = crate::project::pointer_key(crate::project::DEFAULT_PROJECT);
2599 if self.kv.get(&pointer).await?.is_some() {
2600 return Ok(false);
2601 }
2602 let default = Self::default_project_record();
2603 let hash = default.id();
2604 let body = serde_json::to_vec(&default).map_err(|e| DeployError::Serde(e.to_string()))?;
2605 self.kv
2606 .write_batch(vec![
2607 WriteOp::Put(crate::project::spec_key(&hash), body),
2608 WriteOp::Put(pointer, hash.into_bytes()),
2609 ])
2610 .await?;
2611 Ok(true)
2612 }
2613
2614 pub async fn project_exists(&self, name: &str) -> Result<bool, DeployError> {
2619 if name == crate::project::DEFAULT_PROJECT {
2620 return Ok(true);
2621 }
2622 Ok(self
2623 .kv
2624 .get(&crate::project::pointer_key(name))
2625 .await?
2626 .is_some())
2627 }
2628
2629 pub async fn get_project(
2631 &self,
2632 name: &str,
2633 ) -> Result<Option<crate::project::Project>, DeployError> {
2634 let Some(hash) = self.kv.get(&crate::project::pointer_key(name)).await? else {
2635 if name == crate::project::DEFAULT_PROJECT {
2639 return Ok(Some(Self::default_project_record()));
2640 }
2641 return Ok(None);
2642 };
2643 let hash = String::from_utf8_lossy(&hash).into_owned();
2644 match self.kv.get(&crate::project::spec_key(&hash)).await? {
2645 Some(bytes) => Ok(Some(
2646 serde_json::from_slice(&bytes).map_err(|e| DeployError::Serde(e.to_string()))?,
2647 )),
2648 None if name == crate::project::DEFAULT_PROJECT => {
2651 Ok(Some(Self::default_project_record()))
2652 }
2653 None => Ok(None),
2654 }
2655 }
2656
2657 pub async fn list_projects(&self) -> Result<Vec<crate::project::Project>, DeployError> {
2661 let mut out = Vec::new();
2662 for key in self.kv.list_prefix(crate::project::POINTER_PREFIX).await? {
2663 let Some(name) = key.strip_prefix(crate::project::POINTER_PREFIX) else {
2664 continue;
2665 };
2666 if !name.is_empty() {
2667 if let Some(p) = self.get_project(name).await? {
2668 out.push(p);
2669 }
2670 }
2671 }
2672 if !out
2675 .iter()
2676 .any(|p| p.name == crate::project::DEFAULT_PROJECT)
2677 {
2678 out.push(Self::default_project_record());
2679 }
2680 out.sort_by(|a, b| a.name.cmp(&b.name));
2681 Ok(out)
2682 }
2683
2684 pub async fn delete_project(&self, name: &str) -> Result<bool, DeployError> {
2690 if name == crate::project::DEFAULT_PROJECT {
2691 return Err(DeployError::Conflict(
2692 "the `default` project cannot be deleted".to_string(),
2693 ));
2694 }
2695 let existed = self
2696 .kv
2697 .get(&crate::project::pointer_key(name))
2698 .await?
2699 .is_some();
2700 let prefix = crate::project::resource_prefix(name);
2706 let remaining = self.kv.list_prefix(&prefix).await?;
2707 if !remaining.is_empty() {
2708 return Err(DeployError::Conflict(format!(
2709 "project `{name}` still owns resources — delete these first, or \
2710 `project rm --force` to cascade: {}",
2711 Self::summarize_owned_resources(&prefix, &remaining)
2712 )));
2713 }
2714 self.kv
2715 .write_batch(vec![
2716 WriteOp::Delete(crate::project::pointer_key(name)),
2717 WriteOp::Delete(crate::project::history_key(name)),
2718 ])
2719 .await?;
2720 Ok(existed)
2721 }
2722
2723 fn summarize_owned_resources(prefix: &str, keys: &[String]) -> String {
2729 use std::collections::{BTreeMap, BTreeSet};
2730 let mut fam: BTreeMap<&str, (BTreeSet<&str>, usize)> = BTreeMap::new();
2732 for key in keys {
2733 let rest = key
2734 .strip_prefix(prefix)
2735 .unwrap_or(key)
2736 .trim_start_matches('/');
2737 let mut segs = rest.split('/');
2738 let Some(family) = segs.next().filter(|f| !f.is_empty()) else {
2739 continue;
2740 };
2741 let entry = fam.entry(family).or_default();
2742 entry.1 += 1;
2743 if let Some(name) = segs.next().filter(|n| !n.is_empty()) {
2744 entry.0.insert(name);
2745 }
2746 }
2747 fam.iter()
2748 .map(|(family, (names, count))| {
2749 if names.is_empty() {
2750 format!(
2751 "{family} ({count} key{})",
2752 if *count == 1 { "" } else { "s" }
2753 )
2754 } else {
2755 let shown: Vec<&str> = names.iter().take(10).copied().collect();
2756 let more = names.len() - shown.len();
2757 let suffix = if more > 0 {
2758 format!(", +{more} more")
2759 } else {
2760 String::new()
2761 };
2762 format!("{family}: [{}{suffix}]", shown.join(", "))
2763 }
2764 })
2765 .collect::<Vec<_>>()
2766 .join("; ")
2767 }
2768
2769 pub async fn enumerate_project_resources(
2784 &self,
2785 project: &str,
2786 ) -> Result<ProjectTeardownPlan, DeployError> {
2787 let pref = ProjectRef::new(project);
2788
2789 let mut site_set: std::collections::BTreeSet<String> =
2794 self.list_sites(pref).await?.into_iter().collect();
2795 let site_pref = keys::site_prefix(pref);
2796 for key in self.kv.list_prefix(&site_pref).await? {
2797 if let Some(site) = key.strip_prefix(&site_pref).filter(|s| !s.is_empty()) {
2798 site_set.insert(site.to_string());
2799 }
2800 }
2801 let sites: Vec<String> = site_set.into_iter().collect();
2802
2803 let functions: Vec<String> = self
2804 .list_stored_functions(pref)
2805 .await?
2806 .into_iter()
2807 .map(|f| f.name)
2808 .collect();
2809
2810 let mut compute = Vec::new();
2813 for w in self.list_compute_workloads(pref).await? {
2814 let volumes = match self.get_compute_spec(&w.active).await? {
2815 Some(spec) => spec.volumes.into_iter().map(|v| v.name).collect(),
2816 None => Vec::new(),
2817 };
2818 compute.push(ComputeTeardown {
2819 name: w.name,
2820 volumes,
2821 });
2822 }
2823
2824 let secret_prefix = keys::secret_prefix(pref);
2828 let mut secrets: Vec<String> = self
2829 .kv
2830 .list_prefix(&secret_prefix)
2831 .await?
2832 .into_iter()
2833 .filter_map(|k| k.strip_prefix(&secret_prefix).map(str::to_string))
2834 .filter(|n| !n.is_empty())
2835 .collect();
2836 secrets.sort();
2837
2838 let safelist = self
2842 .kv
2843 .list_prefix(&keys::graphql_safelist_prefix(pref))
2844 .await?
2845 .len();
2846
2847 let subgraph_prefix = keys::graphql_subgraph_prefix(pref);
2848 let mut subgraphs: Vec<String> = self
2849 .kv
2850 .list_prefix(&subgraph_prefix)
2851 .await?
2852 .into_iter()
2853 .filter_map(|k| k.strip_prefix(&subgraph_prefix).map(str::to_string))
2854 .filter(|n| !n.is_empty())
2855 .collect();
2856 subgraphs.sort();
2857
2858 let resource_prefix = crate::project::resource_prefix(project);
2861 let residual = self.kv.list_prefix(&resource_prefix).await?;
2862 let covered = ["site", "functions", "compute", "secret"];
2863 let mut other_families: std::collections::BTreeMap<String, usize> = Default::default();
2864 for key in &residual {
2865 let rest = key
2866 .strip_prefix(&resource_prefix)
2867 .unwrap_or(key)
2868 .trim_start_matches('/');
2869 let Some(family) = rest.split('/').next().filter(|f| !f.is_empty()) else {
2870 continue;
2871 };
2872 if covered.contains(&family) {
2873 continue;
2874 }
2875 *other_families.entry(family.to_string()).or_default() += 1;
2876 }
2877
2878 Ok(ProjectTeardownPlan {
2879 project: project.to_string(),
2880 sites,
2881 functions,
2882 compute,
2883 secrets,
2884 safelist,
2885 subgraphs,
2886 other_families,
2887 })
2888 }
2889
2890 pub async fn purge_project(&self, project: &str) -> Result<usize, DeployError> {
2908 if project == crate::project::DEFAULT_PROJECT {
2909 return Err(DeployError::Conflict(
2910 "the `default` project cannot be deleted".to_string(),
2911 ));
2912 }
2913 let _claim = self.domain_claim_lock.lock().await;
2917
2918 let pref = ProjectRef::new(project);
2919 let mut ops: Vec<WriteOp> = Vec::new();
2920
2921 for key in self
2923 .kv
2924 .list_prefix(&crate::project::resource_prefix(project))
2925 .await?
2926 {
2927 ops.push(WriteOp::Delete(key));
2928 }
2929 for key in self
2931 .kv
2932 .list_prefix(&keys::graphql_registry_prefix(pref))
2933 .await?
2934 {
2935 ops.push(WriteOp::Delete(key));
2936 }
2937 for key in self
2938 .kv
2939 .list_prefix(&keys::graphql_safelist_prefix(pref))
2940 .await?
2941 {
2942 ops.push(WriteOp::Delete(key));
2943 }
2944 ops.push(WriteOp::Delete(crate::project::pointer_key(project)));
2946 ops.push(WriteOp::Delete(crate::project::history_key(project)));
2947 for key in self.kv.list_prefix(crate::project::OWNER_PREFIX).await? {
2949 if let Some(bytes) = self.kv.get(&key).await? {
2950 if bytes == project.as_bytes() {
2951 ops.push(WriteOp::Delete(key));
2952 }
2953 }
2954 }
2955
2956 let purged = ops.len();
2957 if !ops.is_empty() {
2958 self.kv.write_batch(ops).await?;
2959 }
2960 self.bump_domain_epoch();
2962 Ok(purged)
2963 }
2964
2965 pub async fn activate(
2969 &self,
2970 project: ProjectRef<'_>,
2971 site: &str,
2972 id: &str,
2973 ) -> Result<(), DeployError> {
2974 let manifest = self
2975 .get_manifest(id)
2976 .await?
2977 .ok_or_else(|| DeployError::NotFound(format!("deployment {id}")))?;
2978 let missing = self.missing_blobs(&manifest).await?;
2979 if !missing.is_empty() {
2980 return Err(DeployError::Incomplete(missing));
2981 }
2982
2983 self.kv
2986 .put(&keys::current(project, site), id.as_bytes().to_vec())
2987 .await?;
2988
2989 let _ = self.record_history(project, site, id).await;
2993 Ok(())
2994 }
2995
2996 async fn record_history(
2999 &self,
3000 project: ProjectRef<'_>,
3001 site: &str,
3002 id: &str,
3003 ) -> Result<(), DeployError> {
3004 let mut history = self.history(project, site).await.unwrap_or_default();
3005 history.retain(|entry| entry.id != id);
3006 history.insert(
3007 0,
3008 HistoryEntry {
3009 id: id.to_string(),
3010 at: now_unix(),
3011 meta: None,
3012 },
3013 );
3014 history.truncate(MAX_HISTORY);
3015 self.kv
3016 .put(&keys::history(project, site), serde_json::to_vec(&history)?)
3017 .await?;
3018 Ok(())
3019 }
3020
3021 pub async fn history(
3023 &self,
3024 project: ProjectRef<'_>,
3025 site: &str,
3026 ) -> Result<Vec<HistoryEntry>, DeployError> {
3027 match self.kv.get(&keys::history(project, site)).await? {
3028 Some(bytes) => Ok(serde_json::from_slice(&bytes)?),
3029 None => Ok(Vec::new()),
3030 }
3031 }
3032
3033 pub async fn deployments(
3036 &self,
3037 project: ProjectRef<'_>,
3038 site: &str,
3039 ) -> Result<DeploymentList, DeployError> {
3040 let mut deployments = self.history(project, site).await?;
3041 for entry in &mut deployments {
3042 entry.meta = self.get_meta(&entry.id).await?;
3043 }
3044 Ok(DeploymentList {
3045 current: self.current_id(project, site).await?,
3046 deployments,
3047 })
3048 }
3049
3050 async fn live_deployment_ids(&self, opts: &GcOptions) -> Result<BTreeSet<String>, DeployError> {
3061 let mut ids = BTreeSet::new();
3062 let now = now_unix();
3063 for project in self.discover_projects().await? {
3064 let pref = ProjectRef::new(&project);
3065 for key in self.kv.list_prefix(&keys::history_prefix(pref)).await? {
3066 if let Some(bytes) = self.kv.get(&key).await? {
3067 if let Ok(history) = serde_json::from_slice::<Vec<HistoryEntry>>(&bytes) {
3068 for (idx, entry) in history.iter().enumerate() {
3069 let within_count = opts.keep_last.is_none_or(|n| idx < n);
3070 let within_age = opts
3071 .keep_age_secs
3072 .is_some_and(|age| now.saturating_sub(entry.at) <= age);
3073 if within_count || within_age {
3074 ids.insert(entry.id.clone());
3075 }
3076 }
3077 }
3078 }
3079 }
3080 for prefix in [keys::current_prefix(pref), keys::alias_project_prefix(pref)] {
3082 for key in self.kv.list_prefix(&prefix).await? {
3083 if let Some(bytes) = self.kv.get(&key).await? {
3084 ids.insert(String::from_utf8_lossy(&bytes).into_owned());
3085 }
3086 }
3087 }
3088 }
3089 Ok(ids)
3090 }
3091
3092 async fn within_grace(
3097 &self,
3098 id: &str,
3099 now: u64,
3100 opts: &GcOptions,
3101 ) -> Result<bool, DeployError> {
3102 if opts.grace_secs == 0 {
3103 return Ok(false);
3104 }
3105 match self.get_meta(id).await? {
3106 Some(meta) => Ok(now.saturating_sub(meta.created_at) < opts.grace_secs),
3107 None => Ok(false),
3108 }
3109 }
3110
3111 pub async fn collect_garbage(&self, prune: bool) -> Result<GcReport, DeployError> {
3118 self.collect_garbage_with(prune, GcOptions::default()).await
3119 }
3120
3121 pub async fn collect_garbage_with(
3132 &self,
3133 prune: bool,
3134 opts: GcOptions,
3135 ) -> Result<GcReport, DeployError> {
3136 let live_ids = self.live_deployment_ids(&opts).await?;
3137 let now = now_unix();
3138
3139 let manifest_keys = self.kv.list_prefix("manifests/").await?;
3140 let manifests_total = manifest_keys.len();
3141 let mut referenced: BTreeSet<String> = BTreeSet::new();
3142 let mut orphan_manifests: Vec<String> = Vec::new();
3143 for key in &manifest_keys {
3144 let id = key.strip_prefix("manifests/").unwrap_or(key);
3145 let protected = live_ids.contains(id) || self.within_grace(id, now, &opts).await?;
3146 if protected {
3147 if let Some(bytes) = self.kv.get(key).await? {
3148 if let Ok(manifest) = Manifest::from_bytes(&bytes) {
3149 referenced.extend(manifest.blob_hashes());
3150 }
3151 }
3152 } else {
3153 orphan_manifests.push(key.clone());
3154 }
3155 }
3156
3157 let blobs = self.storage.list("").await?;
3158 let blobs_total = blobs.len();
3159 let mut blobs_removed = 0;
3160 let mut bytes_reclaimed = 0;
3161 for meta in &blobs {
3162 if !is_blob_key(&meta.key) {
3163 continue;
3164 }
3165 let hash = meta.key.rsplit('/').next().unwrap_or(&meta.key);
3166 if !referenced.contains(hash) {
3167 blobs_removed += 1;
3168 bytes_reclaimed += meta.size.unwrap_or(0);
3169 if prune {
3170 self.storage.delete(&meta.key).await?;
3171 }
3172 }
3173 }
3174
3175 if prune {
3181 let mut referenced_configs: BTreeSet<String> = BTreeSet::new();
3185 for project in self.discover_projects().await? {
3186 let site_prefix = keys::site_prefix(ProjectRef::new(&project));
3187 for pointer in self.kv.list_prefix(&site_prefix).await? {
3188 if let Some(bytes) = self.kv.get(&pointer).await? {
3189 referenced_configs.insert(String::from_utf8_lossy(&bytes).into_owned());
3190 }
3191 }
3192 }
3193 for key in self.kv.list_prefix("siteconfig/").await? {
3194 let hash = key.strip_prefix("siteconfig/").unwrap_or(&key);
3195 if !referenced_configs.contains(hash) {
3196 self.kv.delete(&key).await?;
3197 }
3198 }
3199 }
3200
3201 let manifests_removed = orphan_manifests.len();
3202 if prune {
3203 for key in &orphan_manifests {
3204 self.kv.delete(key).await?;
3205 if let Some(id) = key.strip_prefix("manifests/") {
3207 let _ = self.kv.delete(&keys::meta(id)).await;
3208 }
3209 }
3210 }
3211
3212 Ok(GcReport {
3213 manifests_total,
3214 manifests_removed,
3215 blobs_total,
3216 blobs_removed,
3217 bytes_reclaimed,
3218 })
3219 }
3220
3221 pub async fn scrub_blobs(&self) -> Result<ScrubReport, DeployError> {
3227 let blobs = self.storage.list("").await?;
3228 let mut report = ScrubReport::default();
3229 for meta in &blobs {
3230 if !is_blob_key(&meta.key) {
3231 continue;
3232 }
3233 report.checked += 1;
3234 let expected = meta
3235 .key
3236 .rsplit('/')
3237 .next()
3238 .unwrap_or(meta.key.as_str())
3239 .to_string();
3240 match self.hash_stored_object(&meta.key).await {
3241 Ok(actual) if actual == expected => {}
3242 Ok(actual) => report.mismatched.push(BlobMismatch {
3243 key: meta.key.clone(),
3244 expected,
3245 actual,
3246 }),
3247 Err(err) => report.errors.push(BlobReadError {
3248 key: meta.key.clone(),
3249 error: err.to_string(),
3250 }),
3251 }
3252 }
3253 Ok(report)
3254 }
3255
3256 pub fn invalidate_cache_keys(&self, keys: &[String]) {
3262 self.kv.invalidate_keys(keys);
3263 }
3264
3265 pub fn invalidate_cache(&self) {
3268 self.kv.invalidate_cache();
3269 }
3270
3271 pub async fn cert_status(&self) -> Result<Vec<crate::cert::CertStatus>, DeployError> {
3276 let mut out = Vec::new();
3277 for key in self.kv.list_prefix("cert/").await? {
3278 let domain = key.strip_prefix("cert/").unwrap_or(&key).to_string();
3279 if let Some(bytes) = self.kv.get(&key).await? {
3280 if let Ok(cert) = serde_json::from_slice::<crate::cert::StoredCert>(&bytes) {
3281 out.push(crate::cert::CertStatus {
3282 domain,
3283 not_after_unix: cert.not_after_unix,
3284 });
3285 }
3286 }
3287 }
3288 out.sort_by(|a, b| a.domain.cmp(&b.domain));
3289 Ok(out)
3290 }
3291
3292 async fn hash_stored_object(&self, key: &str) -> Result<String, DeployError> {
3294 let mut body = self.storage.get(key).await?.body;
3295 let mut hasher = Sha256::new();
3296 while let Some(chunk) = body.next().await {
3297 hasher.update(&chunk?);
3298 }
3299 Ok(hex::encode(hasher.finalize()))
3300 }
3301
3302 pub async fn list_sites(&self, project: ProjectRef<'_>) -> Result<Vec<String>, DeployError> {
3306 let prefix = keys::current_prefix(project);
3307 let keys = self.kv.list_prefix(&prefix).await?;
3308 Ok(keys
3309 .into_iter()
3310 .filter_map(|k| k.strip_prefix(&prefix).map(str::to_string))
3311 .collect())
3312 }
3313
3314 pub async fn list_sites_all(&self) -> Result<Vec<(String, String)>, DeployError> {
3317 let mut out = Vec::new();
3318 for project in self.discover_projects().await? {
3319 for site in self.list_sites(ProjectRef::new(&project)).await? {
3320 out.push((project.clone(), site));
3321 }
3322 }
3323 Ok(out)
3324 }
3325
3326 pub async fn delete_site(
3333 &self,
3334 project: ProjectRef<'_>,
3335 site: &str,
3336 ) -> Result<(), DeployError> {
3337 use crate::kv::WriteOp;
3338 let _claim = self.domain_claim_lock.lock().await;
3341 let mut batch = vec![
3342 WriteOp::Delete(keys::site_pointer(project, site)),
3343 WriteOp::Delete(keys::current(project, site)),
3344 WriteOp::Delete(keys::history(project, site)),
3345 ];
3346 if let Some(config) = self.get_site_config(project, site).await? {
3347 for host in config.domains.exact_hosts() {
3348 batch.push(WriteOp::Delete(keys::domain(host)));
3349 }
3350 for wildcard in &config.domains.wildcards {
3351 if let Some(suffix) = wildcard.strip_prefix("*.") {
3352 batch.push(WriteOp::Delete(keys::wildcard(suffix)));
3353 }
3354 }
3355 }
3356 for key in self
3357 .kv
3358 .list_prefix(&keys::alias_prefix(project, site))
3359 .await?
3360 {
3361 batch.push(WriteOp::Delete(key));
3362 }
3363 for key in self
3364 .kv
3365 .list_prefix(&keys::domain_verification_prefix(project, site))
3366 .await?
3367 {
3368 batch.push(WriteOp::Delete(key));
3369 }
3370 self.kv.write_batch(batch).await?;
3371 self.bump_domain_epoch();
3373 Ok(())
3374 }
3375
3376 pub async fn current_id(
3378 &self,
3379 project: ProjectRef<'_>,
3380 site: &str,
3381 ) -> Result<Option<String>, DeployError> {
3382 match self.kv.get(&keys::current(project, site)).await? {
3383 Some(bytes) => Ok(Some(String::from_utf8_lossy(&bytes).into_owned())),
3384 None => Ok(None),
3385 }
3386 }
3387
3388 pub async fn current_manifest(
3390 &self,
3391 project: ProjectRef<'_>,
3392 site: &str,
3393 ) -> Result<Option<Manifest>, DeployError> {
3394 match self.current_id(project, site).await? {
3395 Some(id) => self.get_manifest(&id).await,
3396 None => Ok(None),
3397 }
3398 }
3399
3400 pub async fn resolve(
3405 &self,
3406 project: ProjectRef<'_>,
3407 site: &str,
3408 path: &str,
3409 ) -> Result<Option<FileEntry>, DeployError> {
3410 let Some(manifest) = self.current_manifest(project, site).await? else {
3411 return Ok(None);
3412 };
3413 Ok(lookup(&manifest, path))
3414 }
3415}
3416
3417fn lookup(manifest: &Manifest, path: &str) -> Option<FileEntry> {
3419 let trimmed = path.trim_start_matches('/');
3420 if let Some(entry) = manifest.files.get(trimmed) {
3421 return Some(entry.clone());
3422 }
3423 let index = if trimmed.is_empty() {
3424 "index.html".to_string()
3425 } else {
3426 format!("{}/index.html", trimmed.trim_end_matches('/'))
3427 };
3428 manifest.files.get(&index).cloned()
3429}
3430
3431#[cfg(test)]
3432mod tests {
3433 use super::*;
3434 use crate::config::DeployConfig;
3435 use crate::ObjectMeta;
3436
3437 fn entry(hash: &str) -> FileEntry {
3438 FileEntry {
3439 hash: hash.to_string(),
3440 size: 0,
3441 content_type: None,
3442 variants: BTreeMap::new(),
3443 }
3444 }
3445
3446 #[test]
3447 fn manifest_id_is_deterministic() {
3448 let mut a = Manifest::default();
3449 a.files.insert("index.html".into(), entry("aa"));
3450 a.files.insert("style.css".into(), entry("bb"));
3451
3452 let mut b = Manifest::default();
3453 b.files.insert("style.css".into(), entry("bb"));
3455 b.files.insert("index.html".into(), entry("aa"));
3456
3457 assert_eq!(a.id().unwrap(), b.id().unwrap());
3458 }
3459
3460 #[test]
3461 fn manifest_carries_schema_version_and_reads_legacy() {
3462 let manifest = Manifest::default();
3464 assert_eq!(manifest.version, crate::SCHEMA_VERSION);
3465 assert!(manifest.to_bytes().unwrap().starts_with(b"{\"version\":1"));
3466
3467 let legacy = br#"{"files":{},"config":{}}"#;
3469 assert_eq!(Manifest::from_bytes(legacy).unwrap().version, 1);
3470 }
3471
3472 #[test]
3473 fn directory_index_fallback() {
3474 let mut m = Manifest::default();
3475 m.files.insert("index.html".into(), entry("root"));
3476 m.files.insert("blog/index.html".into(), entry("blog"));
3477
3478 assert_eq!(lookup(&m, "").unwrap().hash, "root");
3479 assert_eq!(lookup(&m, "/").unwrap().hash, "root");
3480 assert_eq!(lookup(&m, "blog").unwrap().hash, "blog");
3481 assert_eq!(lookup(&m, "blog/").unwrap().hash, "blog");
3482 assert!(lookup(&m, "missing.html").is_none());
3483 }
3484
3485 struct NullStorage;
3488
3489 #[async_trait::async_trait]
3490 impl Storage for NullStorage {
3491 async fn get(&self, _: &str) -> Result<GetObject, StorageError> {
3492 Err(StorageError::NotFound(String::new()))
3493 }
3494 async fn get_range(
3495 &self,
3496 _: &str,
3497 _: u64,
3498 _: Option<u64>,
3499 ) -> Result<GetObject, StorageError> {
3500 Err(StorageError::NotFound(String::new()))
3501 }
3502 async fn put(
3503 &self,
3504 _: &str,
3505 _: ByteStream,
3506 _: PutMeta,
3507 ) -> Result<ObjectMeta, StorageError> {
3508 Err(StorageError::unsupported("null"))
3509 }
3510 async fn head(&self, _: &str) -> Result<ObjectMeta, StorageError> {
3511 Err(StorageError::NotFound(String::new()))
3512 }
3513 async fn delete(&self, _: &str) -> Result<(), StorageError> {
3514 Ok(())
3515 }
3516 async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
3517 Ok(Vec::new())
3518 }
3519 }
3520
3521 #[derive(Default)]
3523 struct MemStorage {
3524 objects: Mutex<std::collections::HashMap<String, Vec<u8>>>,
3525 }
3526
3527 #[async_trait::async_trait]
3528 impl Storage for MemStorage {
3529 async fn get(&self, key: &str) -> Result<GetObject, StorageError> {
3530 let bytes = self
3531 .objects
3532 .lock()
3533 .unwrap()
3534 .get(key)
3535 .cloned()
3536 .ok_or_else(|| StorageError::NotFound(key.to_string()))?;
3537 let size = bytes.len() as u64;
3538 let body: ByteStream =
3539 futures::stream::once(async move { Ok(bytes::Bytes::from(bytes)) }).boxed();
3540 Ok(GetObject {
3541 meta: ObjectMeta {
3542 key: key.to_string(),
3543 size: Some(size),
3544 ..Default::default()
3545 },
3546 body,
3547 })
3548 }
3549 async fn get_range(
3550 &self,
3551 key: &str,
3552 _: u64,
3553 _: Option<u64>,
3554 ) -> Result<GetObject, StorageError> {
3555 self.get(key).await
3556 }
3557 async fn put(
3558 &self,
3559 key: &str,
3560 mut body: ByteStream,
3561 _: PutMeta,
3562 ) -> Result<ObjectMeta, StorageError> {
3563 let mut buf = Vec::new();
3564 while let Some(chunk) = body.next().await {
3565 buf.extend_from_slice(&chunk?);
3566 }
3567 let size = buf.len() as u64;
3568 self.objects.lock().unwrap().insert(key.to_string(), buf);
3569 Ok(ObjectMeta {
3570 key: key.to_string(),
3571 size: Some(size),
3572 ..Default::default()
3573 })
3574 }
3575 async fn head(&self, key: &str) -> Result<ObjectMeta, StorageError> {
3576 let map = self.objects.lock().unwrap();
3577 let bytes = map
3578 .get(key)
3579 .ok_or_else(|| StorageError::NotFound(key.to_string()))?;
3580 Ok(ObjectMeta {
3581 key: key.to_string(),
3582 size: Some(bytes.len() as u64),
3583 ..Default::default()
3584 })
3585 }
3586 async fn delete(&self, key: &str) -> Result<(), StorageError> {
3587 self.objects.lock().unwrap().remove(key);
3588 Ok(())
3589 }
3590 async fn list(&self, prefix: &str) -> Result<Vec<ObjectMeta>, StorageError> {
3591 Ok(self
3592 .objects
3593 .lock()
3594 .unwrap()
3595 .keys()
3596 .filter(|k| k.starts_with(prefix))
3597 .map(|k| ObjectMeta {
3598 key: k.clone(),
3599 ..Default::default()
3600 })
3601 .collect())
3602 }
3603 }
3604
3605 fn once_bytes(b: &'static [u8]) -> ByteStream {
3607 futures::stream::once(async move { Ok(bytes::Bytes::from_static(b)) }).boxed()
3608 }
3609
3610 fn manifest_with(files: &[(&str, &str)]) -> Manifest {
3612 let mut m = Manifest::default();
3613 for (path, hash) in files {
3614 m.files.insert(
3615 (*path).to_string(),
3616 FileEntry {
3617 hash: (*hash).to_string(),
3618 size: 1,
3619 content_type: None,
3620 variants: Default::default(),
3621 },
3622 );
3623 }
3624 m
3625 }
3626
3627 #[tokio::test]
3628 async fn function_storage_versioning_alias_rollback() {
3629 use crate::function::{Function, FunctionConfig, Lifecycle, Owner};
3630 use crate::kv::MemoryKv;
3631
3632 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
3633 assert!(store
3634 .list_stored_functions(ProjectRef::DEFAULT)
3635 .await
3636 .unwrap()
3637 .is_empty());
3638
3639 let mut f = Function::new(
3640 "resize",
3641 Owner::Project("acme".into()),
3642 "hashA",
3643 FunctionConfig::default(),
3644 Lifecycle::Independent,
3645 1,
3646 );
3647 store.put_function(ProjectRef::DEFAULT, &f).await.unwrap();
3648 assert_eq!(
3649 store
3650 .get_function(ProjectRef::DEFAULT, "resize")
3651 .await
3652 .unwrap()
3653 .unwrap()
3654 .active,
3655 "hashA"
3656 );
3657 assert_eq!(
3658 store
3659 .list_stored_functions(ProjectRef::DEFAULT)
3660 .await
3661 .unwrap()
3662 .len(),
3663 1
3664 );
3665
3666 f.upsert_version("hashB", Lifecycle::Independent, 2);
3668 f.set_alias("prod", "hashA").unwrap();
3669 store.put_function(ProjectRef::DEFAULT, &f).await.unwrap();
3670 let got = store
3671 .get_function(ProjectRef::DEFAULT, "resize")
3672 .await
3673 .unwrap()
3674 .unwrap();
3675 assert_eq!(got.active, "hashB");
3676 assert_eq!(got.aliases.get("prod").map(String::as_str), Some("hashA"));
3677
3678 assert!(store
3680 .delete_function(ProjectRef::DEFAULT, "resize")
3681 .await
3682 .unwrap());
3683 assert!(store
3684 .get_function(ProjectRef::DEFAULT, "resize")
3685 .await
3686 .unwrap()
3687 .is_none());
3688 assert!(!store
3689 .delete_function(ProjectRef::DEFAULT, "resize")
3690 .await
3691 .unwrap());
3692 }
3693
3694 #[tokio::test]
3695 async fn delete_function_sweeps_the_whole_subtree_only_for_that_function() {
3696 use crate::function::keys as fk;
3697 use crate::kv::{KvStore, MemoryKv};
3698
3699 let kv = Arc::new(MemoryKv::new());
3700 let store = DeployStore::new(Arc::new(NullStorage), kv.clone());
3701 let proj = ProjectRef::new("acme");
3702 for k in [
3705 fk::meta("acme", "greeter"),
3706 fk::version("acme", "greeter", "v1"),
3707 fk::alias("acme", "greeter", "prod"),
3708 fk::trigger("acme", "greeter", "t1"),
3709 fk::invocation("acme", "greeter", "i1"),
3710 fk::metering("acme", "greeter"),
3711 ] {
3712 kv.put(&k, vec![1]).await.unwrap();
3713 }
3714 kv.put(&fk::version("acme", "greeterx", "z1"), vec![1])
3716 .await
3717 .unwrap();
3718
3719 assert!(store.delete_function(proj, "greeter").await.unwrap());
3720
3721 assert!(kv
3723 .get(&fk::meta("acme", "greeter"))
3724 .await
3725 .unwrap()
3726 .is_none());
3727 assert!(kv
3728 .list_prefix(&format!("{}/", fk::meta("acme", "greeter")))
3729 .await
3730 .unwrap()
3731 .is_empty());
3732 assert!(kv
3733 .get(&fk::metering("acme", "greeter"))
3734 .await
3735 .unwrap()
3736 .is_none());
3737 assert!(kv
3739 .get(&fk::version("acme", "greeterx", "z1"))
3740 .await
3741 .unwrap()
3742 .is_some());
3743 }
3744
3745 #[tokio::test]
3746 async fn delete_project_refusal_names_what_remains() {
3747 use crate::function::keys as fk;
3748 use crate::kv::{KvStore, MemoryKv};
3749
3750 let kv = Arc::new(MemoryKv::new());
3751 let store = DeployStore::new(Arc::new(NullStorage), kv.clone());
3752 kv.put(&fk::meta("acme", "greeter"), vec![1]).await.unwrap();
3754 kv.put("project/acme/graphql/safelist/op1", vec![1])
3755 .await
3756 .unwrap();
3757 kv.put("project/acme/graphql/safelist/op2", vec![1])
3758 .await
3759 .unwrap();
3760
3761 let msg = store.delete_project("acme").await.unwrap_err().to_string();
3762 assert!(msg.contains("functions: [greeter]"), "{msg}");
3764 assert!(msg.contains("graphql"), "{msg}");
3765 assert!(msg.contains("--force"), "{msg}");
3766 }
3767
3768 #[tokio::test]
3769 async fn function_invocation_and_idempotency_storage() {
3770 use crate::function::{
3771 Function, FunctionConfig, Invocation, InvocationResult, InvocationStatus, InvokeMode,
3772 Lifecycle, Owner,
3773 };
3774 use crate::kv::MemoryKv;
3775
3776 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
3777 let f = Function::new(
3779 "greeter",
3780 Owner::Project("acme".into()),
3781 "hashA",
3782 FunctionConfig::default(),
3783 Lifecycle::Independent,
3784 1,
3785 );
3786 store.put_function(ProjectRef::DEFAULT, &f).await.unwrap();
3787
3788 let mut inv = Invocation {
3789 id: "inv-1".into(),
3790 function: "greeter".into(),
3791 version: "hashA".into(),
3792 mode: InvokeMode::Async,
3793 status: InvocationStatus::Queued,
3794 idempotency_key: Some("key-1".into()),
3795 attempts: 0,
3796 lease_expires: None,
3797 request_b64: None,
3798 request_content_type: None,
3799 result: None,
3800 created: 10,
3801 updated: 10,
3802 };
3803 store
3804 .put_invocation(ProjectRef::DEFAULT, &inv)
3805 .await
3806 .unwrap();
3807 store
3808 .put_idempotency(ProjectRef::DEFAULT, "greeter", "key-1", "inv-1")
3809 .await
3810 .unwrap();
3811
3812 assert_eq!(
3814 store
3815 .get_invocation(ProjectRef::DEFAULT, "greeter", "inv-1")
3816 .await
3817 .unwrap()
3818 .unwrap()
3819 .status,
3820 InvocationStatus::Queued
3821 );
3822 assert_eq!(
3823 store
3824 .get_idempotency(ProjectRef::DEFAULT, "greeter", "key-1")
3825 .await
3826 .unwrap(),
3827 Some("inv-1".to_string())
3828 );
3829 assert_eq!(
3830 store
3831 .list_invocations(ProjectRef::DEFAULT, "greeter")
3832 .await
3833 .unwrap()
3834 .len(),
3835 1
3836 );
3837 assert_eq!(
3839 store
3840 .list_stored_functions(ProjectRef::DEFAULT)
3841 .await
3842 .unwrap()
3843 .len(),
3844 1
3845 );
3846
3847 inv.status = InvocationStatus::Succeeded;
3849 inv.attempts = 1;
3850 inv.result = Some(InvocationResult {
3851 status: 200,
3852 content_type: Some("text/plain".into()),
3853 body_b64: "aGVsbG8=".into(),
3854 });
3855 inv.updated = 20;
3856 store
3857 .put_invocation(ProjectRef::DEFAULT, &inv)
3858 .await
3859 .unwrap();
3860 let got = store
3861 .get_invocation(ProjectRef::DEFAULT, "greeter", "inv-1")
3862 .await
3863 .unwrap()
3864 .unwrap();
3865 assert!(got.is_terminal());
3866 assert_eq!(got.result.unwrap().status, 200);
3867
3868 assert!(store
3870 .get_idempotency(ProjectRef::DEFAULT, "greeter", "absent")
3871 .await
3872 .unwrap()
3873 .is_none());
3874 }
3875
3876 #[tokio::test]
3877 async fn notification_ledger_provisions_and_retracts_through_the_store() {
3878 use crate::blob_notify::{ManagedResource, ProvisionTier};
3879 use crate::blob_provision::{ensure_watch, retract_watch, ProvisionError, WatchProvider};
3880 use crate::kv::MemoryKv;
3881
3882 struct StoreMock;
3884 #[async_trait::async_trait]
3885 impl WatchProvider for StoreMock {
3886 fn name(&self) -> &str {
3887 "mock"
3888 }
3889 fn recipe(&self, _prefix: &str) -> String {
3890 String::new()
3891 }
3892 async fn provision(
3893 &self,
3894 _prefix: &str,
3895 ) -> Result<Vec<ManagedResource>, ProvisionError> {
3896 Ok(vec![ManagedResource::new("queue", "q-1")])
3897 }
3898 async fn verify(&self, _prefix: &str) -> Result<bool, ProvisionError> {
3899 Ok(true)
3900 }
3901 async fn retract(&self, _res: &[ManagedResource]) -> Result<(), ProvisionError> {
3902 Ok(())
3903 }
3904 }
3905
3906 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
3907 assert!(store
3908 .get_managed_notification(ProjectRef::DEFAULT, "ingest", "uploads/")
3909 .await
3910 .unwrap()
3911 .is_none());
3912
3913 let provider = StoreMock;
3915 let out = ensure_watch(
3916 &provider,
3917 ProvisionTier::Provision,
3918 "ingest",
3919 "uploads/",
3920 &store,
3921 7,
3922 )
3923 .await
3924 .unwrap();
3925 assert!(matches!(
3926 out,
3927 crate::blob_provision::ProvisionOutcome::Ready
3928 ));
3929 let record = store
3930 .get_managed_notification(ProjectRef::DEFAULT, "ingest", "uploads/")
3931 .await
3932 .unwrap()
3933 .expect("the pipeline is recorded in the store ledger");
3934 assert_eq!(record.provider, "mock");
3935 assert_eq!(
3936 store
3937 .list_managed_notifications(ProjectRef::DEFAULT, "ingest")
3938 .await
3939 .unwrap()
3940 .len(),
3941 1
3942 );
3943
3944 retract_watch(&provider, &record, &store).await.unwrap();
3946 assert!(store
3947 .get_managed_notification(ProjectRef::DEFAULT, "ingest", "uploads/")
3948 .await
3949 .unwrap()
3950 .is_none());
3951 }
3952
3953 #[tokio::test]
3954 async fn workflow_definition_and_run_storage() {
3955 use crate::kv::MemoryKv;
3956 use crate::workflow::{Step, Workflow, WorkflowRun};
3957
3958 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
3959 assert!(store
3960 .get_workflow(ProjectRef::DEFAULT, "etl")
3961 .await
3962 .unwrap()
3963 .is_none());
3964 assert!(store
3965 .list_workflows(ProjectRef::DEFAULT)
3966 .await
3967 .unwrap()
3968 .is_empty());
3969
3970 let wf = Workflow {
3971 name: "etl".into(),
3972 steps: vec![
3973 Step {
3974 id: "a".into(),
3975 function: "extract".into(),
3976 depends_on: vec![],
3977 retry: Default::default(),
3978 compensate: None,
3979 },
3980 Step {
3981 id: "b".into(),
3982 function: "load".into(),
3983 depends_on: vec!["a".into()],
3984 retry: Default::default(),
3985 compensate: None,
3986 },
3987 ],
3988 };
3989 store.put_workflow(ProjectRef::DEFAULT, &wf).await.unwrap();
3990 assert_eq!(
3991 store
3992 .get_workflow(ProjectRef::DEFAULT, "etl")
3993 .await
3994 .unwrap()
3995 .unwrap(),
3996 wf
3997 );
3998 assert_eq!(
3999 store
4000 .list_workflows(ProjectRef::DEFAULT)
4001 .await
4002 .unwrap()
4003 .len(),
4004 1
4005 );
4006
4007 let run = WorkflowRun::start(&wf, "r1", None, 5);
4009 store
4010 .put_workflow_run(ProjectRef::DEFAULT, &run)
4011 .await
4012 .unwrap();
4013 assert_eq!(
4014 store
4015 .get_workflow_run(ProjectRef::DEFAULT, "etl", "r1")
4016 .await
4017 .unwrap()
4018 .unwrap(),
4019 run
4020 );
4021 assert_eq!(
4022 store
4023 .list_workflow_runs(ProjectRef::DEFAULT, "etl")
4024 .await
4025 .unwrap()
4026 .len(),
4027 1
4028 );
4029 assert_eq!(
4031 store
4032 .list_workflows(ProjectRef::DEFAULT)
4033 .await
4034 .unwrap()
4035 .len(),
4036 1
4037 );
4038
4039 assert!(store
4041 .delete_workflow(ProjectRef::DEFAULT, "etl")
4042 .await
4043 .unwrap());
4044 assert!(store
4045 .get_workflow(ProjectRef::DEFAULT, "etl")
4046 .await
4047 .unwrap()
4048 .is_none());
4049 assert!(!store
4050 .delete_workflow(ProjectRef::DEFAULT, "etl")
4051 .await
4052 .unwrap());
4053 }
4054
4055 #[tokio::test]
4056 async fn function_trigger_storage_round_trips() {
4057 use crate::function::{
4058 Function, FunctionConfig, FunctionTrigger, Lifecycle, Owner, TriggerKind,
4059 };
4060 use crate::kv::MemoryKv;
4061
4062 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
4063 let f = Function::new(
4064 "worker",
4065 Owner::Project("acme".into()),
4066 "hashA",
4067 FunctionConfig::default(),
4068 Lifecycle::Independent,
4069 1,
4070 );
4071 store.put_function(ProjectRef::DEFAULT, &f).await.unwrap();
4072 assert!(store
4073 .list_triggers(ProjectRef::DEFAULT, "worker")
4074 .await
4075 .unwrap()
4076 .is_empty());
4077
4078 let cron = FunctionTrigger {
4079 id: "tick".into(),
4080 kind: TriggerKind::Cron {
4081 schedule: "* * * * *".into(),
4082 overlap: Default::default(),
4083 },
4084 last_fired_minute: None,
4085 };
4086 store
4087 .put_trigger(ProjectRef::DEFAULT, "worker", &cron)
4088 .await
4089 .unwrap();
4090 let queue = FunctionTrigger {
4091 id: "jobs".into(),
4092 kind: TriggerKind::Queue {
4093 topic: "jobs".into(),
4094 group: String::new(),
4095 start: Default::default(),
4096 },
4097 last_fired_minute: None,
4098 };
4099 store
4100 .put_trigger(ProjectRef::DEFAULT, "worker", &queue)
4101 .await
4102 .unwrap();
4103
4104 assert_eq!(
4105 store
4106 .list_triggers(ProjectRef::DEFAULT, "worker")
4107 .await
4108 .unwrap()
4109 .len(),
4110 2
4111 );
4112 assert_eq!(
4113 store
4114 .get_trigger(ProjectRef::DEFAULT, "worker", "tick")
4115 .await
4116 .unwrap()
4117 .unwrap(),
4118 cron
4119 );
4120 assert_eq!(
4122 store
4123 .list_stored_functions(ProjectRef::DEFAULT)
4124 .await
4125 .unwrap()
4126 .len(),
4127 1
4128 );
4129
4130 let mut fired = cron.clone();
4132 fired.last_fired_minute = Some(42);
4133 store
4134 .put_trigger(ProjectRef::DEFAULT, "worker", &fired)
4135 .await
4136 .unwrap();
4137 assert_eq!(
4138 store
4139 .get_trigger(ProjectRef::DEFAULT, "worker", "tick")
4140 .await
4141 .unwrap()
4142 .unwrap()
4143 .last_fired_minute,
4144 Some(42)
4145 );
4146
4147 assert!(store
4149 .delete_trigger(ProjectRef::DEFAULT, "worker", "tick")
4150 .await
4151 .unwrap());
4152 assert!(!store
4153 .delete_trigger(ProjectRef::DEFAULT, "worker", "tick")
4154 .await
4155 .unwrap());
4156 assert_eq!(
4157 store
4158 .list_triggers(ProjectRef::DEFAULT, "worker")
4159 .await
4160 .unwrap()
4161 .len(),
4162 1
4163 );
4164 }
4165
4166 #[tokio::test]
4167 async fn function_metering_storage_is_tenant_isolated() {
4168 use crate::function::{Metering, MeteringSample};
4169 use crate::kv::MemoryKv;
4170
4171 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
4172 assert!(store
4173 .get_metering(ProjectRef::DEFAULT, "a")
4174 .await
4175 .unwrap()
4176 .is_none());
4177
4178 let mut ma = Metering::new("a");
4179 ma.record(
4180 &MeteringSample {
4181 success: true,
4182 duration_ms: 4,
4183 bytes_in: 1,
4184 bytes_out: 2,
4185 },
4186 10,
4187 );
4188 store.put_metering(ProjectRef::DEFAULT, &ma).await.unwrap();
4189
4190 let mut mb = Metering::new("b");
4191 mb.record(
4192 &MeteringSample {
4193 success: false,
4194 duration_ms: 9,
4195 bytes_in: 0,
4196 bytes_out: 0,
4197 },
4198 11,
4199 );
4200 store.put_metering(ProjectRef::DEFAULT, &mb).await.unwrap();
4201
4202 assert_eq!(
4204 store
4205 .get_metering(ProjectRef::DEFAULT, "a")
4206 .await
4207 .unwrap()
4208 .unwrap()
4209 .successes,
4210 1
4211 );
4212 assert_eq!(
4213 store
4214 .get_metering(ProjectRef::DEFAULT, "b")
4215 .await
4216 .unwrap()
4217 .unwrap()
4218 .failures,
4219 1
4220 );
4221 let all = store.list_metering(ProjectRef::DEFAULT).await.unwrap();
4222 assert_eq!(all.len(), 2);
4223 }
4224
4225 #[tokio::test]
4226 async fn managed_dns_ledger_round_trip_and_retract() {
4227 use crate::dns_managed::{ManagedDns, ManagedRecord};
4228 use crate::kv::MemoryKv;
4229
4230 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
4231 assert!(store
4232 .get_managed_dns(
4233 ProjectRef::DEFAULT,
4234 &SiteName::new("blog"),
4235 "www.example.com"
4236 )
4237 .await
4238 .unwrap()
4239 .is_none());
4240
4241 let ledger = ManagedDns::new(
4242 "www.example.com",
4243 "cloudflare",
4244 vec![ManagedRecord {
4245 kind: "A".into(),
4246 name: "www.example.com".into(),
4247 value: "203.0.113.7".into(),
4248 ttl: 300,
4249 }],
4250 10,
4251 );
4252 store
4253 .set_managed_dns(ProjectRef::DEFAULT, &SiteName::new("blog"), &ledger)
4254 .await
4255 .unwrap();
4256 assert_eq!(
4258 store
4259 .get_managed_dns(
4260 ProjectRef::DEFAULT,
4261 &SiteName::new("blog"),
4262 "WWW.example.com."
4263 )
4264 .await
4265 .unwrap(),
4266 Some(ledger.clone())
4267 );
4268 assert_eq!(
4269 store
4270 .list_managed_dns(ProjectRef::DEFAULT, &SiteName::new("blog"))
4271 .await
4272 .unwrap(),
4273 vec![ledger]
4274 );
4275
4276 store
4277 .remove_managed_dns(
4278 ProjectRef::DEFAULT,
4279 &SiteName::new("blog"),
4280 "www.example.com",
4281 )
4282 .await
4283 .unwrap();
4284 assert!(store
4285 .list_managed_dns(ProjectRef::DEFAULT, &SiteName::new("blog"))
4286 .await
4287 .unwrap()
4288 .is_empty());
4289 }
4290
4291 #[tokio::test]
4292 async fn site_config_round_trip_and_host_routing() {
4293 use crate::config::{DomainConfig, SiteConfig};
4294 use crate::kv::MemoryKv;
4295
4296 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
4297 let config = SiteConfig {
4298 domains: DomainConfig {
4299 primary: Some("example.com".into()),
4300 aliases: vec!["www.example.com".into()],
4301 wildcards: vec!["*.example.com".into()],
4302 ..Default::default()
4303 },
4304 ..Default::default()
4305 };
4306 store
4307 .set_site_config(ProjectRef::DEFAULT, "blog", &config)
4308 .await
4309 .unwrap();
4310
4311 let resolved = |host: &'static str| {
4312 let store = store.clone();
4313 async move {
4314 store
4315 .resolve_site_by_host(host)
4316 .await
4317 .unwrap()
4318 .map(|o| o.site)
4319 }
4320 };
4321 assert_eq!(resolved("example.com").await.as_deref(), Some("blog")); assert_eq!(resolved("www.example.com").await.as_deref(), Some("blog")); assert_eq!(resolved("api.example.com").await.as_deref(), Some("blog")); assert_eq!(resolved("a.b.example.com").await.as_deref(), Some("blog")); assert_eq!(resolved("other.com").await, None);
4326
4327 store
4329 .set_site_config(ProjectRef::DEFAULT, "blog", &SiteConfig::default())
4330 .await
4331 .unwrap();
4332 assert_eq!(resolved("example.com").await, None);
4333 }
4334
4335 #[tokio::test]
4340 async fn resolve_cache_invalidates_negative_entry_on_domain_add() {
4341 use crate::config::{DomainConfig, SiteConfig};
4342 use crate::kv::MemoryKv;
4343
4344 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
4345 assert_eq!(
4347 store.resolve_site_by_host("shop.example").await.unwrap(),
4348 None
4349 );
4350 store
4352 .set_site_config(
4353 ProjectRef::DEFAULT,
4354 "shop",
4355 &SiteConfig {
4356 domains: DomainConfig {
4357 primary: Some("shop.example".into()),
4358 ..Default::default()
4359 },
4360 ..Default::default()
4361 },
4362 )
4363 .await
4364 .unwrap();
4365 assert_eq!(
4367 store
4368 .resolve_site_by_host("shop.example")
4369 .await
4370 .unwrap()
4371 .map(|o| o.site)
4372 .as_deref(),
4373 Some("shop")
4374 );
4375 store
4377 .delete_site(ProjectRef::DEFAULT, "shop")
4378 .await
4379 .unwrap();
4380 assert_eq!(
4381 store.resolve_site_by_host("shop.example").await.unwrap(),
4382 None
4383 );
4384 }
4385
4386 #[tokio::test]
4392 async fn wildcard_vhost_precedence_exact_beats_wildcard() {
4393 use crate::config::{DomainConfig, SiteConfig};
4394 use crate::kv::MemoryKv;
4395
4396 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
4397 let attach = |site: &'static str, domains: DomainConfig| {
4398 let store = store.clone();
4399 async move {
4400 store
4401 .set_site_config(
4402 ProjectRef::DEFAULT,
4403 site,
4404 &SiteConfig {
4405 domains,
4406 ..Default::default()
4407 },
4408 )
4409 .await
4410 .unwrap();
4411 }
4412 };
4413 attach(
4415 "portal",
4416 DomainConfig {
4417 wildcards: vec!["*.construens.com".into()],
4418 ..Default::default()
4419 },
4420 )
4421 .await;
4422 attach(
4424 "console",
4425 DomainConfig {
4426 primary: Some("console.construens.com".into()),
4427 ..Default::default()
4428 },
4429 )
4430 .await;
4431 attach(
4433 "vip",
4434 DomainConfig {
4435 primary: Some("vip.construens.com".into()),
4436 ..Default::default()
4437 },
4438 )
4439 .await;
4440
4441 let resolved = |host: &'static str| {
4442 let store = store.clone();
4443 async move {
4444 store
4445 .resolve_site_by_host(host)
4446 .await
4447 .unwrap()
4448 .map(|o| o.site)
4449 }
4450 };
4451 assert_eq!(
4453 resolved("console.construens.com").await.as_deref(),
4454 Some("console")
4455 );
4456 assert_eq!(resolved("vip.construens.com").await.as_deref(), Some("vip"));
4457 assert_eq!(
4459 resolved("tenant7.construens.com").await.as_deref(),
4460 Some("portal")
4461 );
4462 assert_eq!(
4463 resolved("anything-else.construens.com").await.as_deref(),
4464 Some("portal")
4465 );
4466 assert_eq!(
4467 resolved("deep.team.construens.com").await.as_deref(),
4468 Some("portal")
4469 );
4470 assert_eq!(resolved("construens.com").await, None);
4472 assert_eq!(resolved("console.example.com").await, None);
4473 }
4474
4475 #[tokio::test]
4480 async fn wildcard_attaches_and_routes_without_real_dns_admin_override() {
4481 use crate::domain_verify::VerificationMethod;
4482 use crate::kv::MemoryKv;
4483
4484 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
4485 let site = SiteName::new("portal");
4486 store
4489 .start_domain_verification(
4490 ProjectRef::DEFAULT,
4491 &site,
4492 "*.construens.com",
4493 VerificationMethod::Dns,
4494 0,
4495 )
4496 .await
4497 .unwrap();
4498 store
4499 .mark_domain_verified(ProjectRef::DEFAULT, &site, "*.construens.com")
4500 .await
4501 .unwrap();
4502 store
4503 .attach_verified_domain(ProjectRef::DEFAULT, &site, "*.construens.com")
4504 .await
4505 .unwrap();
4506
4507 assert_eq!(
4509 store
4510 .resolve_site_by_host("tenant7.construens.com")
4511 .await
4512 .unwrap()
4513 .map(|o| o.site)
4514 .as_deref(),
4515 Some("portal")
4516 );
4517 assert_eq!(
4519 store
4520 .resolve_site_by_host("construens.com")
4521 .await
4522 .unwrap()
4523 .map(|o| o.site),
4524 None
4525 );
4526 }
4527
4528 #[tokio::test]
4529 async fn site_config_is_content_addressed_and_dedups() {
4530 use crate::config::SiteConfig;
4531 use crate::kv::MemoryKv;
4532
4533 let kv = Arc::new(MemoryKv::new());
4534 let store = DeployStore::new(Arc::new(NullStorage), kv.clone());
4535
4536 let mut cfg = SiteConfig::default();
4540 cfg.security.https_redirect = true;
4541 store
4543 .set_site_config(ProjectRef::DEFAULT, "s1", &cfg)
4544 .await
4545 .unwrap();
4546 let mut cfg2 = cfg.clone();
4547 cfg2.security.https_redirect = false;
4548 store
4549 .set_site_config(ProjectRef::DEFAULT, "s2", &cfg2)
4550 .await
4551 .unwrap();
4552 store
4553 .set_site_config(ProjectRef::DEFAULT, "s3", &cfg)
4554 .await
4555 .unwrap(); let bodies = kv.list_prefix("siteconfig/").await.unwrap();
4559 assert_eq!(bodies.len(), 2, "s1/s3 share a body; s2 distinct");
4560 let pointers = kv.list_prefix("project/default/site/").await.unwrap();
4561 assert_eq!(pointers.len(), 3);
4562
4563 assert!(
4565 store
4566 .get_site_config(ProjectRef::DEFAULT, "s1")
4567 .await
4568 .unwrap()
4569 .unwrap()
4570 .security
4571 .https_redirect
4572 );
4573 assert_eq!(
4574 store
4575 .get_site_config(ProjectRef::DEFAULT, "missing")
4576 .await
4577 .unwrap(),
4578 None
4579 );
4580
4581 let mut edited = cfg.clone();
4584 edited.security.frame_options = Some("DENY".into());
4585 store
4586 .set_site_config(ProjectRef::DEFAULT, "s1", &edited)
4587 .await
4588 .unwrap();
4589 store.collect_garbage(true).await.unwrap();
4590 assert_eq!(kv.list_prefix("siteconfig/").await.unwrap().len(), 3);
4592 store
4594 .set_site_config(ProjectRef::DEFAULT, "s3", &edited)
4595 .await
4596 .unwrap();
4597 store.collect_garbage(true).await.unwrap();
4598 let remaining = kv.list_prefix("siteconfig/").await.unwrap();
4599 assert_eq!(remaining.len(), 2, "orphaned shared body reclaimed");
4600 assert!(
4602 store
4603 .get_site_config(ProjectRef::DEFAULT, "s1")
4604 .await
4605 .unwrap()
4606 .unwrap()
4607 .security
4608 .https_redirect
4609 );
4610 assert!(
4611 !store
4612 .get_site_config(ProjectRef::DEFAULT, "s2")
4613 .await
4614 .unwrap()
4615 .unwrap()
4616 .security
4617 .https_redirect
4618 );
4619 }
4620
4621 #[tokio::test]
4622 async fn cert_status_lists_domains_and_expiry_without_keys() {
4623 use crate::cert::StoredCert;
4624 use crate::kv::MemoryKv;
4625
4626 let kv = Arc::new(MemoryKv::new());
4627 let store = DeployStore::new(Arc::new(NullStorage), kv.clone());
4628 for (domain, not_after) in [("b.example.com", 2000u64), ("a.example.com", 1000u64)] {
4630 let cert = StoredCert::new("CHAINPEM", "KEYPEM", not_after);
4631 kv.put(
4632 &crate::cert::cert_key(domain),
4633 serde_json::to_vec(&cert).unwrap(),
4634 )
4635 .await
4636 .unwrap();
4637 }
4638 kv.put("site/x/config", b"{}".to_vec()).await.unwrap();
4639
4640 let status = store.cert_status().await.unwrap();
4641 assert_eq!(status.len(), 2);
4642 assert_eq!(status[0].domain, "a.example.com");
4644 assert_eq!(status[0].not_after_unix, 1000);
4645 assert_eq!(status[1].domain, "b.example.com");
4646 }
4647
4648 #[tokio::test]
4649 async fn domain_verification_gates_attachment() {
4650 use crate::domain_verify::VerificationMethod;
4651
4652 let store = store();
4653
4654 let v1 = store
4656 .start_domain_verification(
4657 ProjectRef::DEFAULT,
4658 &SiteName::new("blog"),
4659 "example.com",
4660 VerificationMethod::Dns,
4661 100,
4662 )
4663 .await
4664 .unwrap();
4665 let v2 = store
4666 .start_domain_verification(
4667 ProjectRef::DEFAULT,
4668 &SiteName::new("blog"),
4669 "example.com",
4670 VerificationMethod::Dns,
4671 200,
4672 )
4673 .await
4674 .unwrap();
4675 assert_eq!(v1.token, v2.token, "same method → same pending token");
4676 assert!(!store
4677 .is_domain_verified(ProjectRef::DEFAULT, &SiteName::new("blog"), "example.com")
4678 .await
4679 .unwrap());
4680
4681 assert!(store
4683 .attach_verified_domain(ProjectRef::DEFAULT, &SiteName::new("blog"), "example.com")
4684 .await
4685 .is_err());
4686
4687 store
4689 .mark_domain_verified(ProjectRef::DEFAULT, &SiteName::new("blog"), "example.com")
4690 .await
4691 .unwrap();
4692 assert!(store
4693 .is_domain_verified(ProjectRef::DEFAULT, &SiteName::new("blog"), "example.com")
4694 .await
4695 .unwrap());
4696 store
4697 .attach_verified_domain(ProjectRef::DEFAULT, &SiteName::new("blog"), "example.com")
4698 .await
4699 .unwrap();
4700 assert_eq!(
4701 store
4702 .resolve_site_by_host("example.com")
4703 .await
4704 .unwrap()
4705 .map(|o| o.site)
4706 .as_deref(),
4707 Some("blog")
4708 );
4709 store
4711 .start_domain_verification(
4712 ProjectRef::DEFAULT,
4713 &SiteName::new("blog"),
4714 "www.example.com",
4715 VerificationMethod::Http,
4716 300,
4717 )
4718 .await
4719 .unwrap();
4720 store
4721 .mark_domain_verified(
4722 ProjectRef::DEFAULT,
4723 &SiteName::new("blog"),
4724 "www.example.com",
4725 )
4726 .await
4727 .unwrap();
4728 let config = store
4729 .attach_verified_domain(
4730 ProjectRef::DEFAULT,
4731 &SiteName::new("blog"),
4732 "www.example.com",
4733 )
4734 .await
4735 .unwrap();
4736 assert_eq!(config.domains.primary.as_deref(), Some("example.com"));
4737 assert_eq!(config.domains.aliases, vec!["www.example.com".to_string()]);
4738
4739 store
4741 .start_domain_verification(
4742 ProjectRef::DEFAULT,
4743 &SiteName::new("blog"),
4744 "*.example.com",
4745 VerificationMethod::Dns,
4746 400,
4747 )
4748 .await
4749 .unwrap();
4750 assert!(store
4752 .is_domain_verified(ProjectRef::DEFAULT, &SiteName::new("blog"), "*.example.com")
4753 .await
4754 .unwrap());
4755 let config = store
4756 .attach_verified_domain(ProjectRef::DEFAULT, &SiteName::new("blog"), "*.example.com")
4757 .await
4758 .unwrap();
4759 assert_eq!(config.domains.wildcards, vec!["*.example.com".to_string()]);
4760
4761 assert_eq!(
4763 store
4764 .list_domain_verifications(ProjectRef::DEFAULT, &SiteName::new("blog"))
4765 .await
4766 .unwrap()
4767 .len(),
4768 2
4769 );
4770 assert!(store
4771 .remove_domain_verification(ProjectRef::DEFAULT, &SiteName::new("blog"), "example.com")
4772 .await
4773 .unwrap());
4774 assert!(!store
4775 .is_domain_verified(ProjectRef::DEFAULT, &SiteName::new("blog"), "example.com")
4776 .await
4777 .unwrap());
4778 }
4779
4780 #[tokio::test]
4784 async fn host_cannot_be_hijacked_across_sites() {
4785 use crate::config::{DomainConfig, SiteConfig};
4786 use crate::domain_verify::VerificationMethod;
4787
4788 let store = store();
4789
4790 store
4792 .start_domain_verification(
4793 ProjectRef::DEFAULT,
4794 &SiteName::new("a"),
4795 "shared.example",
4796 VerificationMethod::Http,
4797 100,
4798 )
4799 .await
4800 .unwrap();
4801 store
4802 .mark_domain_verified(ProjectRef::DEFAULT, &SiteName::new("a"), "shared.example")
4803 .await
4804 .unwrap();
4805 store
4806 .attach_verified_domain(ProjectRef::DEFAULT, &SiteName::new("a"), "shared.example")
4807 .await
4808 .unwrap();
4809 assert_eq!(
4810 store
4811 .resolve_site_by_host("shared.example")
4812 .await
4813 .unwrap()
4814 .map(|o| o.site)
4815 .as_deref(),
4816 Some("a")
4817 );
4818
4819 store
4823 .start_domain_verification(
4824 ProjectRef::DEFAULT,
4825 &SiteName::new("b"),
4826 "shared.example",
4827 VerificationMethod::Http,
4828 200,
4829 )
4830 .await
4831 .unwrap();
4832 store
4833 .mark_domain_verified(ProjectRef::DEFAULT, &SiteName::new("b"), "shared.example")
4834 .await
4835 .unwrap();
4836 let err = store
4837 .attach_verified_domain(ProjectRef::DEFAULT, &SiteName::new("b"), "shared.example")
4838 .await
4839 .expect_err("second site must not hijack an attached host");
4840 assert!(matches!(err, DeployError::Conflict(_)), "got {err:?}");
4841
4842 let stolen = SiteConfig {
4845 domains: DomainConfig {
4846 primary: Some("shared.example".into()),
4847 ..Default::default()
4848 },
4849 ..Default::default()
4850 };
4851 let err = store
4852 .set_site_config(ProjectRef::DEFAULT, "b", &stolen)
4853 .await
4854 .expect_err("set_site_config must refuse another site's host");
4855 assert!(matches!(err, DeployError::Conflict(_)), "got {err:?}");
4856
4857 assert_eq!(
4859 store
4860 .resolve_site_by_host("shared.example")
4861 .await
4862 .unwrap()
4863 .map(|o| o.site)
4864 .as_deref(),
4865 Some("a")
4866 );
4867
4868 let readd = SiteConfig {
4871 domains: DomainConfig {
4872 primary: Some("shared.example".into()),
4873 aliases: vec!["www.shared.example".into()],
4874 ..Default::default()
4875 },
4876 ..Default::default()
4877 };
4878 store
4879 .set_site_config(ProjectRef::DEFAULT, "a", &readd)
4880 .await
4881 .unwrap();
4882 assert_eq!(
4883 store
4884 .resolve_site_by_host("www.shared.example")
4885 .await
4886 .unwrap()
4887 .map(|o| o.site)
4888 .as_deref(),
4889 Some("a")
4890 );
4891 }
4892
4893 #[tokio::test]
4897 async fn domain_context_tag_is_written_and_inherited() {
4898 use crate::config::{DomainConfig, SiteConfig};
4899
4900 let store = store();
4901 let cfg = SiteConfig {
4902 domains: DomainConfig {
4903 primary: Some("acme-store.com".into()),
4904 aliases: vec!["www.acme-store.com".into()],
4905 wildcards: vec!["*.globex-store.com".into()],
4906 contexts: std::collections::BTreeMap::from([
4907 ("acme-store.com".into(), "acme".into()),
4908 ("*.globex-store.com".into(), "globex".into()),
4909 ]),
4910 ..Default::default()
4911 },
4912 ..Default::default()
4913 };
4914 store
4915 .set_site_config(ProjectRef::DEFAULT, "shop", &cfg)
4916 .await
4917 .unwrap();
4918
4919 async fn ctx(store: &DeployStore, host: &str) -> Option<String> {
4920 store
4921 .resolve_site_by_host(host)
4922 .await
4923 .unwrap()
4924 .and_then(|o| o.context)
4925 }
4926 assert_eq!(ctx(&store, "acme-store.com").await.as_deref(), Some("acme"));
4929 assert_eq!(
4930 ctx(&store, "www.acme-store.com").await.as_deref(),
4931 Some("acme")
4932 );
4933 assert_eq!(
4934 ctx(&store, "tenant7.globex-store.com").await.as_deref(),
4935 Some("globex")
4936 );
4937 }
4938
4939 #[tokio::test]
4943 async fn host_uniqueness_is_case_and_dot_insensitive() {
4944 use crate::config::{DomainConfig, SiteConfig};
4945 use crate::domain_verify::VerificationMethod;
4946
4947 let store = store();
4948 store
4950 .start_domain_verification(
4951 ProjectRef::DEFAULT,
4952 &SiteName::new("a"),
4953 "example.com",
4954 VerificationMethod::Http,
4955 100,
4956 )
4957 .await
4958 .unwrap();
4959 store
4960 .mark_domain_verified(ProjectRef::DEFAULT, &SiteName::new("a"), "example.com")
4961 .await
4962 .unwrap();
4963 store
4964 .attach_verified_domain(ProjectRef::DEFAULT, &SiteName::new("a"), "example.com")
4965 .await
4966 .unwrap();
4967
4968 for variant in ["Example.COM", "example.com.", "EXAMPLE.com."] {
4970 let cfg = SiteConfig {
4971 domains: DomainConfig {
4972 primary: Some(variant.into()),
4973 ..Default::default()
4974 },
4975 ..Default::default()
4976 };
4977 let err = store
4978 .set_site_config(ProjectRef::DEFAULT, "b", &cfg)
4979 .await
4980 .expect_err("variant claim must be refused");
4981 assert!(
4982 matches!(err, DeployError::Conflict(_)),
4983 "variant {variant:?} must Conflict, got {err:?}"
4984 );
4985 }
4986
4987 for h in ["example.com", "Example.com", "EXAMPLE.COM", "example.com."] {
4989 assert_eq!(
4990 store
4991 .resolve_site_by_host(h)
4992 .await
4993 .unwrap()
4994 .map(|o| o.site)
4995 .as_deref(),
4996 Some("a"),
4997 "host {h:?} must resolve to site a"
4998 );
4999 }
5000 }
5001
5002 #[tokio::test]
5006 async fn self_serve_challenge_lookup_matches_pending_http_only() {
5007 use crate::domain_verify::{VerificationMethod, CHALLENGE_TTL_SECS};
5008
5009 let store = store();
5010 let v = store
5011 .start_domain_verification(
5012 ProjectRef::DEFAULT,
5013 &SiteName::new("docs"),
5014 "docs.example",
5015 VerificationMethod::Http,
5016 1_000,
5017 )
5018 .await
5019 .unwrap();
5020
5021 let found = store
5023 .find_pending_http_challenge("docs.example", &v.token, 1_000)
5024 .await
5025 .unwrap();
5026 assert_eq!(
5027 found.as_ref().map(|f| f.token.clone()),
5028 Some(v.token.clone())
5029 );
5030 assert!(store
5032 .find_pending_http_challenge("Docs.Example.", &v.token, 1_000)
5033 .await
5034 .unwrap()
5035 .is_some());
5036
5037 assert!(store
5039 .find_pending_http_challenge("docs.example", "not-the-token", 1_000)
5040 .await
5041 .unwrap()
5042 .is_none());
5043 assert!(store
5044 .find_pending_http_challenge("other.example", &v.token, 1_000)
5045 .await
5046 .unwrap()
5047 .is_none());
5048
5049 assert!(store
5051 .find_pending_http_challenge("docs.example", &v.token, 1_000 + CHALLENGE_TTL_SECS + 1)
5052 .await
5053 .unwrap()
5054 .is_none());
5055
5056 let dv = store
5058 .start_domain_verification(
5059 ProjectRef::DEFAULT,
5060 &SiteName::new("dns-site"),
5061 "dns.example",
5062 VerificationMethod::Dns,
5063 1_000,
5064 )
5065 .await
5066 .unwrap();
5067 assert!(store
5068 .find_pending_http_challenge("dns.example", &dv.token, 1_000)
5069 .await
5070 .unwrap()
5071 .is_none());
5072 }
5073
5074 #[tokio::test]
5078 async fn wildcard_requires_dns_and_stale_index_is_safe() {
5079 use crate::domain_verify::VerificationMethod;
5080
5081 let store = store();
5082
5083 let http = store
5085 .start_domain_verification(
5086 ProjectRef::DEFAULT,
5087 &SiteName::new("s"),
5088 "*.example.com",
5089 VerificationMethod::Http,
5090 100,
5091 )
5092 .await
5093 .unwrap();
5094 store
5095 .mark_domain_verified(ProjectRef::DEFAULT, &SiteName::new("s"), "*.example.com")
5096 .await
5097 .unwrap();
5098 let err = store
5099 .attach_verified_domain(ProjectRef::DEFAULT, &SiteName::new("s"), "*.example.com")
5100 .await
5101 .expect_err("wildcard with only HTTP proof must be refused");
5102 assert!(matches!(err, DeployError::Conflict(_)), "got {err:?}");
5103
5104 store
5108 .remove_domain_verification(ProjectRef::DEFAULT, &SiteName::new("s"), "*.example.com")
5109 .await
5110 .unwrap();
5111 assert!(store
5113 .find_pending_http_challenge("example.com", &http.token, 100)
5114 .await
5115 .unwrap()
5116 .is_none());
5117 store
5118 .start_domain_verification(
5119 ProjectRef::DEFAULT,
5120 &SiteName::new("s"),
5121 "*.example.com",
5122 VerificationMethod::Dns,
5123 200,
5124 )
5125 .await
5126 .unwrap();
5127 store
5128 .mark_domain_verified(ProjectRef::DEFAULT, &SiteName::new("s"), "*.example.com")
5129 .await
5130 .unwrap();
5131 let cfg = store
5132 .attach_verified_domain(ProjectRef::DEFAULT, &SiteName::new("s"), "*.example.com")
5133 .await
5134 .unwrap();
5135 assert_eq!(cfg.domains.wildcards, vec!["*.example.com".to_string()]);
5136
5137 let h2 = store
5141 .start_domain_verification(
5142 ProjectRef::DEFAULT,
5143 &SiteName::new("s2"),
5144 "host.example",
5145 VerificationMethod::Http,
5146 300,
5147 )
5148 .await
5149 .unwrap();
5150 assert!(store
5151 .find_pending_http_challenge("host.example", &h2.token, 300)
5152 .await
5153 .unwrap()
5154 .is_some());
5155 store
5156 .start_domain_verification(
5157 ProjectRef::DEFAULT,
5158 &SiteName::new("s2"),
5159 "host.example",
5160 VerificationMethod::Dns,
5161 300,
5162 )
5163 .await
5164 .unwrap();
5165 assert!(
5166 store
5167 .find_pending_http_challenge("host.example", &h2.token, 300)
5168 .await
5169 .unwrap()
5170 .is_none(),
5171 "a stale HTTP index must not serve a token whose record is now DNS"
5172 );
5173 }
5174
5175 #[tokio::test]
5176 async fn pending_verifications_are_capped_per_site() {
5177 let store = store();
5178 for i in 0..64 {
5180 store
5181 .start_domain_verification(
5182 ProjectRef::DEFAULT,
5183 &SiteName::new("site"),
5184 &format!("h{i}.example"),
5185 VerificationMethod::Http,
5186 100,
5187 )
5188 .await
5189 .unwrap();
5190 }
5191 let err = store
5193 .start_domain_verification(
5194 ProjectRef::DEFAULT,
5195 &SiteName::new("site"),
5196 "overflow.example",
5197 VerificationMethod::Http,
5198 100,
5199 )
5200 .await
5201 .expect_err("the 65th pending host must be rejected");
5202 assert!(matches!(err, DeployError::Conflict(_)), "got {err:?}");
5203 store
5205 .start_domain_verification(
5206 ProjectRef::DEFAULT,
5207 &SiteName::new("site"),
5208 "h0.example",
5209 VerificationMethod::Http,
5210 100,
5211 )
5212 .await
5213 .expect("re-running an existing challenge is not capped");
5214 store
5216 .start_domain_verification(
5217 ProjectRef::DEFAULT,
5218 &SiteName::new("other"),
5219 "fresh.example",
5220 VerificationMethod::Http,
5221 100,
5222 )
5223 .await
5224 .expect("a different site is unaffected");
5225 }
5226
5227 fn store() -> DeployStore {
5228 use crate::kv::MemoryKv;
5229 DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()))
5230 }
5231
5232 #[tokio::test]
5233 async fn default_project_is_visible_on_a_fresh_store_and_ensure_is_idempotent() {
5234 let s = store();
5235 let default = crate::project::DEFAULT_PROJECT;
5236
5237 assert!(
5239 s.get_project(default).await.unwrap().is_some(),
5240 "`project show default` must never 404, even on a fresh store"
5241 );
5242 let listed = s.list_projects().await.unwrap();
5243 assert_eq!(
5244 listed.iter().filter(|p| p.name == default).count(),
5245 1,
5246 "`project ls` must show exactly one `default` on a fresh store"
5247 );
5248 assert!(s.get_project("nope").await.unwrap().is_none());
5250
5251 assert!(s.project_exists(default).await.unwrap());
5254 assert!(!s.project_exists("nope").await.unwrap());
5255
5256 assert!(
5258 s.ensure_default_project().await.unwrap(),
5259 "first ensure creates"
5260 );
5261 assert!(
5262 !s.ensure_default_project().await.unwrap(),
5263 "second ensure is idempotent (presence-checked)"
5264 );
5265
5266 let listed = s.list_projects().await.unwrap();
5268 assert_eq!(
5269 listed.iter().filter(|p| p.name == default).count(),
5270 1,
5271 "materializing `default` must not duplicate it in the listing"
5272 );
5273 assert_eq!(s.get_project(default).await.unwrap().unwrap().name, default);
5274 }
5275
5276 #[tokio::test]
5277 async fn daemon_config_store_round_trips_and_rolls_back() {
5278 use crate::daemon_config::DaemonConfig;
5279 let s = store();
5280 assert!(s.get_daemon_config().await.unwrap().is_none());
5282 assert!(s.daemon_config_generation().await.unwrap().is_none());
5283
5284 let g1cfg = DaemonConfig {
5286 default_site: Some("one".into()),
5287 ..Default::default()
5288 };
5289 let g1 = s.set_daemon_config(&g1cfg).await.unwrap();
5290 assert_eq!(
5291 s.daemon_config_generation().await.unwrap().as_deref(),
5292 Some(g1.as_str())
5293 );
5294 assert_eq!(s.get_daemon_config().await.unwrap().unwrap(), g1cfg);
5295 assert!(s.daemon_config_history().await.unwrap().is_empty());
5296
5297 let g2cfg = DaemonConfig {
5299 default_site: Some("two".into()),
5300 ..Default::default()
5301 };
5302 let g2 = s.set_daemon_config(&g2cfg).await.unwrap();
5303 assert_ne!(g1, g2);
5304 assert_eq!(s.daemon_config_history().await.unwrap(), vec![g1.clone()]);
5305
5306 let rolled = s.rollback_daemon_config().await.unwrap();
5308 assert_eq!(rolled.as_deref(), Some(g1.as_str()));
5309 assert_eq!(
5310 s.daemon_config_generation().await.unwrap().as_deref(),
5311 Some(g1.as_str())
5312 );
5313 assert_eq!(s.get_daemon_config().await.unwrap().unwrap(), g1cfg);
5314 assert!(s.rollback_daemon_config().await.unwrap().is_none());
5316 }
5317
5318 #[tokio::test]
5319 async fn compute_store_round_trips() {
5320 use crate::compute::{
5321 ComputeSpec, ComputeWorkload, PlacementConstraints, RestartPolicy, RootSource,
5322 };
5323 let s = store();
5324 let spec = ComputeSpec {
5325 version: crate::SCHEMA_VERSION,
5326 root: RootSource::Rootfs("r".repeat(64)),
5327 kernel: "k".repeat(64),
5328 kernel_cmdline: None,
5329 vcpus: 1,
5330 mem_mib: 256,
5331 entrypoint: vec!["/app".into()],
5332 env: Default::default(),
5333 port: 8080,
5334 restart: RestartPolicy::Always,
5335 startup_grace_secs: 30,
5336 scale_to_zero: false,
5337 volumes: vec![],
5338 writable_root: false,
5339 cap_add: Vec::new(),
5340 user: None,
5341 isolation: Default::default(),
5342 prefer_backend: None,
5343 bindings: vec![],
5344 };
5345 let hash = s.put_compute_spec(&spec).await.unwrap();
5347 assert_eq!(hash, spec.id());
5348 assert_eq!(s.get_compute_spec(&hash).await.unwrap(), Some(spec));
5349 assert!(s.get_compute_spec("deadbeef").await.unwrap().is_none());
5350
5351 let workload = ComputeWorkload {
5353 version: crate::SCHEMA_VERSION,
5354 name: "api".into(),
5355 active: hash.clone(),
5356 replicas: 3,
5357 placement: PlacementConstraints::default(),
5358 };
5359 s.set_compute_workload(ProjectRef::DEFAULT, &workload)
5360 .await
5361 .unwrap();
5362 assert_eq!(
5363 s.get_compute_workload(ProjectRef::DEFAULT, "api")
5364 .await
5365 .unwrap(),
5366 Some(workload)
5367 );
5368 assert_eq!(
5369 s.list_compute_workloads(ProjectRef::DEFAULT)
5370 .await
5371 .unwrap()
5372 .len(),
5373 1
5374 );
5375 assert!(s
5376 .delete_compute_workload(ProjectRef::DEFAULT, "api")
5377 .await
5378 .unwrap());
5379 assert!(!s
5380 .delete_compute_workload(ProjectRef::DEFAULT, "api")
5381 .await
5382 .unwrap());
5383 assert!(s
5384 .list_compute_workloads(ProjectRef::DEFAULT)
5385 .await
5386 .unwrap()
5387 .is_empty());
5388 }
5389
5390 fn empty_manifest(clean_urls: bool) -> Manifest {
5393 Manifest {
5394 config: DeployConfig {
5395 clean_urls,
5396 ..DeployConfig::default()
5397 },
5398 ..Default::default()
5399 }
5400 }
5401
5402 #[tokio::test]
5403 async fn deploy_meta_records_sizes_and_merges_provenance() {
5404 let store = store();
5405 let mut manifest = Manifest::default();
5406 manifest.files.insert("index.html".into(), {
5407 let mut e = entry("aa");
5408 e.size = 10;
5409 e
5410 });
5411 manifest.files.insert("app.js".into(), {
5412 let mut e = entry("bb");
5413 e.size = 32;
5414 e
5415 });
5416
5417 let id = store
5419 .put_manifest_with(
5420 &manifest,
5421 DeployMetaInput {
5422 source: Some("abc123".into()),
5423 message: Some("first".into()),
5424 ..Default::default()
5425 },
5426 )
5427 .await
5428 .unwrap();
5429 let meta = store.get_meta(&id).await.unwrap().unwrap();
5430 assert_eq!(meta.file_count, 2);
5431 assert_eq!(meta.total_size, 42);
5432 assert_eq!(meta.source.as_deref(), Some("abc123"));
5433 let created = meta.created_at;
5434
5435 store
5437 .put_manifest_with(&manifest, DeployMetaInput::default())
5438 .await
5439 .unwrap();
5440 let meta = store.get_meta(&id).await.unwrap().unwrap();
5441 assert_eq!(meta.created_at, created);
5442 assert_eq!(meta.source.as_deref(), Some("abc123"));
5443 assert_eq!(meta.message.as_deref(), Some("first"));
5444 }
5445
5446 #[tokio::test]
5447 async fn aliases_round_trip_and_guard_completeness() {
5448 let store = store();
5449 let manifest = empty_manifest(false);
5450 let id = store.put_manifest(&manifest).await.unwrap();
5451
5452 assert!(matches!(
5454 store
5455 .set_alias(ProjectRef::DEFAULT, "blog", "staging", "deadbeef")
5456 .await,
5457 Err(DeployError::NotFound(_))
5458 ));
5459
5460 store
5461 .set_alias(ProjectRef::DEFAULT, "blog", "staging", &id)
5462 .await
5463 .unwrap();
5464 assert_eq!(
5465 store
5466 .get_alias(ProjectRef::DEFAULT, "blog", "staging")
5467 .await
5468 .unwrap(),
5469 Some(id.clone())
5470 );
5471 let aliases = store
5472 .list_aliases(ProjectRef::DEFAULT, "blog")
5473 .await
5474 .unwrap();
5475 assert_eq!(aliases.get("staging"), Some(&id));
5476
5477 assert!(store
5478 .remove_alias(ProjectRef::DEFAULT, "blog", "staging")
5479 .await
5480 .unwrap());
5481 assert!(!store
5482 .remove_alias(ProjectRef::DEFAULT, "blog", "staging")
5483 .await
5484 .unwrap());
5485 assert_eq!(
5486 store
5487 .get_alias(ProjectRef::DEFAULT, "blog", "staging")
5488 .await
5489 .unwrap(),
5490 None
5491 );
5492 }
5493
5494 #[tokio::test]
5495 async fn retention_keep_last_collects_older_history() {
5496 let store = store();
5497 let m1 = empty_manifest(false);
5499 let m2 = empty_manifest(true);
5500 let mut m3 = empty_manifest(true);
5501 m3.config.trailing_slash = crate::config::TrailingSlash::Always;
5502 let id1 = store.put_manifest(&m1).await.unwrap();
5503 let id2 = store.put_manifest(&m2).await.unwrap();
5504 let id3 = store.put_manifest(&m3).await.unwrap();
5505 store
5506 .activate(ProjectRef::DEFAULT, "blog", &id1)
5507 .await
5508 .unwrap();
5509 store
5510 .activate(ProjectRef::DEFAULT, "blog", &id2)
5511 .await
5512 .unwrap();
5513 store
5514 .activate(ProjectRef::DEFAULT, "blog", &id3)
5515 .await
5516 .unwrap();
5517
5518 let report = store.collect_garbage(false).await.unwrap();
5520 assert_eq!(report.manifests_removed, 0);
5521
5522 store
5525 .set_alias(ProjectRef::DEFAULT, "blog", "pinned", &id1)
5526 .await
5527 .unwrap();
5528 let report = store
5529 .collect_garbage_with(
5530 false,
5531 GcOptions {
5532 keep_last: Some(1),
5533 ..Default::default()
5534 },
5535 )
5536 .await
5537 .unwrap();
5538 assert_eq!(report.manifests_removed, 1); }
5540
5541 #[tokio::test]
5542 async fn grace_window_protects_in_flight_manifest() {
5543 let store = store();
5544 let id = store.put_manifest(&empty_manifest(false)).await.unwrap();
5546 assert!(store.get_manifest(&id).await.unwrap().is_some());
5547
5548 let report = store
5550 .collect_garbage_with(
5551 false,
5552 GcOptions {
5553 grace_secs: 3600,
5554 ..Default::default()
5555 },
5556 )
5557 .await
5558 .unwrap();
5559 assert_eq!(report.manifests_removed, 0);
5560
5561 let report = store.collect_garbage(false).await.unwrap();
5563 assert_eq!(report.manifests_removed, 1);
5564 }
5565
5566 #[tokio::test]
5567 async fn resolve_manifest_id_exact_prefix_and_missing() {
5568 let store = store();
5569 let id = store.put_manifest(&empty_manifest(true)).await.unwrap();
5570
5571 assert_eq!(
5573 store.resolve_manifest_id(&id).await.unwrap().as_deref(),
5574 Some(id.as_str())
5575 );
5576 assert_eq!(
5578 store
5579 .resolve_manifest_id(&id[..16])
5580 .await
5581 .unwrap()
5582 .as_deref(),
5583 Some(id.as_str())
5584 );
5585 assert!(store
5587 .resolve_manifest_id("ffffffffffffffff")
5588 .await
5589 .unwrap()
5590 .is_none());
5591 }
5592
5593 #[tokio::test]
5594 async fn delete_site_removes_config_routing_and_aliases() {
5595 let store = store();
5596 let mut cfg = SiteConfig::default();
5597 cfg.domains.primary = Some("blog.example".into());
5598 cfg.domains.wildcards = vec!["*.preview.blog.example".into()];
5599 store
5600 .set_site_config(ProjectRef::DEFAULT, "blog", &cfg)
5601 .await
5602 .unwrap();
5603 let id = store.put_manifest(&empty_manifest(true)).await.unwrap();
5605 store
5606 .set_alias(ProjectRef::DEFAULT, "blog", "stable", &id)
5607 .await
5608 .unwrap();
5609
5610 assert!(store
5612 .get_site_config(ProjectRef::DEFAULT, "blog")
5613 .await
5614 .unwrap()
5615 .is_some());
5616 assert_eq!(
5617 store
5618 .resolve_site_by_host("blog.example")
5619 .await
5620 .unwrap()
5621 .map(|o| o.site)
5622 .as_deref(),
5623 Some("blog")
5624 );
5625
5626 store
5627 .delete_site(ProjectRef::DEFAULT, "blog")
5628 .await
5629 .unwrap();
5630
5631 assert!(store
5633 .get_site_config(ProjectRef::DEFAULT, "blog")
5634 .await
5635 .unwrap()
5636 .is_none());
5637 assert!(store
5638 .resolve_site_by_host("blog.example")
5639 .await
5640 .unwrap()
5641 .is_none());
5642 assert!(store
5643 .list_aliases(ProjectRef::DEFAULT, "blog")
5644 .await
5645 .unwrap()
5646 .is_empty());
5647
5648 store
5650 .delete_site(ProjectRef::DEFAULT, "blog")
5651 .await
5652 .unwrap();
5653 }
5654
5655 #[tokio::test]
5656 async fn gc_blob_reachability_is_a_cross_project_union() {
5657 use crate::kv::MemoryKv;
5658 let store = DeployStore::new(Arc::new(MemStorage::default()), Arc::new(MemoryKv::new()));
5660
5661 let shared = sha256_hex(b"shared-blob");
5665 let keep = sha256_hex(b"keep-blob");
5666 let dead = sha256_hex(b"dead-blob");
5667 store
5668 .put_blob(&shared, once_bytes(b"shared-blob"))
5669 .await
5670 .unwrap();
5671 store
5672 .put_blob(&keep, once_bytes(b"keep-blob"))
5673 .await
5674 .unwrap();
5675 store
5676 .put_blob(&dead, once_bytes(b"dead-blob"))
5677 .await
5678 .unwrap();
5679
5680 let m_shared = manifest_with(&[("index.html", &shared)]);
5681 let m_keep = manifest_with(&[("index.html", &keep)]);
5682 let m_dead = manifest_with(&[("old.html", &dead)]);
5683 let id_shared = store.put_manifest(&m_shared).await.unwrap();
5684 let id_keep = store.put_manifest(&m_keep).await.unwrap();
5685 let id_dead = store.put_manifest(&m_dead).await.unwrap();
5686
5687 let acme = ProjectRef::new("acme");
5689 store.activate(acme, "site", &id_shared).await.unwrap();
5690 let shop = ProjectRef::new("shop");
5694 store.activate(shop, "site", &id_shared).await.unwrap();
5695 store.activate(shop, "site", &id_keep).await.unwrap();
5696 let report = store
5699 .collect_garbage_with(
5700 true,
5701 GcOptions {
5702 keep_last: Some(1),
5703 ..Default::default()
5704 },
5705 )
5706 .await
5707 .unwrap();
5708
5709 assert!(store.get_manifest(&id_dead).await.unwrap().is_none());
5711 assert!(
5712 !store.has_blob(&dead).await.unwrap(),
5713 "dead blob is reclaimed"
5714 );
5715 assert_eq!(report.manifests_removed, 1, "only the dead manifest goes");
5716
5717 assert!(
5720 store.get_manifest(&id_shared).await.unwrap().is_some(),
5721 "shared manifest kept by acme's reference"
5722 );
5723 assert!(
5724 store.has_blob(&shared).await.unwrap(),
5725 "shared blob kept — reachability is the union across projects"
5726 );
5727 assert!(store.has_blob(&keep).await.unwrap());
5729 assert_eq!(
5730 store.current_id(acme, "site").await.unwrap().as_deref(),
5731 Some(id_shared.as_str())
5732 );
5733 }
5734
5735 #[tokio::test]
5736 async fn same_site_name_in_two_projects_is_isolated() {
5737 use crate::kv::MemoryKv;
5738 let store = DeployStore::new(Arc::new(NullStorage), Arc::new(MemoryKv::new()));
5739 let acme = ProjectRef::new("acme");
5740 let shop = ProjectRef::new("shop");
5741
5742 let id_a = store.put_manifest(&empty_manifest(false)).await.unwrap();
5744 let id_b = store.put_manifest(&empty_manifest(true)).await.unwrap();
5745 store.activate(acme, "www", &id_a).await.unwrap();
5746 store.activate(shop, "www", &id_b).await.unwrap();
5747
5748 assert_eq!(
5750 store.current_id(acme, "www").await.unwrap().as_deref(),
5751 Some(id_a.as_str())
5752 );
5753 assert_eq!(
5754 store.current_id(shop, "www").await.unwrap().as_deref(),
5755 Some(id_b.as_str())
5756 );
5757 assert_ne!(id_a, id_b);
5758
5759 let mut cfg_a = SiteConfig::default();
5761 cfg_a.domains.primary = Some("acme.example".into());
5762 store.set_site_config(acme, "www", &cfg_a).await.unwrap();
5763 let mut cfg_b = SiteConfig::default();
5764 cfg_b.domains.primary = Some("shop.example".into());
5765 store.set_site_config(shop, "www", &cfg_b).await.unwrap();
5766 store
5767 .set_alias(acme, "www", "staging", &id_a)
5768 .await
5769 .unwrap();
5770
5771 assert_eq!(
5772 store
5773 .get_site_config(acme, "www")
5774 .await
5775 .unwrap()
5776 .unwrap()
5777 .domains
5778 .primary
5779 .as_deref(),
5780 Some("acme.example")
5781 );
5782 assert_eq!(
5783 store
5784 .get_site_config(shop, "www")
5785 .await
5786 .unwrap()
5787 .unwrap()
5788 .domains
5789 .primary
5790 .as_deref(),
5791 Some("shop.example")
5792 );
5793 assert!(store
5795 .get_alias(acme, "www", "staging")
5796 .await
5797 .unwrap()
5798 .is_some());
5799 assert!(store
5800 .get_alias(shop, "www", "staging")
5801 .await
5802 .unwrap()
5803 .is_none());
5804 assert!(store.list_aliases(shop, "www").await.unwrap().is_empty());
5805
5806 assert_eq!(
5808 store.resolve_site_by_host("acme.example").await.unwrap(),
5809 Some(DomainOwner::new("acme", "www"))
5810 );
5811 assert_eq!(
5812 store.resolve_site_by_host("shop.example").await.unwrap(),
5813 Some(DomainOwner::new("shop", "www"))
5814 );
5815 store.delete_site(acme, "www").await.unwrap();
5816 assert!(store.get_site_config(acme, "www").await.unwrap().is_none());
5817 assert!(
5818 store.get_site_config(shop, "www").await.unwrap().is_some(),
5819 "deleting acme/www must not touch shop/www"
5820 );
5821 assert!(store.current_id(shop, "www").await.unwrap().is_some());
5822 }
5823
5824 #[tokio::test]
5825 async fn project_entity_crud_and_delete_guard() {
5826 use crate::project::Project;
5827 let store = store();
5828 let acme = Project {
5829 version: crate::SCHEMA_VERSION,
5830 name: "acme".into(),
5831 created_at: 1,
5832 meta: Default::default(),
5833 config: Default::default(),
5834 secrets_ref: None,
5835 };
5836 let hash = store.put_project(&acme).await.unwrap();
5838 assert_eq!(hash, acme.id());
5839 assert_eq!(store.get_project("acme").await.unwrap(), Some(acme.clone()));
5840 assert!(store.get_project("ghost").await.unwrap().is_none());
5841 let names: Vec<String> = store
5843 .list_projects()
5844 .await
5845 .unwrap()
5846 .into_iter()
5847 .map(|p| p.name)
5848 .collect();
5849 assert_eq!(names, vec!["acme".to_string(), "default".to_string()]);
5850
5851 store
5853 .set_site_config(ProjectRef::new("acme"), "www", &SiteConfig::default())
5854 .await
5855 .unwrap();
5856 assert!(matches!(
5857 store.delete_project("acme").await,
5858 Err(DeployError::Conflict(_))
5859 ));
5860 store
5862 .delete_site(ProjectRef::new("acme"), "www")
5863 .await
5864 .unwrap();
5865 assert!(store.delete_project("acme").await.unwrap());
5866 assert!(store.get_project("acme").await.unwrap().is_none());
5867
5868 assert!(matches!(
5870 store.delete_project("default").await,
5871 Err(DeployError::Conflict(_))
5872 ));
5873 }
5874
5875 #[tokio::test]
5876 async fn open_blob_cached_caches_small_and_streams_large() {
5877 use crate::kv::MemoryKv;
5878 let store = DeployStore::new(Arc::new(MemStorage::default()), Arc::new(MemoryKv::new()));
5879
5880 let small = b"hello, static hot path";
5883 let small_hash = sha256_hex(small);
5884 store
5885 .put_blob(&small_hash, once_bytes(small))
5886 .await
5887 .unwrap();
5888 match store
5889 .open_blob_cached(&small_hash, small.len() as u64)
5890 .await
5891 .unwrap()
5892 {
5893 BlobBody::Cached(bytes) => assert_eq!(&bytes[..], small),
5894 BlobBody::Stream(_) => panic!("small blob should be cached, not streamed"),
5895 BlobBody::Mapped(_) => panic!("small blob should be cached, not mapped"),
5896 }
5897 {
5898 let cache = store.blob_body_cache.read().unwrap();
5899 assert!(cache.map.contains_key(&small_hash));
5900 assert_eq!(cache.bytes, small.len());
5901 }
5902 assert!(matches!(
5903 store
5904 .open_blob_cached(&small_hash, small.len() as u64)
5905 .await
5906 .unwrap(),
5907 BlobBody::Cached(_)
5908 ));
5909
5910 let large = vec![7u8; SMALL_BLOB_CACHE_MAX as usize + 1];
5912 let large_hash = sha256_hex(&large);
5913 let put_body = {
5914 let large = large.clone();
5915 futures::stream::once(async move { Ok(bytes::Bytes::from(large)) }).boxed()
5916 };
5917 store.put_blob(&large_hash, put_body).await.unwrap();
5918 match store
5919 .open_blob_cached(&large_hash, large.len() as u64)
5920 .await
5921 .unwrap()
5922 {
5923 BlobBody::Stream(object) => {
5924 let mut body = object.body;
5925 let mut got = Vec::new();
5926 while let Some(chunk) = body.next().await {
5927 got.extend_from_slice(&chunk.unwrap());
5928 }
5929 assert_eq!(got, large);
5930 }
5931 BlobBody::Cached(_) => panic!("large blob must stream, not cache"),
5932 BlobBody::Mapped(_) => panic!("MemStorage can't memory-map; expected a stream"),
5934 }
5935 assert!(!store
5936 .blob_body_cache
5937 .read()
5938 .unwrap()
5939 .map
5940 .contains_key(&large_hash));
5941 }
5942
5943 #[tokio::test]
5944 async fn enumerate_and_purge_project_are_scoped() {
5945 use crate::compute::{
5946 ComputeSpec, ComputeWorkload, PlacementConstraints, RestartPolicy, RootSource,
5947 VolumeRef,
5948 };
5949 use crate::function::{Function, FunctionConfig, Lifecycle, Owner};
5950 use crate::kv::MemoryKv;
5951 use crate::project::Project;
5952
5953 let kv = Arc::new(MemoryKv::new());
5954 let store = DeployStore::new(Arc::new(NullStorage), kv.clone());
5955 let acme = ProjectRef::new("acme");
5956 let other = ProjectRef::new("other");
5957
5958 store
5960 .put_project(&Project {
5961 version: crate::SCHEMA_VERSION,
5962 name: "acme".into(),
5963 created_at: 1,
5964 meta: Default::default(),
5965 config: Default::default(),
5966 secrets_ref: None,
5967 })
5968 .await
5969 .unwrap();
5970
5971 let mut cfg = SiteConfig::default();
5973 cfg.domains.primary = Some("acme.example".into());
5974 store.set_site_config(acme, "www", &cfg).await.unwrap();
5975 kv.put(
5977 &keys::domain("acme.example"),
5978 crate::project::DomainOwner::new("acme", "www").to_bytes(),
5979 )
5980 .await
5981 .unwrap();
5982
5983 store
5985 .put_function(
5986 acme,
5987 &Function::new(
5988 "worker",
5989 Owner::Project("acme".into()),
5990 "component-hash",
5991 FunctionConfig::default(),
5992 Lifecycle::default(),
5993 0,
5994 ),
5995 )
5996 .await
5997 .unwrap();
5998
5999 let spec = ComputeSpec {
6001 version: crate::SCHEMA_VERSION,
6002 root: RootSource::Rootfs("r".repeat(64)),
6003 kernel: "k".repeat(64),
6004 kernel_cmdline: None,
6005 vcpus: 1,
6006 mem_mib: 256,
6007 entrypoint: vec!["/app".into()],
6008 env: Default::default(),
6009 port: 8080,
6010 restart: RestartPolicy::Always,
6011 startup_grace_secs: 30,
6012 scale_to_zero: false,
6013 volumes: vec![VolumeRef {
6014 mount: "/data".into(),
6015 name: "pg-data".into(),
6016 size_mib: 1024,
6017 }],
6018 writable_root: false,
6019 cap_add: Vec::new(),
6020 user: None,
6021 isolation: Default::default(),
6022 prefer_backend: None,
6023 bindings: vec![],
6024 };
6025 let hash = store.put_compute_spec(&spec).await.unwrap();
6026 store
6027 .set_compute_workload(
6028 acme,
6029 &ComputeWorkload {
6030 version: crate::SCHEMA_VERSION,
6031 name: "pg".into(),
6032 active: hash,
6033 replicas: 1,
6034 placement: PlacementConstraints::default(),
6035 },
6036 )
6037 .await
6038 .unwrap();
6039
6040 kv.put(&keys::secret(acme, "api-key"), b"sealed".to_vec())
6042 .await
6043 .unwrap();
6044 kv.put(
6046 &format!("{}deadbeef", keys::graphql_safelist_prefix(acme)),
6047 b"query{me}".to_vec(),
6048 )
6049 .await
6050 .unwrap();
6051 kv.put(
6052 &format!("{}users", keys::graphql_subgraph_prefix(acme)),
6053 b"type Query{me:ID}".to_vec(),
6054 )
6055 .await
6056 .unwrap();
6057 kv.put("graphql/acme/version", 3u64.to_be_bytes().to_vec())
6058 .await
6059 .unwrap();
6060 kv.put(
6062 &crate::project::owner_key(crate::project::owner_kind::COMPUTE, "pg"),
6063 b"acme".to_vec(),
6064 )
6065 .await
6066 .unwrap();
6067 kv.put(
6068 &crate::project::owner_key(crate::project::owner_kind::SITE, "elsewhere"),
6069 b"other".to_vec(),
6070 )
6071 .await
6072 .unwrap();
6073
6074 store
6076 .set_site_config(other, "shop", &SiteConfig::default())
6077 .await
6078 .unwrap();
6079 kv.put("siteconfig/shared-cas", b"body".to_vec())
6080 .await
6081 .unwrap();
6082
6083 let plan = store.enumerate_project_resources("acme").await.unwrap();
6085 assert_eq!(plan.project, "acme");
6086 assert_eq!(plan.sites, vec!["www".to_string()]);
6087 assert_eq!(plan.functions, vec!["worker".to_string()]);
6088 assert_eq!(plan.compute.len(), 1);
6089 assert_eq!(plan.compute[0].name, "pg");
6090 assert_eq!(plan.compute[0].volumes, vec!["pg-data".to_string()]);
6091 assert_eq!(plan.all_volumes(), vec!["pg-data".to_string()]);
6092 assert_eq!(plan.secrets, vec!["api-key".to_string()]);
6093 assert_eq!(plan.safelist, 1);
6094 assert_eq!(plan.subgraphs, vec!["users".to_string()]);
6095 assert!(!plan.is_empty());
6096
6097 assert!(store.get_project("acme").await.unwrap().is_some());
6099 assert!(store.get_site_config(acme, "www").await.unwrap().is_some());
6100
6101 store.delete_site(acme, "www").await.unwrap();
6103 store.delete_function(acme, "worker").await.unwrap();
6104 store.delete_compute_workload(acme, "pg").await.unwrap();
6105 let purged = store.purge_project("acme").await.unwrap();
6106 assert!(purged > 0, "purge removed residual keys");
6107
6108 assert!(store.get_project("acme").await.unwrap().is_none());
6110 assert!(kv
6111 .list_prefix(&crate::project::resource_prefix("acme"))
6112 .await
6113 .unwrap()
6114 .is_empty());
6115 assert!(kv.list_prefix("graphql/acme/").await.unwrap().is_empty());
6116 assert!(kv
6117 .list_prefix(&keys::graphql_safelist_prefix(acme))
6118 .await
6119 .unwrap()
6120 .is_empty());
6121 assert!(kv
6122 .get(&crate::project::pointer_key("acme"))
6123 .await
6124 .unwrap()
6125 .is_none());
6126 assert!(kv
6127 .get(&crate::project::history_key("acme"))
6128 .await
6129 .unwrap()
6130 .is_none());
6131 assert!(kv
6132 .get(&crate::project::owner_key(
6133 crate::project::owner_kind::COMPUTE,
6134 "pg"
6135 ))
6136 .await
6137 .unwrap()
6138 .is_none());
6139 assert!(kv
6141 .get(&keys::domain("acme.example"))
6142 .await
6143 .unwrap()
6144 .is_none());
6145
6146 assert!(store
6148 .get_site_config(other, "shop")
6149 .await
6150 .unwrap()
6151 .is_some());
6152 assert!(!kv
6153 .list_prefix(&crate::project::resource_prefix("other"))
6154 .await
6155 .unwrap()
6156 .is_empty());
6157 assert_eq!(
6158 kv.get(&crate::project::owner_key(
6159 crate::project::owner_kind::SITE,
6160 "elsewhere"
6161 ))
6162 .await
6163 .unwrap(),
6164 Some(b"other".to_vec())
6165 );
6166 assert!(kv.get("siteconfig/shared-cas").await.unwrap().is_some());
6167
6168 assert!(matches!(
6170 store.purge_project("default").await,
6171 Err(DeployError::Conflict(_))
6172 ));
6173 }
6174
6175 #[tokio::test]
6182 async fn read_backfills_project_into_a_legacy_replica_record() {
6183 use crate::compute::{replica_state_key, Endpoint, ReplicaPhase, Scheme, Snapshot};
6184 use crate::kv::MemoryKv;
6185 let kv = Arc::new(MemoryKv::new());
6186 let store = DeployStore::new(Arc::new(NullStorage), kv.clone());
6187
6188 let legacy = serde_json::json!({
6192 "handle": { "workload": "web", "replica": 0, "backend_ref": "10.0.0.5:8080" },
6193 "node": 1,
6194 "backend": "container",
6195 "endpoint": { "scheme": "http", "host": "10.0.0.5", "port": 8080 },
6196 "healthy": false,
6197 "phase": "Zero",
6198 "snapshot": { "workload": "web", "replica": 0, "data_ref": "img|/rootfs|10.0.0.5|8080" }
6199 });
6200 assert!(legacy["handle"].get("project").is_none());
6202 kv.put(
6203 &replica_state_key("acme", "web", 0),
6204 serde_json::to_vec(&legacy).unwrap(),
6205 )
6206 .await
6207 .unwrap();
6208
6209 let states = store
6211 .list_replica_states(ProjectRef::new("acme"), "web")
6212 .await
6213 .unwrap();
6214 assert_eq!(states.len(), 1);
6215 assert_eq!(
6216 states[0].handle.project, "acme",
6217 "handle project backfilled"
6218 );
6219 assert_eq!(states[0].handle.workload, "web");
6220 assert_eq!(
6221 states[0].snapshot.as_ref().unwrap().project,
6222 "acme",
6223 "parked snapshot's project backfilled too"
6224 );
6225 assert_eq!(states[0].phase, ReplicaPhase::Zero);
6226
6227 let all = store.list_all_replica_states().await.unwrap();
6229 assert_eq!(all.len(), 1);
6230 assert_eq!(all[0].handle.project, "acme");
6231
6232 store
6235 .set_replica_state(
6236 ProjectRef::new("beta"),
6237 &crate::compute::ObservedInstance {
6238 handle: crate::compute::InstanceHandle {
6239 project: "beta".into(),
6240 workload: "web".into(),
6241 replica: 0,
6242 backend_ref: "10.0.0.6:8080".into(),
6243 },
6244 node: 1,
6245 backend: "container".into(),
6246 endpoint: Endpoint {
6247 scheme: Scheme::Http,
6248 host: "10.0.0.6".into(),
6249 port: 8080,
6250 },
6251 region: None,
6252 healthy: true,
6253 started_at: None,
6254 phase: ReplicaPhase::Running,
6255 snapshot: None::<Snapshot>,
6256 },
6257 )
6258 .await
6259 .unwrap();
6260 let beta = store
6261 .list_replica_states(ProjectRef::new("beta"), "web")
6262 .await
6263 .unwrap();
6264 assert_eq!(beta[0].handle.project, "beta");
6265 }
6266}