use std::collections::{HashMap, HashSet};
use std::sync::Arc;
use parking_lot::RwLock;
use crate::computation_graph::stream_backend::{
StreamBackendFactory, StreamBackendFuture, StreamConfig,
};
use crate::computation_graph::triggerless::TriggerlessGraphRegistration;
use crate::task::{Task, TaskNamespace};
use crate::tenant_scope::{resolve_tenant_key, visible_keys, TenantKey, TenantOwner, TenantScope};
use crate::trigger::Trigger;
use crate::workflow::Workflow;
use cloacina_computation_graph::{
ComputationGraphConstructor, ComputationGraphRegistration, ReactorConstructor,
ReactorRegistration,
};
pub(crate) type TriggerlessGraphConstructor =
Box<dyn Fn() -> TriggerlessGraphRegistration + Send + Sync>;
pub(crate) type TaskConstructorFn = Box<dyn Fn() -> Arc<dyn Task> + Send + Sync>;
pub(crate) type WorkflowConstructorFn = Box<dyn Fn() -> Workflow + Send + Sync>;
pub(crate) type TriggerConstructorFn = Box<dyn Fn() -> Arc<dyn Trigger> + Send + Sync>;
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
pub enum RuntimeRegistrationError {
#[error(
"{kind} '{name}' is already registered in tenant '{tenant}' by {existing}; \
{incoming} cannot claim the same name — rename the {kind} in one of the two packages"
)]
OwnershipConflict {
kind: &'static str,
name: String,
tenant: String,
existing: String,
incoming: String,
},
}
struct ScopedEntry<V> {
owner: TenantOwner,
value: V,
}
struct ScopedRegistry<V> {
kind: &'static str,
entries: RwLock<HashMap<TenantKey, ScopedEntry<V>>>,
}
impl<V> ScopedRegistry<V> {
fn new(kind: &'static str) -> Self {
Self {
kind,
entries: RwLock::new(HashMap::new()),
}
}
fn insert(
&self,
scope: TenantScope<'_>,
name: String,
owner: TenantOwner,
value: V,
) -> Result<(), RuntimeRegistrationError> {
let key = scope.own_key(&name);
let mut guard = self.entries.write();
if let Some(existing) = guard.get(&key) {
if !existing.owner.may_replace(&owner) {
return Err(RuntimeRegistrationError::OwnershipConflict {
kind: self.kind,
name,
tenant: key.tenant_label().to_string(),
existing: existing.owner.label(),
incoming: owner.label(),
});
}
}
guard.insert(key, ScopedEntry { owner, value });
Ok(())
}
fn with<R>(&self, scope: TenantScope<'_>, name: &str, f: impl FnOnce(&V) -> R) -> Option<R> {
let guard = self.entries.read();
let key = resolve_tenant_key(&*guard, scope, name).ok()?;
guard.get(&key).map(|entry| f(&entry.value))
}
fn remove(&self, scope: TenantScope<'_>, name: &str) -> bool {
let mut guard = self.entries.write();
let Ok(key) = resolve_tenant_key(&*guard, scope, name) else {
return false;
};
guard.remove(&key).is_some()
}
fn names(&self, scope: TenantScope<'_>) -> Vec<String> {
let guard = self.entries.read();
let mut seen = HashSet::new();
visible_keys(&*guard, scope)
.filter(|k| seen.insert(k.name.clone()))
.map(|k| k.name.clone())
.collect()
}
fn all_keys(&self) -> Vec<TenantKey> {
self.entries.read().keys().cloned().collect()
}
fn owner(&self, scope: TenantScope<'_>, name: &str) -> Option<TenantOwner> {
let guard = self.entries.read();
let key = resolve_tenant_key(&*guard, scope, name).ok()?;
guard.get(&key).map(|entry| entry.owner.clone())
}
fn conflict(
&self,
scope: TenantScope<'_>,
name: &str,
existing: &TenantOwner,
incoming: &TenantOwner,
) -> RuntimeRegistrationError {
RuntimeRegistrationError::OwnershipConflict {
kind: self.kind,
name: name.to_string(),
tenant: scope.own_key(name).tenant_label().to_string(),
existing: existing.label(),
incoming: incoming.label(),
}
}
fn check_claim(
&self,
scope: TenantScope<'_>,
name: &str,
incoming: &TenantOwner,
) -> Result<bool, RuntimeRegistrationError> {
match self.owner(scope, name) {
None => Ok(true),
Some(existing) if existing.may_replace(incoming) => Ok(false),
Some(existing) => Err(self.conflict(scope, name, &existing, incoming)),
}
}
fn len(&self) -> usize {
self.entries.read().len()
}
}
#[derive(Clone)]
pub struct Runtime {
inner: Arc<RuntimeInner>,
scope: RuntimeScope,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
struct RuntimeScope {
tenant_id: Option<String>,
is_admin: bool,
}
impl RuntimeScope {
fn as_tenant_scope(&self) -> TenantScope<'_> {
TenantScope {
tenant_id: self.tenant_id.as_deref(),
is_admin: self.is_admin,
}
}
}
struct RuntimeInner {
tasks: RwLock<HashMap<TaskNamespace, TaskConstructorFn>>,
workflows: ScopedRegistry<WorkflowConstructorFn>,
triggers: ScopedRegistry<TriggerConstructorFn>,
computation_graphs: ScopedRegistry<ComputationGraphConstructor>,
triggerless_graphs: ScopedRegistry<TriggerlessGraphConstructor>,
reactors: ScopedRegistry<ReactorConstructor>,
stream_backends: RwLock<HashMap<String, StreamBackendFactory>>,
}
impl Runtime {
pub fn new() -> Self {
let rt = Self::empty();
rt.seed_from_inventory();
rt
}
pub fn empty() -> Self {
Self {
inner: Arc::new(RuntimeInner {
tasks: RwLock::new(HashMap::new()),
workflows: ScopedRegistry::new("workflow"),
triggers: ScopedRegistry::new("trigger"),
computation_graphs: ScopedRegistry::new("computation graph"),
triggerless_graphs: ScopedRegistry::new("trigger-less computation graph"),
reactors: ScopedRegistry::new("reactor"),
stream_backends: RwLock::new(HashMap::new()),
}),
scope: RuntimeScope::default(),
}
}
fn scope(&self) -> TenantScope<'_> {
self.scope.as_tenant_scope()
}
pub fn scoped_to_tenant(&self, tenant_id: impl Into<String>) -> Self {
Self {
inner: Arc::clone(&self.inner),
scope: RuntimeScope {
tenant_id: Some(tenant_id.into()),
is_admin: false,
},
}
}
pub fn untenanted_view(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
scope: RuntimeScope::default(),
}
}
pub fn admin_view(&self) -> Self {
Self {
inner: Arc::clone(&self.inner),
scope: RuntimeScope {
tenant_id: None,
is_admin: true,
},
}
}
pub fn tenant_id(&self) -> Option<&str> {
self.scope.tenant_id.as_deref()
}
pub fn is_admin(&self) -> bool {
self.scope.is_admin
}
pub fn shares_registries_with(&self, other: &Runtime) -> bool {
Arc::ptr_eq(&self.inner, &other.inner)
}
pub fn seed_from_inventory(&self) {
use crate::inventory_entries::{
ComputationGraphEntry, ReactorEntry, StreamBackendEntry, TaskEntry, TriggerEntry,
TriggerlessGraphEntry, WorkflowEntry,
};
let global = self.untenanted_view();
for entry in inventory::iter::<TaskEntry> {
let ns = (entry.namespace)();
let ctor = entry.constructor;
self.register_task(ns, move || ctor());
}
for entry in inventory::iter::<WorkflowEntry> {
global.register_workflow(entry.name.to_string(), entry.constructor);
}
for entry in inventory::iter::<TriggerEntry> {
global.register_trigger(entry.name.to_string(), entry.constructor);
}
for entry in inventory::iter::<ComputationGraphEntry> {
global.register_computation_graph(entry.name.to_string(), entry.constructor);
}
for entry in inventory::iter::<TriggerlessGraphEntry> {
global.register_triggerless_graph(entry.name.to_string(), entry.constructor);
}
for entry in inventory::iter::<ReactorEntry> {
global.register_reactor(entry.name.to_string(), entry.constructor);
}
for entry in inventory::iter::<StreamBackendEntry> {
let factory = entry.factory;
self.register_stream_backend(
entry.type_name.to_string(),
Box::new(move |config| factory(config)),
);
}
}
pub fn register_task<F>(&self, namespace: TaskNamespace, factory: F)
where
F: Fn() -> Arc<dyn Task> + Send + Sync + 'static,
{
self.inner
.tasks
.write()
.insert(namespace, Box::new(factory));
}
pub fn unregister_task(&self, namespace: &TaskNamespace) -> bool {
self.inner.tasks.write().remove(namespace).is_some()
}
pub fn get_task(&self, namespace: &TaskNamespace) -> Option<Arc<dyn Task>> {
self.inner.tasks.read().get(namespace).map(|ctor| ctor())
}
#[cfg(test)]
pub(crate) fn has_task(&self, namespace: &TaskNamespace) -> bool {
self.inner.tasks.read().contains_key(namespace)
}
pub fn task_namespaces(&self) -> Vec<TaskNamespace> {
self.inner.tasks.read().keys().cloned().collect()
}
pub fn register_workflow<F>(&self, name: String, constructor: F)
where
F: Fn() -> Workflow + Send + Sync + 'static,
{
let _ = self.inner.workflows.insert(
self.scope(),
name,
TenantOwner::unknown(),
Box::new(constructor),
);
}
pub fn try_register_workflow<F>(
&self,
owner: &TenantOwner,
name: String,
constructor: F,
) -> Result<(), RuntimeRegistrationError>
where
F: Fn() -> Workflow + Send + Sync + 'static,
{
self.inner
.workflows
.insert(self.scope(), name, owner.clone(), Box::new(constructor))
}
pub fn unregister_workflow(&self, name: &str) -> bool {
self.inner.workflows.remove(self.scope(), name)
}
pub fn get_workflow(&self, name: &str) -> Option<Workflow> {
self.inner.workflows.with(self.scope(), name, |ctor| ctor())
}
pub fn workflow_names(&self) -> Vec<String> {
self.inner.workflows.names(self.scope())
}
pub fn workflow_keys(&self) -> Vec<TenantKey> {
self.inner.workflows.all_keys()
}
pub fn register_trigger<F>(&self, name: String, factory: F)
where
F: Fn() -> Arc<dyn Trigger> + Send + Sync + 'static,
{
let _ = self.inner.triggers.insert(
self.scope(),
name,
TenantOwner::unknown(),
Box::new(factory),
);
}
pub fn try_register_trigger<F>(
&self,
owner: &TenantOwner,
name: String,
factory: F,
) -> Result<(), RuntimeRegistrationError>
where
F: Fn() -> Arc<dyn Trigger> + Send + Sync + 'static,
{
self.inner
.triggers
.insert(self.scope(), name, owner.clone(), Box::new(factory))
}
pub fn may_claim_trigger(
&self,
owner: &TenantOwner,
name: &str,
) -> Result<bool, RuntimeRegistrationError> {
self.inner.triggers.check_claim(self.scope(), name, owner)
}
pub fn unregister_trigger(&self, name: &str) -> bool {
self.inner.triggers.remove(self.scope(), name)
}
pub fn get_trigger(&self, name: &str) -> Option<Arc<dyn Trigger>> {
self.inner.triggers.with(self.scope(), name, |ctor| ctor())
}
pub fn trigger_names(&self) -> Vec<String> {
self.inner.triggers.names(self.scope())
}
pub fn trigger_keys(&self) -> Vec<TenantKey> {
self.inner.triggers.all_keys()
}
pub fn register_computation_graph<F>(&self, name: String, constructor: F)
where
F: Fn() -> ComputationGraphRegistration + Send + Sync + 'static,
{
let _ = self.inner.computation_graphs.insert(
self.scope(),
name,
TenantOwner::unknown(),
Box::new(constructor),
);
}
pub fn try_register_computation_graph<F>(
&self,
owner: &TenantOwner,
name: String,
constructor: F,
) -> Result<(), RuntimeRegistrationError>
where
F: Fn() -> ComputationGraphRegistration + Send + Sync + 'static,
{
self.inner.computation_graphs.insert(
self.scope(),
name,
owner.clone(),
Box::new(constructor),
)
}
pub fn unregister_computation_graph(&self, name: &str) -> bool {
self.inner.computation_graphs.remove(self.scope(), name)
}
pub fn get_computation_graph(&self, name: &str) -> Option<ComputationGraphRegistration> {
self.inner
.computation_graphs
.with(self.scope(), name, |ctor| ctor())
}
pub fn computation_graph_names(&self) -> Vec<String> {
self.inner.computation_graphs.names(self.scope())
}
pub fn computation_graph_keys(&self) -> Vec<TenantKey> {
self.inner.computation_graphs.all_keys()
}
pub fn register_triggerless_graph<F>(&self, name: String, constructor: F)
where
F: Fn() -> TriggerlessGraphRegistration + Send + Sync + 'static,
{
let _ = self.inner.triggerless_graphs.insert(
self.scope(),
name,
TenantOwner::unknown(),
Box::new(constructor),
);
}
pub fn try_register_triggerless_graph<F>(
&self,
owner: &TenantOwner,
name: String,
constructor: F,
) -> Result<(), RuntimeRegistrationError>
where
F: Fn() -> TriggerlessGraphRegistration + Send + Sync + 'static,
{
self.inner.triggerless_graphs.insert(
self.scope(),
name,
owner.clone(),
Box::new(constructor),
)
}
pub fn may_claim_triggerless_graph(
&self,
owner: &TenantOwner,
name: &str,
) -> Result<bool, RuntimeRegistrationError> {
self.inner
.triggerless_graphs
.check_claim(self.scope(), name, owner)
}
pub fn unregister_triggerless_graph(&self, name: &str) -> bool {
self.inner.triggerless_graphs.remove(self.scope(), name)
}
pub fn get_triggerless_graph(&self, name: &str) -> Option<TriggerlessGraphRegistration> {
self.inner
.triggerless_graphs
.with(self.scope(), name, |ctor| ctor())
}
pub fn triggerless_graph_names(&self) -> Vec<String> {
self.inner.triggerless_graphs.names(self.scope())
}
pub fn triggerless_graph_keys(&self) -> Vec<TenantKey> {
self.inner.triggerless_graphs.all_keys()
}
pub fn register_reactor<F>(&self, name: String, constructor: F)
where
F: Fn() -> ReactorRegistration + Send + Sync + 'static,
{
let _ = self.inner.reactors.insert(
self.scope(),
name,
TenantOwner::unknown(),
Box::new(constructor),
);
}
pub fn try_register_reactor<F>(
&self,
owner: &TenantOwner,
name: String,
constructor: F,
) -> Result<(), RuntimeRegistrationError>
where
F: Fn() -> ReactorRegistration + Send + Sync + 'static,
{
self.inner
.reactors
.insert(self.scope(), name, owner.clone(), Box::new(constructor))
}
pub fn unregister_reactor(&self, name: &str) -> bool {
self.inner.reactors.remove(self.scope(), name)
}
pub fn get_reactor(&self, name: &str) -> Option<ReactorRegistration> {
self.inner.reactors.with(self.scope(), name, |ctor| ctor())
}
pub fn reactor_names(&self) -> Vec<String> {
self.inner.reactors.names(self.scope())
}
pub fn reactor_keys(&self) -> Vec<TenantKey> {
self.inner.reactors.all_keys()
}
pub fn register_stream_backend(&self, type_name: String, factory: StreamBackendFactory) {
self.inner
.stream_backends
.write()
.insert(type_name, factory);
}
pub fn unregister_stream_backend(&self, type_name: &str) -> bool {
self.inner
.stream_backends
.write()
.remove(type_name)
.is_some()
}
#[cfg(test)]
pub(crate) fn has_stream_backend(&self, type_name: &str) -> bool {
self.inner.stream_backends.read().contains_key(type_name)
}
pub fn create_stream_backend(
&self,
type_name: &str,
config: StreamConfig,
) -> Option<StreamBackendFuture> {
let guard = self.inner.stream_backends.read();
let factory = guard.get(type_name)?;
Some(factory(config))
}
#[cfg(test)]
pub(crate) fn stream_backend_names(&self) -> Vec<String> {
self.inner.stream_backends.read().keys().cloned().collect()
}
}
impl Default for Runtime {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Debug for Runtime {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let tasks = self.inner.tasks.read().len();
let workflows = self.inner.workflows.len();
let triggers = self.inner.triggers.len();
let cgs = self.inner.computation_graphs.len();
let sbs = self.inner.stream_backends.read().len();
f.debug_struct("Runtime")
.field("tenant", &self.scope.tenant_id)
.field("is_admin", &self.scope.is_admin)
.field("tasks", &tasks)
.field("workflows", &workflows)
.field("triggers", &triggers)
.field("computation_graphs", &cgs)
.field("stream_backends", &sbs)
.finish()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::task::TaskNamespace;
#[test]
fn register_and_unregister_workflow() {
let rt = Runtime::empty();
assert!(!rt.unregister_workflow("nope"));
let wf = crate::workflow::Workflow::new("unit-test-wf");
rt.register_workflow("unit-test-wf".to_string(), move || wf.clone());
assert!(rt.get_workflow("unit-test-wf").is_some());
assert_eq!(rt.workflow_names(), vec!["unit-test-wf".to_string()]);
assert!(rt.unregister_workflow("unit-test-wf"));
assert!(rt.get_workflow("unit-test-wf").is_none());
assert!(rt.workflow_names().is_empty());
}
#[test]
fn register_and_unregister_trigger_by_name() {
let rt = Runtime::empty();
assert!(!rt.unregister_trigger("missing"));
assert!(rt.get_trigger("missing").is_none());
assert!(rt.trigger_names().is_empty());
}
#[test]
fn register_and_unregister_task() {
let rt = Runtime::empty();
let ns = TaskNamespace::new("t", "p", "w", "task_a");
assert!(!rt.unregister_task(&ns));
assert!(!rt.has_task(&ns));
}
#[test]
fn stream_backend_roundtrip_names_only() {
let rt = Runtime::empty();
assert!(!rt.has_stream_backend("mock"));
assert!(rt.stream_backend_names().is_empty());
assert!(!rt.unregister_stream_backend("mock"));
}
#[test]
fn runtimes_are_independent() {
let rt1 = Runtime::empty();
let rt2 = Runtime::empty();
let wf = crate::workflow::Workflow::new("iso");
rt1.register_workflow("iso".to_string(), move || wf.clone());
assert!(rt1.get_workflow("iso").is_some());
assert!(rt2.get_workflow("iso").is_none());
}
#[test]
fn debug_format_reports_sizes() {
let rt = Runtime::empty();
let debug = format!("{:?}", rt);
assert!(debug.contains("computation_graphs: 0"));
assert!(debug.contains("stream_backends: 0"));
}
fn wf(name: &str) -> crate::workflow::Workflow {
crate::workflow::Workflow::new(name)
}
#[test]
fn two_tenants_same_workflow_name_do_not_collide() {
let shared = Runtime::empty();
let acme = shared.scoped_to_tenant("acme");
let globex = shared.scoped_to_tenant("globex");
assert!(acme.shares_registries_with(&globex));
let a = wf("acme-only-desc");
acme.register_workflow("pipeline".to_string(), move || a.clone());
let g = wf("globex-only-desc");
globex.register_workflow("pipeline".to_string(), move || g.clone());
assert_eq!(
acme.get_workflow("pipeline").unwrap().name(),
"acme-only-desc"
);
assert_eq!(
globex.get_workflow("pipeline").unwrap().name(),
"globex-only-desc"
);
assert_eq!(shared.workflow_keys().len(), 2);
}
#[test]
fn unregister_is_tenant_scoped() {
let shared = Runtime::empty();
let acme = shared.scoped_to_tenant("acme");
let globex = shared.scoped_to_tenant("globex");
let a = wf("a");
acme.register_workflow("pipeline".to_string(), move || a.clone());
let g = wf("g");
globex.register_workflow("pipeline".to_string(), move || g.clone());
assert!(acme.unregister_workflow("pipeline"));
assert!(acme.get_workflow("pipeline").is_none());
assert_eq!(globex.get_workflow("pipeline").unwrap().name(), "g");
assert_eq!(shared.workflow_keys().len(), 1);
}
#[test]
fn other_tenants_entries_are_invisible() {
let shared = Runtime::empty();
let acme = shared.scoped_to_tenant("acme");
let a = wf("a");
acme.register_workflow("private".to_string(), move || a.clone());
let globex = shared.scoped_to_tenant("globex");
assert!(globex.get_workflow("private").is_none());
assert!(globex.workflow_names().is_empty());
assert!(!globex.unregister_workflow("private"));
assert!(acme.get_workflow("private").is_some());
}
#[test]
fn same_tenant_cross_package_collision_is_loud() {
let rt = Runtime::empty().scoped_to_tenant("acme");
let first = TenantOwner::package("pkg-a");
let second = TenantOwner::package("pkg-b");
let a = wf("from-a");
rt.try_register_workflow(&first, "reports".to_string(), move || a.clone())
.expect("first claim");
let b = wf("from-b");
let err = rt
.try_register_workflow(&second, "reports".to_string(), move || b.clone())
.expect_err("second package must not silently replace");
let msg = err.to_string();
assert!(msg.contains("pkg-a"), "{msg}");
assert!(msg.contains("pkg-b"), "{msg}");
assert!(msg.contains("acme"), "{msg}");
assert_eq!(rt.get_workflow("reports").unwrap().name(), "from-a");
let a2 = wf("from-a-v2");
rt.try_register_workflow(&first, "reports".to_string(), move || a2.clone())
.expect("same package may replace");
assert_eq!(rt.get_workflow("reports").unwrap().name(), "from-a-v2");
}
#[test]
fn cross_tenant_same_name_is_not_a_conflict() {
let shared = Runtime::empty();
let owner = TenantOwner::package("pkg-a");
for tenant in ["acme", "globex"] {
let w = wf(tenant);
shared
.scoped_to_tenant(tenant)
.try_register_workflow(&owner, "reports".to_string(), move || w.clone())
.expect("distinct tenants never collide");
}
assert_eq!(shared.workflow_keys().len(), 2);
}
#[test]
fn untenanted_path_is_unchanged() {
let rt = Runtime::empty();
assert_eq!(rt.tenant_id(), None);
let w = wf("embedded");
rt.register_workflow("embedded".to_string(), move || w.clone());
assert!(rt.get_workflow("embedded").is_some());
assert_eq!(rt.workflow_names(), vec!["embedded".to_string()]);
assert!(rt.unregister_workflow("embedded"));
assert!(rt.get_workflow("embedded").is_none());
assert!(rt.workflow_names().is_empty());
}
#[test]
fn untenanted_entries_are_visible_from_a_tenant_view() {
let shared = Runtime::empty();
let w = wf("inventory");
shared.register_workflow("inventory".to_string(), move || w.clone());
let acme = shared.scoped_to_tenant("acme");
assert!(acme.get_workflow("inventory").is_some());
assert_eq!(acme.workflow_names(), vec!["inventory".to_string()]);
}
#[test]
fn own_tenant_shadows_untenanted() {
let shared = Runtime::empty();
let global = wf("global");
shared.register_workflow("dup".to_string(), move || global.clone());
let acme = shared.scoped_to_tenant("acme");
let own = wf("own");
acme.register_workflow("dup".to_string(), move || own.clone());
assert_eq!(acme.get_workflow("dup").unwrap().name(), "own");
assert_eq!(shared.get_workflow("dup").unwrap().name(), "global");
assert_eq!(acme.workflow_names(), vec!["dup".to_string()]);
}
#[test]
fn admin_view_resolves_unique_and_refuses_ambiguous() {
let shared = Runtime::empty();
let a = wf("a");
shared
.scoped_to_tenant("acme")
.register_workflow("only-acme".to_string(), move || a.clone());
let admin = shared.admin_view();
assert!(admin.is_admin());
assert!(admin.get_workflow("only-acme").is_some());
for tenant in ["acme", "globex"] {
let w = wf(tenant);
shared
.scoped_to_tenant(tenant)
.register_workflow("shared-name".to_string(), move || w.clone());
}
assert!(
admin.get_workflow("shared-name").is_none(),
"an ambiguous name must not resolve to a guess"
);
let mut names = admin.workflow_names();
names.sort();
assert_eq!(
names,
vec!["only-acme".to_string(), "shared-name".to_string()]
);
}
#[test]
fn reactors_and_graphs_are_tenant_scoped_too() {
let shared = Runtime::empty();
let acme = shared.scoped_to_tenant("acme");
let globex = shared.scoped_to_tenant("globex");
for rt in [&acme, &globex] {
rt.register_triggerless_graph("g".to_string(), || unreachable!());
rt.register_computation_graph("cg".to_string(), || unreachable!());
}
assert_eq!(shared.triggerless_graph_keys().len(), 2);
assert_eq!(shared.computation_graph_keys().len(), 2);
assert!(acme.unregister_triggerless_graph("g"));
assert!(!acme.unregister_triggerless_graph("g"));
assert_eq!(globex.triggerless_graph_names(), vec!["g".to_string()]);
assert!(acme.unregister_computation_graph("cg"));
assert_eq!(globex.computation_graph_names(), vec!["cg".to_string()]);
}
#[test]
fn inventory_seeding_is_always_untenanted() {
let shared = Runtime::empty();
shared.scoped_to_tenant("acme").seed_from_inventory();
assert!(
shared.workflow_keys().iter().all(|k| k.tenant_id.is_none()),
"inventory entries must never be stamped with a tenant"
);
assert!(shared.reactor_keys().iter().all(|k| k.tenant_id.is_none()));
assert!(shared.trigger_keys().iter().all(|k| k.tenant_id.is_none()));
}
}