#![cfg(any(feature = "sql-postgres", feature = "sql-mysql"))]
use std::collections::BTreeMap;
use std::net::Ipv4Addr;
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use boatramp_container::dns::{Decision, ResolvedAddrs, Resolver, DEFAULT_INTERNAL_DOMAIN};
use boatramp_core::compute::{
reconcile_once, AlwaysActive, Artifact, BackendError, BackendPolicy, Capabilities,
ComputeBackend, ComputeSpec, Endpoint, Health, Instance, InstanceHandle, IsolationClass,
LaunchRequest, Node, ObservedInstance, ReplicaPhase, Scheme, Snapshot,
};
use boatramp_core::deploy::DeployStore;
use boatramp_core::envelope::{EnvelopeError, KeyEnvelope};
use boatramp_core::ipam::IpAuthority;
use boatramp_core::kv::{KvStore, MemoryKv};
use boatramp_core::project::ProjectRef;
use boatramp_core::{ByteStream, GetObject, ObjectMeta, PutMeta, Storage, StorageError};
use boatramp_node::compute::adopt_running_replica_ips;
use boatramp_node::config::{ExternalDatabaseConfig, TenantIsolation, TenantScope};
use boatramp_node::managed_sql::{
auto_register_managed_db_workloads, DeployEndpointResolver, ManagedSqlCredentials,
};
use boatramp_node::tenant_sql::provision_tenant;
use boatramp_storage::sql_compute::ComputeEndpointResolver;
use boatramp_storage::tenant_provision::sanitize_ident;
const COMPUTE: &str = "pg";
const SUBNET: &str = "10.0.0.0/24";
const NODE_ID: u64 = 1;
type ReplicaKey = (String, String, u32);
#[derive(Debug, Clone, PartialEq, Eq)]
struct FakeInstance {
container_id: String,
ip: Ipv4Addr,
port: u16,
}
struct FakeComputeBackend {
authority: IpAuthority,
scheme: Scheme,
instances: Mutex<BTreeMap<ReplicaKey, FakeInstance>>,
not_ready_polls: Mutex<u32>,
}
impl FakeComputeBackend {
fn new(authority: IpAuthority) -> Self {
Self {
authority,
scheme: Scheme::Http,
instances: Mutex::new(BTreeMap::new()),
not_ready_polls: Mutex::new(0),
}
}
fn instances(&self) -> Vec<(String, Ipv4Addr)> {
self.instances
.lock()
.expect("instances")
.values()
.map(|i| (i.container_id.clone(), i.ip))
.collect()
}
fn ip_of(&self, project: &str, workload: &str, replica: u32) -> Option<Ipv4Addr> {
self.instances
.lock()
.expect("instances")
.get(&(project.to_string(), workload.to_string(), replica))
.map(|i| i.ip)
}
fn container_id_of(&self, project: &str, workload: &str, replica: u32) -> Option<String> {
self.instances
.lock()
.expect("instances")
.get(&(project.to_string(), workload.to_string(), replica))
.map(|i| i.container_id.clone())
}
fn release_if_last(&self, instances: &BTreeMap<ReplicaKey, FakeInstance>, ip: Ipv4Addr) {
if !instances.values().any(|i| i.ip == ip) {
self.authority.release(ip);
}
}
}
#[async_trait]
impl ComputeBackend for FakeComputeBackend {
fn id(&self) -> &'static str {
"container"
}
fn capabilities(&self) -> Capabilities {
Capabilities {
isolation: IsolationClass::Namespace,
scale_to_zero: true,
persistent_volumes: true,
max_vcpus: None,
max_mem_mib: None,
}
}
async fn materialize(&self, spec: &ComputeSpec) -> Result<Artifact, BackendError> {
Ok(Artifact::Image {
reference: format!("fake:{}", spec.port),
})
}
async fn reserve_in_use(&self, replicas: &[(String, String, u32, Ipv4Addr)]) {
let ips: Vec<Ipv4Addr> = replicas.iter().map(|(_, _, _, ip)| *ip).collect();
self.authority.reserve_in_use(&ips);
let mut instances = self.instances.lock().expect("instances");
for (p, w, r, ip) in replicas {
if self.authority.manages(*ip) {
instances.insert(
(p.clone(), w.clone(), *r),
FakeInstance {
container_id: boatramp_core::compute::compute_instance_id(p, w, *r),
ip: *ip,
port: 0,
},
);
}
}
}
async fn gc_ip_pool(&self, parked: &[(String, String, u32)]) {
let _parked: std::collections::BTreeSet<ReplicaKey> = parked.iter().cloned().collect();
}
async fn launch(&self, req: &LaunchRequest) -> Result<Instance, BackendError> {
let key = (req.project.clone(), req.workload.clone(), req.replica);
let mut instances = self.instances.lock().expect("instances");
let recorded = instances.get(&key).map(|i| i.ip);
let owns_recorded =
recorded.is_some_and(|ip| !instances.iter().any(|(k, i)| k != &key && i.ip == ip));
let ip = if let Some(ip) = recorded.filter(|_| owns_recorded) {
self.authority.reserve(ip); ip
} else {
self.authority
.allocate_stable(recorded)
.map_err(|e| BackendError::Launch(e.to_string()))?
};
let port = req.spec.port;
instances.insert(
key,
FakeInstance {
container_id: boatramp_core::compute::compute_instance_id(
&req.project,
&req.workload,
req.replica,
),
ip,
port,
},
);
Ok(Instance {
handle: InstanceHandle {
project: req.project.clone(),
workload: req.workload.clone(),
replica: req.replica,
backend_ref: format!("{ip}:{port}"),
},
endpoint: Endpoint {
scheme: self.scheme,
host: ip.to_string(),
port,
},
})
}
async fn stop(&self, handle: &InstanceHandle) -> Result<(), BackendError> {
let key = (
handle.project.clone(),
handle.workload.clone(),
handle.replica,
);
let mut instances = self.instances.lock().expect("instances");
let ip = instances.remove(&key).map(|i| i.ip).or_else(|| {
handle
.backend_ref
.split(':')
.next()
.and_then(|s| s.parse().ok())
});
if let Some(ip) = ip {
self.release_if_last(&instances, ip);
}
Ok(())
}
async fn health(&self, _handle: &InstanceHandle) -> Result<Health, BackendError> {
let mut left = self.not_ready_polls.lock().expect("not_ready_polls");
if *left > 0 {
*left -= 1;
Ok(Health::Unhealthy)
} else {
Ok(Health::Healthy)
}
}
async fn snapshot(&self, handle: &InstanceHandle) -> Result<Option<Snapshot>, BackendError> {
let key = (
handle.project.clone(),
handle.workload.clone(),
handle.replica,
);
let instances = self.instances.lock().expect("instances");
let Some(inst) = instances.get(&key) else {
return Ok(None);
};
Ok(Some(Snapshot {
project: handle.project.clone(),
workload: handle.workload.clone(),
replica: handle.replica,
data_ref: format!("snap|{}|{}", inst.ip, inst.port),
}))
}
async fn restore(&self, snapshot: &Snapshot) -> Result<Instance, BackendError> {
let mut parts = snapshot.data_ref.rsplitn(3, '|');
let port: u16 = parts.next().and_then(|s| s.parse().ok()).unwrap_or(0);
let ip: Ipv4Addr = parts
.next()
.and_then(|s| s.parse().ok())
.ok_or_else(|| BackendError::Other("bad snapshot ref".into()))?;
self.authority.reserve(ip);
let mut instances = self.instances.lock().expect("instances");
instances.insert(
(
snapshot.project.clone(),
snapshot.workload.clone(),
snapshot.replica,
),
FakeInstance {
container_id: boatramp_core::compute::compute_instance_id(
&snapshot.project,
&snapshot.workload,
snapshot.replica,
),
ip,
port,
},
);
Ok(Instance {
handle: InstanceHandle {
project: snapshot.project.clone(),
workload: snapshot.workload.clone(),
replica: snapshot.replica,
backend_ref: format!("{ip}:{port}"),
},
endpoint: Endpoint {
scheme: self.scheme,
host: ip.to_string(),
port,
},
})
}
}
struct NullStorage;
#[async_trait]
impl Storage for NullStorage {
async fn get(&self, _: &str) -> Result<GetObject, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn get_range(&self, _: &str, _: u64, _: Option<u64>) -> Result<GetObject, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn put(&self, _: &str, _: ByteStream, _: PutMeta) -> Result<ObjectMeta, StorageError> {
Err(StorageError::unsupported("null"))
}
async fn head(&self, _: &str) -> Result<ObjectMeta, StorageError> {
Err(StorageError::NotFound(String::new()))
}
async fn delete(&self, _: &str) -> Result<(), StorageError> {
Ok(())
}
async fn list(&self, _: &str) -> Result<Vec<ObjectMeta>, StorageError> {
Ok(Vec::new())
}
}
struct RevEnvelope;
#[async_trait]
impl KeyEnvelope for RevEnvelope {
async fn wrap(&self, p: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
Ok(p.iter().rev().copied().collect())
}
async fn unwrap(&self, w: &[u8]) -> Result<Vec<u8>, EnvelopeError> {
Ok(w.iter().rev().copied().collect())
}
}
fn single_project_binding() -> ExternalDatabaseConfig {
ExternalDatabaseConfig {
kind: "postgres".into(),
compute: Some(COMPUTE.into()),
database: Some("appdb".into()),
user: Some("app".into()),
tenant: TenantIsolation::Single,
tenant_scope: TenantScope::Project,
..Default::default()
}
}
fn databases() -> BTreeMap<String, ExternalDatabaseConfig> {
BTreeMap::from([("main".to_string(), single_project_binding())])
}
fn tenant_workload_name(project: &str) -> String {
format!("{COMPUTE}-{}", sanitize_ident(project))
}
fn nodes() -> Vec<Node> {
vec![Node {
id: NODE_ID,
region: Some("eu".into()),
labels: BTreeMap::new(),
free_vcpus: 32,
free_mem_mib: 65536,
backends: vec![boatramp_core::compute::BackendKind {
id: "container".into(),
isolation: IsolationClass::Namespace,
persistent_volumes: true,
scale_to_zero: true,
}],
}]
}
fn control_plane() -> (DeployStore, Arc<dyn KvStore>, IpAuthority) {
let kv: Arc<dyn KvStore> = Arc::new(MemoryKv::new());
let deploy = DeployStore::new(Arc::new(NullStorage), kv.clone());
let authority = IpAuthority::new(SUBNET).expect("valid subnet");
(deploy, kv, authority)
}
fn registry(backend: Arc<FakeComputeBackend>) -> boatramp_core::compute::BackendRegistry {
let mut reg: boatramp_core::compute::BackendRegistry = BTreeMap::new();
reg.insert("container".into(), backend as Arc<dyn ComputeBackend>);
reg
}
async fn seed_project(deploy: &DeployStore, name: &str) {
deploy
.put_project(&boatramp_core::project::Project {
version: 1,
name: name.to_string(),
created_at: 0,
meta: Default::default(),
config: Default::default(),
secrets_ref: None,
})
.await
.expect("seed project pointer");
}
async fn seed_static_site(kv: &Arc<dyn KvStore>, project: &str, site: &str) {
kv.put(
&format!("project/{project}/current/{site}"),
b"deadbeef".to_vec(),
)
.await
.expect("seed static-site pointer");
}
async fn reconcile(
deploy: &DeployStore,
reg: &boatramp_core::compute::BackendRegistry,
) -> boatramp_core::compute::ReconcileReport {
reconcile_once(
deploy,
reg,
&nodes(),
&BackendPolicy::default(),
&AlwaysActive,
None,
None,
)
.await
.expect("reconcile pass")
}
async fn all_states(deploy: &DeployStore) -> Vec<ObservedInstance> {
deploy
.list_all_replica_states()
.await
.expect("list all replica states")
}
#[tokio::test]
async fn single_binding_does_not_overwarm_a_static_only_default() {
let (deploy, kv, _authority) = control_plane();
let kv_arc = kv.clone();
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
seed_project(&deploy, "default").await;
seed_static_site(&kv, "default", "www").await;
seed_project(&deploy, "acme").await;
auto_register_managed_db_workloads(&deploy, &databases()).await;
let default_after_boot = deploy
.list_compute_workloads(ProjectRef::DEFAULT)
.await
.expect("list default workloads");
assert!(
default_after_boot.is_empty(),
"a Single binding must NOT auto-register any workload under a static-only default; \
found {default_after_boot:?}"
);
provision_tenant(
&deploy,
&kv_arc,
&envelope,
&single_project_binding(),
"acme",
"",
)
.await
.expect("provision acme");
let default_final = deploy
.list_compute_workloads(ProjectRef::DEFAULT)
.await
.expect("list default workloads");
assert!(
default_final.is_empty(),
"no bare `pg` (or any managed workload) may exist under default; found {default_final:?}"
);
let acme = deploy
.list_compute_workloads(ProjectRef::new("acme"))
.await
.expect("list acme workloads");
assert_eq!(
acme.len(),
1,
"acme must have exactly one managed workload, found {acme:?}"
);
assert_eq!(
acme[0].name,
tenant_workload_name("acme"),
"acme's managed workload must be the tenant-aware `pg-<ident>`"
);
}
#[tokio::test]
async fn reboot_adoption_keeps_ips_unique_and_stable() {
let (deploy, kv, authority) = control_plane();
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
for project in ["acme", "globex"] {
seed_project(&deploy, project).await;
provision_tenant(
&deploy,
&kv,
&envelope,
&single_project_binding(),
project,
"",
)
.await
.unwrap_or_else(|e| panic!("provision {project}: {e}"));
}
let backend = Arc::new(FakeComputeBackend::new(authority.clone()));
let reg = registry(backend.clone());
let report = reconcile(&deploy, ®).await;
assert_eq!(
report.launched, 2,
"both tenants launch: {:?}",
report.errors
);
let acme_wl = tenant_workload_name("acme");
let globex_wl = tenant_workload_name("globex");
let acme_ip = backend
.ip_of("acme", &acme_wl, 0)
.expect("acme replica launched");
let globex_ip = backend
.ip_of("globex", &globex_wl, 0)
.expect("globex replica launched");
assert_ne!(
acme_ip, globex_ip,
"two projects' managed DBs must get distinct IPs (no 10.0.0.2-twice collision)"
);
let rebooted_authority = IpAuthority::new(SUBNET).expect("subnet");
let rebooted = Arc::new(FakeComputeBackend::new(rebooted_authority));
let reg2 = registry(rebooted.clone());
adopt_running_replica_ips(&deploy, ®2).await;
let next_fresh = rebooted
.authority
.allocate()
.expect("pool still has room after adopting two addresses");
assert!(
next_fresh != acme_ip && next_fresh != globex_ip,
"a fresh allocation after adoption re-handed a LIVE address ({next_fresh}) — the \
10.0.0.2-twice collision; adoption must reserve {acme_ip} and {globex_ip}"
);
rebooted.authority.release(next_fresh);
let report2 = reconcile(&deploy, ®2).await;
assert_eq!(
report2.launched, 0,
"a boot reconcile of already-running replicas relaunches nothing: {:?}",
report2.errors
);
let acme_after = rebooted.ip_of("acme", &acme_wl, 0).expect("acme adopted");
let globex_after = rebooted
.ip_of("globex", &globex_wl, 0)
.expect("globex adopted");
assert_eq!(acme_after, acme_ip, "acme keeps its IP across the reboot");
assert_eq!(
globex_after, globex_ip,
"globex keeps its IP across the reboot"
);
assert_ne!(
acme_after, globex_after,
"the two replicas' IPs remain distinct after the reboot"
);
let acme_states = deploy
.list_replica_states(ProjectRef::new("acme"), &acme_wl)
.await
.unwrap();
assert_eq!(acme_states[0].endpoint.host, acme_ip.to_string());
}
#[tokio::test]
async fn two_projects_same_workload_name_get_distinct_ids_and_ips() {
let (deploy, _kv, authority) = control_plane();
let spec = plain_web_spec();
let spec_id = deploy.put_compute_spec(&spec).await.expect("put spec");
for project in ["acme", "globex"] {
seed_project(&deploy, project).await;
deploy
.set_compute_workload(
ProjectRef::new(project),
&boatramp_core::compute::ComputeWorkload {
version: 1,
name: "web".into(),
active: spec_id.clone(),
replicas: 1,
placement: Default::default(),
},
)
.await
.expect("register web");
}
let backend = Arc::new(FakeComputeBackend::new(authority));
let reg = registry(backend.clone());
let report = reconcile(&deploy, ®).await;
assert_eq!(
report.launched, 2,
"both `web`s launch: {:?}",
report.errors
);
let live = backend.instances();
assert_eq!(live.len(), 2, "two replicas launched: {live:?}");
let acme_id = backend.container_id_of("acme", "web", 0).expect("acme web");
let globex_id = backend
.container_id_of("globex", "web", 0)
.expect("globex web");
assert_ne!(
acme_id, globex_id,
"same-named workloads in two projects must derive DISTINCT identities"
);
assert_eq!(acme_id, "acme-web-0");
assert_eq!(globex_id, "globex-web-0");
let acme_ip = backend.ip_of("acme", "web", 0).unwrap();
let globex_ip = backend.ip_of("globex", "web", 0).unwrap();
assert_ne!(
acme_ip, globex_ip,
"same-named workloads in two projects must get DISTINCT IPs"
);
let ids: std::collections::BTreeSet<_> = live.into_iter().map(|(id, _)| id).collect();
assert!(ids.contains("acme-web-0") && ids.contains("globex-web-0"));
}
fn plain_web_spec() -> ComputeSpec {
ComputeSpec {
version: 1,
root: boatramp_core::compute::RootSource::Image("web:latest".into()),
kernel: String::new(),
kernel_cmdline: None,
vcpus: 1,
mem_mib: 128,
entrypoint: vec![],
env: BTreeMap::new(),
port: 8080,
restart: boatramp_core::compute::RestartPolicy::Always,
startup_grace_secs: 30,
scale_to_zero: false,
volumes: vec![],
writable_root: false,
cap_add: Vec::new(),
user: None,
isolation: boatramp_core::compute::IsolationRequirement::Trusted,
prefer_backend: None,
bindings: vec![],
}
}
#[tokio::test]
async fn operator_sql_targets_the_tenant_not_default() {
let (deploy, kv, authority) = control_plane();
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
seed_project(&deploy, "acme").await;
provision_tenant(
&deploy,
&kv,
&envelope,
&single_project_binding(),
"acme",
"",
)
.await
.expect("provision acme");
let backend = Arc::new(FakeComputeBackend::new(authority));
let reg = registry(backend);
reconcile(&deploy, ®).await;
let acme_wl = tenant_workload_name("acme");
let acme_resolver = DeployEndpointResolver::new(deploy.clone(), "acme");
let acme_eps = acme_resolver
.endpoints(&acme_wl)
.await
.expect("acme endpoints");
assert_eq!(
acme_eps.len(),
1,
"operator sql for acme must resolve acme's own `pg-<ident>` replica"
);
let default_resolver = DeployEndpointResolver::new(deploy.clone(), "default");
assert!(
default_resolver
.endpoints(COMPUTE)
.await
.expect("default endpoints")
.is_empty(),
"operator sql must NOT find a bare `pg` under default (the routing bug)"
);
let creds = ManagedSqlCredentials::new(kv.clone(), envelope.clone());
let tenant_pw = creds
.password("acme", &acme_wl)
.await
.expect("tenant credential exists");
assert!(
!tenant_pw.is_empty(),
"the tenant credential is materialized"
);
let default_workloads = deploy
.list_compute_workloads(ProjectRef::DEFAULT)
.await
.unwrap();
assert!(
!default_workloads.iter().any(|w| w.name == COMPUTE),
"no bare `{COMPUTE}` workload may exist under default"
);
}
#[tokio::test]
async fn compute_upstream_resolves_per_project() {
let (deploy, _kv, authority) = control_plane();
let spec = plain_web_spec();
let spec_id = deploy.put_compute_spec(&spec).await.unwrap();
for project in ["acme", "default"] {
seed_project(&deploy, project).await;
deploy
.set_compute_workload(
ProjectRef::new(project),
&boatramp_core::compute::ComputeWorkload {
version: 1,
name: "api".into(),
active: spec_id.clone(),
replicas: 1,
placement: Default::default(),
},
)
.await
.unwrap();
}
let backend = Arc::new(FakeComputeBackend::new(authority));
let reg = registry(backend.clone());
let report = reconcile(&deploy, ®).await;
assert_eq!(
report.launched, 2,
"both `api`s launch: {:?}",
report.errors
);
let acme_ip = backend.ip_of("acme", "api", 0).unwrap();
let default_ip = backend.ip_of("default", "api", 0).unwrap();
let acme = DeployEndpointResolver::new(deploy.clone(), "acme");
let acme_eps = acme.endpoints("api").await.unwrap();
assert_eq!(acme_eps, vec![(acme_ip.to_string(), 8080)]);
assert!(
!acme_eps.iter().any(|(h, _)| *h == default_ip.to_string()),
"acme's compute upstream must not see default's replica"
);
let default = DeployEndpointResolver::new(deploy.clone(), "default");
let default_eps = default.endpoints("api").await.unwrap();
assert_eq!(default_eps, vec![(default_ip.to_string(), 8080)]);
assert!(
!default_eps.iter().any(|(h, _)| *h == acme_ip.to_string()),
"default's compute upstream must not see acme's replica (the hardcoded-DEFAULT bug)"
);
let beta = DeployEndpointResolver::new(deploy.clone(), "beta");
assert!(beta.endpoints("api").await.unwrap().is_empty());
}
#[tokio::test]
async fn internal_dns_refuses_cross_project() {
let (deploy, _kv, authority) = control_plane();
let spec = plain_web_spec();
let spec_id = deploy.put_compute_spec(&spec).await.unwrap();
for project in ["acme", "globex"] {
seed_project(&deploy, project).await;
deploy
.set_compute_workload(
ProjectRef::new(project),
&boatramp_core::compute::ComputeWorkload {
version: 1,
name: "web".into(),
active: spec_id.clone(),
replicas: 1,
placement: Default::default(),
},
)
.await
.unwrap();
}
let backend = Arc::new(FakeComputeBackend::new(authority));
let reg = registry(backend);
reconcile(&deploy, ®).await;
let states = all_states(&deploy).await;
let mut owners: BTreeMap<Ipv4Addr, (String, String)> = BTreeMap::new();
let mut addrs: BTreeMap<(String, String), Vec<Ipv4Addr>> = BTreeMap::new();
for st in &states {
let ip: Ipv4Addr = st.endpoint.host.parse().expect("endpoint host is an IPv4");
let key = (st.handle.project.clone(), st.handle.workload.clone());
owners.insert(ip, key.clone());
if st.phase == ReplicaPhase::Running && st.healthy {
addrs.entry(key).or_default().push(ip);
}
}
let acme_ip = *owners
.iter()
.find(|(_, (p, _))| p == "acme")
.map(|(ip, _)| ip)
.expect("acme replica IP");
let globex_ip = *owners
.iter()
.find(|(_, (p, _))| p == "globex")
.map(|(ip, _)| ip)
.expect("globex replica IP");
let resolver = Resolver::new(
DEFAULT_INTERNAL_DOMAIN,
move |ip| owners.get(&ip).cloned(),
move |p, w| ResolvedAddrs {
v4: addrs
.get(&(p.to_string(), w.to_string()))
.cloned()
.unwrap_or_default(),
v6: Vec::new(),
},
);
let q = encode_query(0x1, "web", 1, 1);
let Decision::Reply(reply) = resolver.handle_query(acme_ip, &q) else {
panic!("acme must resolve its own workload, got forward");
};
assert_eq!(rcode_of(&reply), 0, "NOERROR for an in-project name");
assert_eq!(
a_records(&reply),
vec![acme_ip],
"the bare name resolves within the caller's own project only"
);
let q = encode_query(0x2, "web.globex.boatramp.internal", 1, 1);
let Decision::Reply(reply) = resolver.handle_query(acme_ip, &q) else {
panic!("a cross-project internal name must NOT be forwarded");
};
assert_eq!(
rcode_of(&reply),
5,
"cross-project resolution must be REFUSED"
);
assert!(
a_records(&reply).is_empty() && !a_records(&reply).contains(&globex_ip),
"globex's address must never leak into acme's answer"
);
let q = encode_query(0x3, "web.acme.boatramp.internal", 1, 1);
assert_eq!(
resolver.handle_query(Ipv4Addr::new(203, 0, 113, 7), &q),
Decision::Forward,
"an unknown source must be forwarded, not answered an internal name"
);
}
#[tokio::test]
async fn health_recovery_is_persisted_so_the_resolver_sees_it() {
let (deploy, kv, authority) = control_plane();
let envelope: Arc<dyn KeyEnvelope> = Arc::new(RevEnvelope);
seed_project(&deploy, "acme").await;
provision_tenant(
&deploy,
&kv,
&envelope,
&single_project_binding(),
"acme",
"",
)
.await
.expect("provision acme");
let wl = tenant_workload_name("acme");
let backend = Arc::new(FakeComputeBackend::new(authority));
*backend.not_ready_polls.lock().expect("not_ready_polls") = 1;
let reg = registry(backend.clone());
let resolver = DeployEndpointResolver::new(deploy.clone(), "acme");
reconcile(&deploy, ®).await;
assert!(
resolver.endpoints(&wl).await.unwrap().is_empty(),
"a replica whose launch-time probe failed must not resolve as healthy yet"
);
reconcile(&deploy, ®).await;
assert!(
!resolver.endpoints(&wl).await.unwrap().is_empty(),
"recovered health was not persisted — the resolver still reports no healthy \
replica for a reachable DB (the v0.3.12 regression this locks)"
);
}
fn encode_query(id: u16, name: &str, qtype: u16, qclass: u16) -> Vec<u8> {
let mut buf = Vec::new();
buf.extend_from_slice(&id.to_be_bytes());
buf.extend_from_slice(&0x0100u16.to_be_bytes()); buf.extend_from_slice(&1u16.to_be_bytes()); buf.extend_from_slice(&[0, 0, 0, 0, 0, 0]); for label in name.trim_end_matches('.').split('.') {
buf.push(label.len() as u8);
buf.extend_from_slice(label.as_bytes());
}
buf.push(0); buf.extend_from_slice(&qtype.to_be_bytes());
buf.extend_from_slice(&qclass.to_be_bytes());
buf
}
fn rcode_of(reply: &[u8]) -> u8 {
reply[3] & 0x0F
}
fn ancount_of(reply: &[u8]) -> u16 {
u16::from_be_bytes([reply[6], reply[7]])
}
fn a_records(reply: &[u8]) -> Vec<Ipv4Addr> {
let mut pos = 12;
while reply[pos] != 0 {
pos += 1 + reply[pos] as usize;
}
pos += 1 + 4; let mut out = Vec::new();
let mut i = 0;
while i < ancount_of(reply) {
let rtype = u16::from_be_bytes([reply[pos + 2], reply[pos + 3]]);
let rdlen = u16::from_be_bytes([reply[pos + 10], reply[pos + 11]]) as usize;
let rdata = &reply[pos + 12..pos + 12 + rdlen];
if rtype == 1 && rdlen == 4 {
out.push(Ipv4Addr::new(rdata[0], rdata[1], rdata[2], rdata[3]));
}
pos += 12 + rdlen;
i += 1;
}
out
}