use std::collections::{BTreeMap, HashMap, HashSet};
use std::future::Future;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use anyhow::{Context, Result, anyhow, bail};
use arc_swap::ArcSwap;
use parking_lot::Mutex;
use reqwest::Client;
use serde_json::Value;
use tokio::runtime::{Handle, Runtime};
use tokio::task::JoinHandle;
use crate::config::HostConfig;
use crate::engine::host::{SessionHost, StateHost};
use crate::engine::runtime::StateMachineRuntime;
use crate::oauth::{OAuthBrokerConfig, request_resource_token};
use crate::operator_metrics::OperatorMetrics;
use crate::operator_registry::OperatorRegistry;
use crate::pack::{ComponentResolution, PackRuntime};
use crate::runner::adapt_events_email::{
EmailExecutionPlan, EmailSendRequest, build_email_execution_plan, execute_email_request,
};
use crate::runner::contract_cache::{ContractCache, ContractCacheStats};
use crate::runner::engine::FlowEngine;
use crate::runner::mocks::MockLayer;
use crate::secrets::{DynSecretsManager, canonicalize_secret_key, read_secret_blocking};
use crate::storage::session::DynSessionStore;
use crate::storage::state::DynStateStore;
use crate::telemetry::RolloutIds;
use crate::trace::PackTraceInfo;
use crate::wasi::RunnerWasiPolicy;
use greentic_deploy_spec::ids::{BundleId, DeploymentId, RevisionId};
use greentic_types::SecretRequirement;
use runner_core::packs::PackDigest;
const RUNTIME_SECRETS_PACK_ID: &str = "_runner";
#[derive(Clone, Debug, PartialEq, Eq, Hash)]
pub struct RuntimeKey {
pub tenant: String,
pub deployment_id: Option<DeploymentId>,
pub bundle_id: Option<BundleId>,
pub revision_id: Option<RevisionId>,
}
impl RuntimeKey {
pub fn legacy(tenant: impl Into<String>) -> Self {
Self {
tenant: tenant.into(),
deployment_id: None,
bundle_id: None,
revision_id: None,
}
}
pub fn revision(
tenant: impl Into<String>,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
) -> Self {
Self {
tenant: tenant.into(),
deployment_id: Some(deployment_id),
bundle_id: Some(bundle_id),
revision_id: Some(revision_id),
}
}
pub fn is_legacy(&self) -> bool {
self.deployment_id.is_none() && self.bundle_id.is_none() && self.revision_id.is_none()
}
}
fn merge_legacy_reload<V: Clone>(
prev: &HashMap<RuntimeKey, V>,
mut legacy: HashMap<RuntimeKey, V>,
) -> HashMap<RuntimeKey, V> {
for (key, value) in prev {
if !key.is_legacy() {
legacy.insert(key.clone(), value.clone());
}
}
legacy
}
fn remove_keyed_entry<V: Clone>(
prev: &HashMap<RuntimeKey, V>,
key: &RuntimeKey,
) -> Option<(HashMap<RuntimeKey, V>, V)> {
let removed = prev.get(key)?.clone();
let mut next = prev.clone();
next.remove(key);
Some((next, removed))
}
fn insert_keyed_entry<V: Clone>(
prev: &HashMap<RuntimeKey, V>,
key: RuntimeKey,
value: V,
) -> HashMap<RuntimeKey, V> {
let mut next = prev.clone();
next.insert(key, value);
next
}
pub struct ActivePacks {
inner: ArcSwap<HashMap<RuntimeKey, Arc<TenantRuntime>>>,
write_lock: Mutex<()>,
}
impl ActivePacks {
pub fn new() -> Self {
Self {
inner: ArcSwap::from_pointee(HashMap::new()),
write_lock: Mutex::new(()),
}
}
pub fn load_pack(&self, tenant: &str) -> Option<Arc<TenantRuntime>> {
self.inner.load().get(&RuntimeKey::legacy(tenant)).cloned()
}
pub fn load_revision(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
) -> Option<Arc<TenantRuntime>> {
self.inner
.load()
.get(&RuntimeKey::revision(
tenant,
deployment_id,
bundle_id,
revision_id,
))
.cloned()
}
pub fn snapshot(&self) -> Arc<HashMap<RuntimeKey, Arc<TenantRuntime>>> {
self.inner.load_full()
}
pub fn insert_pack(&self, tenant: &str, runtime: Arc<TenantRuntime>) {
let _guard = self.write_lock.lock();
let mut next = (*self.inner.load_full()).clone();
next.insert(RuntimeKey::legacy(tenant), runtime);
self.inner.store(Arc::new(next));
}
pub fn insert_revision(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
runtime: Arc<TenantRuntime>,
) -> Result<()> {
if runtime.tenant() != tenant {
bail!(
"revision runtime tenant `{}` does not match key tenant `{tenant}`",
runtime.tenant()
);
}
let ids = runtime.engine().rollout_ids();
let key_deployment = deployment_id.to_string();
let key_bundle = bundle_id.as_str();
let key_revision = revision_id.to_string();
if ids.deployment_id.as_deref() != Some(key_deployment.as_str())
|| ids.bundle_id.as_deref() != Some(key_bundle)
|| ids.revision_id.as_deref() != Some(key_revision.as_str())
{
bail!(
"revision runtime rollout identity (deployment={:?}, bundle={:?}, revision={:?}) \
does not match key (deployment=`{key_deployment}`, bundle=`{key_bundle}`, \
revision=`{key_revision}`)",
ids.deployment_id,
ids.bundle_id,
ids.revision_id
);
}
let _guard = self.write_lock.lock();
let key = RuntimeKey::revision(tenant, deployment_id, bundle_id, revision_id);
let next = insert_keyed_entry(&self.inner.load_full(), key, runtime);
self.inner.store(Arc::new(next));
Ok(())
}
pub fn replace_legacy(&self, legacy: HashMap<RuntimeKey, Arc<TenantRuntime>>) {
let _guard = self.write_lock.lock();
let prev = self.inner.load_full();
let next = merge_legacy_reload(&prev, legacy);
self.inner.store(Arc::new(next));
}
pub fn remove_revision(
&self,
tenant: &str,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
) -> Option<Arc<TenantRuntime>> {
let _guard = self.write_lock.lock();
let prev = self.inner.load_full();
let key = RuntimeKey::revision(tenant, deployment_id, bundle_id, revision_id);
let (next, removed) = remove_keyed_entry(&prev, &key)?;
self.inner.store(Arc::new(next));
Some(removed)
}
pub fn replace(&self, next: HashMap<RuntimeKey, Arc<TenantRuntime>>) {
let _guard = self.write_lock.lock();
self.inner.store(Arc::new(next));
}
pub fn len(&self) -> usize {
self.inner.load().len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
}
impl Default for ActivePacks {
fn default() -> Self {
Self::new()
}
}
pub struct TenantRuntime {
tenant: String,
config: Arc<HostConfig>,
packs: Vec<Arc<PackRuntime>>,
digests: Vec<Option<String>>,
engine: Arc<FlowEngine>,
state_machine: Arc<StateMachineRuntime>,
session_store: DynSessionStore,
http_client: Client,
mocks: Option<Arc<MockLayer>>,
timer_handles: Mutex<Vec<JoinHandle<()>>>,
secrets: DynSecretsManager,
operator_registry: OperatorRegistry,
operator_metrics: Arc<OperatorMetrics>,
contract_cache: ContractCache,
}
#[derive(Clone)]
pub struct ResolvedComponent {
pub digest: String,
pub component_ref: String,
pub pack: Arc<PackRuntime>,
}
#[derive(Clone, Debug)]
pub struct RevisionPackRef {
pub path: PathBuf,
pub digest: String,
}
pub fn block_on<F: Future<Output = R>, R>(future: F) -> R {
if let Ok(handle) = Handle::try_current() {
handle.block_on(future)
} else {
Runtime::new()
.expect("failed to create tokio runtime")
.block_on(future)
}
}
impl TenantRuntime {
#[allow(clippy::too_many_arguments)]
pub async fn load(
pack_path: &Path,
config: Arc<HostConfig>,
mocks: Option<Arc<MockLayer>>,
archive_source: Option<&Path>,
digest: Option<String>,
wasi_policy: Arc<RunnerWasiPolicy>,
session_host: Arc<dyn SessionHost>,
session_store: DynSessionStore,
state_store: DynStateStore,
state_host: Arc<dyn StateHost>,
secrets_manager: DynSecretsManager,
) -> Result<Arc<Self>> {
let pack = Self::load_pack_runtime(
pack_path,
&config,
mocks.clone(),
archive_source,
&wasi_policy,
&session_store,
&state_store,
&secrets_manager,
&BTreeMap::new(),
&BTreeMap::new(),
None,
)
.await?;
Self::from_packs(
config,
vec![(pack, digest)],
mocks,
session_host,
session_store,
state_store,
state_host,
secrets_manager,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub async fn load_revision(
pack_refs: &[RevisionPackRef],
config: Arc<HostConfig>,
mocks: Option<Arc<MockLayer>>,
wasi_policy: Arc<RunnerWasiPolicy>,
session_host: Arc<dyn SessionHost>,
session_store: DynSessionStore,
state_store: DynStateStore,
state_host: Arc<dyn StateHost>,
secrets_manager: DynSecretsManager,
deployment_id: DeploymentId,
bundle_id: BundleId,
revision_id: RevisionId,
customer_id: Option<String>,
runtime_configs_by_pack_id: &BTreeMap<String, Arc<BTreeMap<String, Value>>>,
runtime_refs_by_pack_id: &BTreeMap<String, Arc<BTreeMap<String, String>>>,
runtime_ref_resolver: Option<Arc<dyn crate::runtime_refs::RuntimeRefResolver>>,
) -> Result<Arc<Self>> {
if pack_refs.is_empty() {
bail!(
"revision runtime for tenant {} requires at least one pack",
config.tenant
);
}
let mut packs = Vec::with_capacity(pack_refs.len());
let mut seen_pack_ids = HashSet::with_capacity(pack_refs.len());
for pack_ref in pack_refs {
let expected = PackDigest::parse(&pack_ref.digest).with_context(|| {
format!(
"revision pack `{}` has an invalid digest `{}`",
pack_ref.path.display(),
pack_ref.digest
)
})?;
if expected.algorithm() != "sha256" {
bail!(
"revision pack `{}` pins unsupported digest algorithm `{}`; only sha256 is supported",
pack_ref.path.display(),
expected.algorithm()
);
}
if !expected.matches_file(&pack_ref.path).with_context(|| {
format!(
"hashing revision pack `{}` for digest verification",
pack_ref.path.display()
)
})? {
bail!(
"revision pack `{}` does not match pinned digest `{}`",
pack_ref.path.display(),
pack_ref.digest
);
}
let pack = Self::load_pack_runtime(
&pack_ref.path,
&config,
mocks.clone(),
None,
&wasi_policy,
&session_store,
&state_store,
&secrets_manager,
runtime_configs_by_pack_id,
runtime_refs_by_pack_id,
runtime_ref_resolver.as_ref(),
)
.await?;
let pack_id = pack.metadata().pack_id.clone();
if !seen_pack_ids.insert(pack_id.clone()) {
bail!(
"revision for tenant {} contains duplicate pack_id `{}` (path `{}`)",
config.tenant,
pack_id,
pack_ref.path.display(),
);
}
packs.push((pack, Some(expected.raw_string())));
}
let rollout = RolloutIds {
customer_id,
deployment_id: Some(deployment_id.to_string()),
bundle_id: Some(bundle_id.as_str().to_string()),
revision_id: Some(revision_id.to_string()),
};
Self::from_packs_with_rollout(
config,
packs,
mocks,
session_host,
session_store,
state_store,
state_host,
secrets_manager,
rollout,
)
.await
}
#[allow(clippy::too_many_arguments)]
async fn load_pack_runtime(
pack_path: &Path,
config: &Arc<HostConfig>,
mocks: Option<Arc<MockLayer>>,
archive_source: Option<&Path>,
wasi_policy: &Arc<RunnerWasiPolicy>,
session_store: &DynSessionStore,
state_store: &DynStateStore,
secrets_manager: &DynSecretsManager,
runtime_configs_by_pack_id: &BTreeMap<String, Arc<BTreeMap<String, Value>>>,
runtime_refs_by_pack_id: &BTreeMap<String, Arc<BTreeMap<String, String>>>,
runtime_ref_resolver: Option<&Arc<dyn crate::runtime_refs::RuntimeRefResolver>>,
) -> Result<Arc<PackRuntime>> {
let oauth_config = config.oauth_broker_config();
let mut pack = PackRuntime::load(
pack_path,
Arc::clone(config),
mocks,
archive_source,
Some(Arc::clone(session_store)),
Some(Arc::clone(state_store)),
Arc::clone(wasi_policy),
Arc::clone(secrets_manager),
oauth_config,
true,
ComponentResolution::default(),
)
.await
.with_context(|| {
format!(
"failed to load pack {} for tenant {}",
pack_path.display(),
config.tenant
)
})?;
let pack_id = pack.metadata().pack_id.clone();
if let Some(non_secret) = runtime_configs_by_pack_id.get(pack_id.as_str()) {
pack.set_runtime_config_non_secret(Some(Arc::clone(non_secret)));
}
if let Some(refs) = runtime_refs_by_pack_id.get(pack_id.as_str()) {
let resolver = runtime_ref_resolver.ok_or_else(|| {
anyhow!(
"pack `{}` has runtime_refs bound but no RuntimeRefResolver was provided",
pack_id,
)
})?;
pack.set_runtime_refs(Some(crate::runtime_refs::RuntimeRefsInjection {
refs: Arc::clone(refs),
resolver: Arc::clone(resolver),
}));
}
Ok(Arc::new(pack))
}
#[allow(clippy::too_many_arguments)]
pub async fn from_packs(
config: Arc<HostConfig>,
packs: Vec<(Arc<PackRuntime>, Option<String>)>,
mocks: Option<Arc<MockLayer>>,
session_host: Arc<dyn SessionHost>,
session_store: DynSessionStore,
state_store: DynStateStore,
state_host: Arc<dyn StateHost>,
secrets_manager: DynSecretsManager,
) -> Result<Arc<Self>> {
Self::from_packs_with_rollout(
config,
packs,
mocks,
session_host,
session_store,
state_store,
state_host,
secrets_manager,
RolloutIds::default(),
)
.await
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn from_packs_with_rollout(
config: Arc<HostConfig>,
packs: Vec<(Arc<PackRuntime>, Option<String>)>,
mocks: Option<Arc<MockLayer>>,
session_host: Arc<dyn SessionHost>,
session_store: DynSessionStore,
_state_store: DynStateStore,
state_host: Arc<dyn StateHost>,
secrets_manager: DynSecretsManager,
rollout: RolloutIds,
) -> Result<Arc<Self>> {
let operator_registry = OperatorRegistry::build(&packs)?;
let operator_metrics = Arc::new(OperatorMetrics::default());
let pack_runtimes = packs
.iter()
.map(|(pack, _)| Arc::clone(pack))
.collect::<Vec<_>>();
let digests = packs
.iter()
.map(|(_, digest)| digest.clone())
.collect::<Vec<_>>();
let mut pack_trace = HashMap::new();
for (pack, digest) in &packs {
let pack_id = pack.metadata().pack_id.clone();
let pack_ref = config
.pack_bindings
.iter()
.find(|binding| binding.pack_id == pack_id)
.map(|binding| binding.pack_ref.clone())
.unwrap_or_else(|| pack_id.clone());
pack_trace.insert(
pack_id,
PackTraceInfo {
pack_ref,
resolved_digest: digest.clone(),
},
);
}
let engine = Arc::new(
FlowEngine::new(pack_runtimes.clone(), Arc::clone(&config))
.await
.context("failed to prime flow engine")?
.with_rollout_ids(rollout),
);
let state_machine = Arc::new(
StateMachineRuntime::from_flow_engine(
Arc::clone(&config),
Arc::clone(&engine),
pack_trace,
session_host,
Arc::clone(&session_store),
state_host,
Arc::clone(&secrets_manager),
mocks.clone(),
)
.context("failed to initialise state machine runtime")?,
);
let http_client = Client::builder().build()?;
Ok(Arc::new(Self {
tenant: config.tenant.clone(),
config,
packs: pack_runtimes,
digests,
engine,
state_machine,
session_store,
http_client,
mocks,
timer_handles: Mutex::new(Vec::new()),
secrets: secrets_manager,
operator_registry,
operator_metrics,
contract_cache: ContractCache::from_env(),
}))
}
pub fn tenant(&self) -> &str {
&self.tenant
}
pub fn config(&self) -> &Arc<HostConfig> {
&self.config
}
pub fn operator_registry(&self) -> &OperatorRegistry {
&self.operator_registry
}
pub fn operator_metrics(&self) -> &OperatorMetrics {
&self.operator_metrics
}
pub fn contract_cache(&self) -> &ContractCache {
&self.contract_cache
}
pub fn contract_cache_stats(&self) -> ContractCacheStats {
self.contract_cache.stats()
}
pub fn main_pack(&self) -> &Arc<PackRuntime> {
self.packs
.first()
.expect("tenant runtime must contain at least one pack")
}
pub fn pack(&self) -> Arc<PackRuntime> {
Arc::clone(self.main_pack())
}
pub fn overlays(&self) -> Vec<Arc<PackRuntime>> {
self.packs.iter().skip(1).cloned().collect()
}
pub fn all_packs(&self) -> &[Arc<PackRuntime>] {
&self.packs
}
pub fn pack_digests(&self) -> &[Option<String>] {
&self.digests
}
pub fn engine(&self) -> &Arc<FlowEngine> {
&self.engine
}
pub fn state_machine(&self) -> &Arc<StateMachineRuntime> {
&self.state_machine
}
pub fn session_store(&self) -> &DynSessionStore {
&self.session_store
}
pub fn http_client(&self) -> &Client {
&self.http_client
}
pub fn oauth_config(&self) -> Option<OAuthBrokerConfig> {
self.config.oauth_broker_config()
}
pub fn digest(&self) -> Option<&str> {
self.digests.first().and_then(|d| d.as_deref())
}
pub fn overlay_digests(&self) -> Vec<Option<String>> {
self.digests.iter().skip(1).cloned().collect()
}
pub fn required_secrets(&self) -> Vec<SecretRequirement> {
self.packs
.iter()
.flat_map(|pack| pack.required_secrets().iter().cloned())
.collect()
}
pub fn missing_secrets(&self) -> Vec<SecretRequirement> {
self.packs
.iter()
.flat_map(|pack| pack.missing_secrets(&self.config.tenant_ctx()))
.collect()
}
pub fn mocks(&self) -> Option<&Arc<MockLayer>> {
self.mocks.as_ref()
}
pub fn register_timers(&self, handles: Vec<JoinHandle<()>>) {
self.timer_handles.lock().extend(handles);
}
pub fn get_secret(&self, key: &str) -> Result<String> {
if crate::provider_core_only::is_enabled() {
bail!(crate::provider_core_only::blocked_message("secrets"))
}
if !self.config.secrets_policy.is_allowed(key) {
bail!("secret {key} is not permitted by bindings policy");
}
let ctx = self.config.tenant_ctx();
let canonical_key = canonicalize_secret_key(key);
let bytes =
read_secret_blocking(&self.secrets, &ctx, RUNTIME_SECRETS_PACK_ID, &canonical_key)
.context("failed to read secret from manager")?;
let value = String::from_utf8(bytes).context("secret value is not valid UTF-8")?;
Ok(value)
}
pub fn build_events_email_execution_plan(
&self,
tenant: &greentic_types::TenantCtx,
request: &EmailSendRequest,
) -> Result<EmailExecutionPlan> {
let oauth = self
.oauth_config()
.ok_or_else(|| anyhow!("oauth broker config is not configured for tenant runtime"))?;
build_email_execution_plan(&oauth, tenant, request)
}
pub async fn execute_events_email_request(
&self,
access_token: &str,
request: &EmailSendRequest,
) -> Result<()> {
execute_email_request(self.http_client(), access_token, request).await
}
pub async fn execute_events_email_with_oauth(
&self,
tenant: &greentic_types::TenantCtx,
request: &EmailSendRequest,
) -> Result<()> {
let plan = self.build_events_email_execution_plan(tenant, request)?;
let token = request_resource_token(self.http_client(), &plan.token_request).await?;
self.execute_events_email_request(&token.access_token, request)
.await
}
pub fn pack_for_component(&self, component_ref: &str) -> Option<Arc<PackRuntime>> {
self.packs
.iter()
.find(|pack| pack.contains_component(component_ref))
.cloned()
}
pub fn pack_for_component_with_digest(
&self,
component_ref: &str,
) -> Option<(Arc<PackRuntime>, Option<String>)> {
self.packs
.iter()
.zip(self.digests.iter())
.find(|(pack, _)| pack.contains_component(component_ref))
.map(|(pack, digest)| (Arc::clone(pack), digest.clone()))
}
pub fn resolve_component(&self, component_ref: &str) -> Option<ResolvedComponent> {
self.pack_for_component_with_digest(component_ref)
.map(|(pack, digest)| ResolvedComponent {
digest: digest
.or_else(|| self.digest().map(ToString::to_string))
.unwrap_or_else(|| "unknown".to_string()),
component_ref: component_ref.to_string(),
pack,
})
}
}
impl Drop for TenantRuntime {
fn drop(&mut self) {
for handle in self.timer_handles.lock().drain(..) {
handle.abort();
}
}
}
#[cfg(test)]
mod runtime_key_tests {
use super::*;
#[test]
fn legacy_keys_match_only_on_tenant() {
assert_eq!(RuntimeKey::legacy("acme"), RuntimeKey::legacy("acme"));
assert_ne!(RuntimeKey::legacy("acme"), RuntimeKey::legacy("other"));
}
#[test]
fn legacy_and_revision_keys_never_collide() {
let key = RuntimeKey::revision(
"acme",
DeploymentId::new(),
BundleId::from("bundle-a"),
RevisionId::new(),
);
assert_ne!(RuntimeKey::legacy("acme"), key);
}
#[test]
fn revision_keys_distinguish_revision_id() {
let deployment = DeploymentId::new();
let bundle = BundleId::from("bundle-a");
let rev_a = RevisionId::new();
let rev_b = RevisionId::new();
let key_a = RuntimeKey::revision("acme", deployment, bundle.clone(), rev_a);
let key_b = RuntimeKey::revision("acme", deployment, bundle.clone(), rev_b);
let same = RuntimeKey::revision("acme", deployment, bundle, rev_a);
assert_ne!(key_a, key_b);
assert_eq!(key_a, same);
}
#[test]
fn legacy_reload_preserves_revision_entries() {
let revision_key = RuntimeKey::revision(
"acme",
DeploymentId::new(),
BundleId::from("bundle-a"),
RevisionId::new(),
);
let mut prev: HashMap<RuntimeKey, u32> = HashMap::new();
prev.insert(RuntimeKey::legacy("acme"), 1); prev.insert(RuntimeKey::legacy("retired"), 2); prev.insert(revision_key.clone(), 99);
let mut legacy: HashMap<RuntimeKey, u32> = HashMap::new();
legacy.insert(RuntimeKey::legacy("acme"), 10);
legacy.insert(RuntimeKey::legacy("newcomer"), 20);
let next = merge_legacy_reload(&prev, legacy);
assert_eq!(next.get(&RuntimeKey::legacy("acme")), Some(&10));
assert_eq!(next.get(&RuntimeKey::legacy("newcomer")), Some(&20));
assert_eq!(next.get(&RuntimeKey::legacy("retired")), None);
assert_eq!(next.get(&revision_key), Some(&99));
}
#[test]
fn remove_keyed_entry_pops_only_the_targeted_key() {
let deployment = DeploymentId::new();
let bundle = BundleId::from("bundle-a");
let rev_a = RevisionId::new();
let rev_b = RevisionId::new();
let key_a = RuntimeKey::revision("acme", deployment, bundle.clone(), rev_a);
let key_b = RuntimeKey::revision("acme", deployment, bundle, rev_b);
let mut prev: HashMap<RuntimeKey, u32> = HashMap::new();
prev.insert(RuntimeKey::legacy("acme"), 1);
prev.insert(key_a.clone(), 10);
prev.insert(key_b.clone(), 20);
let (next, removed) = remove_keyed_entry(&prev, &key_a).expect("present");
assert_eq!(removed, 10);
assert_eq!(next.get(&key_a), None);
assert_eq!(next.get(&key_b), Some(&20));
assert_eq!(next.get(&RuntimeKey::legacy("acme")), Some(&1));
}
#[test]
fn remove_keyed_entry_returns_none_for_missing_key() {
let prev: HashMap<RuntimeKey, u32> = HashMap::new();
let ghost = RuntimeKey::revision(
"acme",
DeploymentId::new(),
BundleId::from("bundle-a"),
RevisionId::new(),
);
assert!(remove_keyed_entry(&prev, &ghost).is_none());
}
#[test]
fn remove_keyed_entry_leaves_other_deployments_alone() {
let bundle = BundleId::from("bundle-a");
let dep_a = DeploymentId::new();
let dep_b = DeploymentId::new();
let rev = RevisionId::new();
let key_a = RuntimeKey::revision("acme", dep_a, bundle.clone(), rev);
let key_b = RuntimeKey::revision("acme", dep_b, bundle, rev);
let mut prev: HashMap<RuntimeKey, u32> = HashMap::new();
prev.insert(key_a.clone(), 100);
prev.insert(key_b.clone(), 200);
let (next, removed) = remove_keyed_entry(&prev, &key_a).expect("present");
assert_eq!(removed, 100);
assert_eq!(next.get(&key_b), Some(&200));
assert_eq!(next.len(), 1);
}
#[test]
fn map_lookup_separates_legacy_from_revision() {
let deployment = DeploymentId::new();
let bundle = BundleId::from("bundle-a");
let revision = RevisionId::new();
let mut map: HashMap<RuntimeKey, u32> = HashMap::new();
map.insert(RuntimeKey::legacy("acme"), 1);
map.insert(
RuntimeKey::revision("acme", deployment, bundle.clone(), revision),
2,
);
assert_eq!(map.get(&RuntimeKey::legacy("acme")), Some(&1));
assert_eq!(
map.get(&RuntimeKey::revision("acme", deployment, bundle, revision)),
Some(&2)
);
assert_eq!(map.get(&RuntimeKey::legacy("ghost")), None);
}
#[test]
fn insert_keyed_entry_adds_revision_preserving_legacy_and_siblings() {
let deployment = DeploymentId::new();
let bundle = BundleId::from("bundle-a");
let rev_a = RevisionId::new();
let rev_b = RevisionId::new();
let key_a = RuntimeKey::revision("acme", deployment, bundle.clone(), rev_a);
let key_b = RuntimeKey::revision("acme", deployment, bundle, rev_b);
let mut prev: HashMap<RuntimeKey, u32> = HashMap::new();
prev.insert(RuntimeKey::legacy("acme"), 1);
prev.insert(key_a.clone(), 10);
let next = insert_keyed_entry(&prev, key_b.clone(), 20);
assert_eq!(next.get(&key_b), Some(&20));
assert_eq!(next.get(&key_a), Some(&10));
assert_eq!(next.get(&RuntimeKey::legacy("acme")), Some(&1));
assert_eq!(next.len(), 3);
}
#[test]
fn insert_keyed_entry_replaces_existing_revision() {
let key = RuntimeKey::revision(
"acme",
DeploymentId::new(),
BundleId::from("bundle-a"),
RevisionId::new(),
);
let mut prev: HashMap<RuntimeKey, u32> = HashMap::new();
prev.insert(key.clone(), 10);
let next = insert_keyed_entry(&prev, key.clone(), 99);
assert_eq!(next.get(&key), Some(&99));
assert_eq!(next.len(), 1);
}
}