use aion_core::WorkflowId;
use aion_store_haematite::HaematiteStore;
use super::directory::{NodeRef, OwnerView, ShardDirectory};
const REMINT_ATTEMPTS_PER_SHARD: usize = 16;
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RouteDecision {
Local,
Forward {
owner: NodeRef,
shard: usize,
},
NotOwner {
shard: usize,
},
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum RemintOutcome {
EngineMint,
UseId(WorkflowId),
}
const STEER_ATTEMPTS_PER_SHARD: usize = 16;
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum SteerDecision {
Local(WorkflowId),
Forward {
owner: NodeRef,
shard: usize,
},
NotOwner {
shard: usize,
},
}
#[must_use]
pub fn route_mutation(
cluster_store: Option<&HaematiteStore>,
directory: Option<&dyn ShardDirectory>,
workflow_id: &WorkflowId,
) -> RouteDecision {
let Some(store) = cluster_store else {
return RouteDecision::Local;
};
let shard = store.shard_for_workflow(workflow_id);
let Some(directory) = directory else {
return if store.owns_workflow_shard(workflow_id) {
RouteDecision::Local
} else {
RouteDecision::NotOwner { shard }
};
};
match directory.owner_of(shard) {
OwnerView::Local | OwnerView::Unknown => RouteDecision::Local,
OwnerView::Remote(owner) => RouteDecision::Forward { owner, shard },
}
}
#[must_use]
pub fn route_start(cluster_store: Option<&HaematiteStore>) -> RemintOutcome {
let Some(store) = cluster_store else {
return RemintOutcome::EngineMint;
};
let budget = store.shard_count().max(1) * REMINT_ATTEMPTS_PER_SHARD;
match store.remint_for_owned_shard(budget) {
Some(workflow_id) => RemintOutcome::UseId(workflow_id),
None => RemintOutcome::EngineMint,
}
}
#[must_use]
pub fn route_start_steered(
store: &HaematiteStore,
directory: Option<&dyn ShardDirectory>,
routing_key: &str,
) -> SteerDecision {
let shard = store.shard_for_routing_key(routing_key);
let owner = directory.map_or(OwnerView::Unknown, |directory| directory.owner_of(shard));
match owner {
OwnerView::Remote(owner) if owner.grpc_addr.is_some() => {
SteerDecision::Forward { owner, shard }
}
OwnerView::Remote(_) => SteerDecision::NotOwner { shard },
OwnerView::Local | OwnerView::Unknown => SteerDecision::Local(mint_on_shard(store, shard)),
}
}
fn mint_on_shard(store: &HaematiteStore, shard: usize) -> WorkflowId {
let budget = store.shard_count().max(1) * STEER_ATTEMPTS_PER_SHARD;
store
.mint_for_shard(shard, budget)
.unwrap_or_else(WorkflowId::new_v4)
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{SystemTime, UNIX_EPOCH};
use super::super::directory::{DirectoryPeer, StaticShardDirectory};
use super::{
RemintOutcome, RouteDecision, SteerDecision, route_mutation, route_start,
route_start_steered,
};
use aion_core::WorkflowId;
use aion_store::StoreError;
use aion_store_haematite::HaematiteStore;
type TestResult = Result<(), StoreError>;
fn unique_dir(name: &str) -> PathBuf {
static COUNTER: AtomicU64 = AtomicU64::new(0);
let nanos = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_or(0, |duration| duration.as_nanos());
let counter = COUNTER.fetch_add(1, Ordering::Relaxed);
std::env::temp_dir().join(format!(
"aion-routing-edge-{name}-{}-{nanos}-{counter}",
std::process::id()
))
}
fn test_node_cache_budget() -> Result<haematite::NodeCacheBudget, StoreError> {
haematite::NodeCacheBudget::bytes(1 << 30)
.map_err(|error| StoreError::Backend(error.to_string()))
}
fn store_owning(
name: &str,
shard_count: usize,
owned: &[usize],
) -> Result<HaematiteStore, StoreError> {
let store = HaematiteStore::create_with_shard_count(
unique_dir(name),
shard_count,
test_node_cache_budget()?,
)?;
store.set_owned_shards(owned.iter().copied());
Ok(store)
}
#[test]
fn mutation_without_cluster_store_is_local() {
let workflow_id = WorkflowId::new_v4();
assert_eq!(
route_mutation(None, None, &workflow_id),
RouteDecision::Local
);
}
#[test]
fn mutation_without_directory_falls_back_to_bare_ownership() -> TestResult {
let store = store_owning("mutation", 4, &[0])?;
let mut owned_id = None;
let mut foreign_id = None;
for _ in 0..10_000 {
let candidate = WorkflowId::new_v4();
if store.owns_workflow_shard(&candidate) {
owned_id.get_or_insert(candidate);
} else {
foreign_id.get_or_insert(candidate);
}
if owned_id.is_some() && foreign_id.is_some() {
break;
}
}
let (Some(owned_id), Some(foreign_id)) = (owned_id, foreign_id) else {
return Err(StoreError::Backend(
"expected both an owned and a non-owned shard id".to_owned(),
));
};
assert_eq!(
route_mutation(Some(&store), None, &owned_id),
RouteDecision::Local
);
let shard = store.shard_for_workflow(&foreign_id);
assert_eq!(
route_mutation(Some(&store), None, &foreign_id),
RouteDecision::NotOwner { shard }
);
Ok(())
}
#[test]
fn mutation_with_directory_routes_unknown_owner_locally() -> TestResult {
use super::super::directory::{DirectoryPeer, StaticShardDirectory};
let store = std::sync::Arc::new(store_owning("dir-unknown", 4, &[0])?);
let directory = StaticShardDirectory::new(
std::sync::Arc::clone(&store),
vec![DirectoryPeer {
name: "peer-1".to_owned(),
owned_shards: vec![1, 2, 3],
grpc_addr: None,
}],
None,
);
let mut foreign_id = None;
for _ in 0..10_000 {
let candidate = WorkflowId::new_v4();
if !store.owns_workflow_shard(&candidate) {
foreign_id = Some(candidate);
break;
}
}
let Some(foreign_id) = foreign_id else {
return Err(StoreError::Backend("expected a non-owned id".to_owned()));
};
assert_eq!(
route_mutation(Some(store.as_ref()), Some(&directory), &foreign_id),
RouteDecision::Local
);
Ok(())
}
#[test]
fn start_without_cluster_store_uses_engine_mint() {
assert_eq!(route_start(None), RemintOutcome::EngineMint);
}
#[test]
fn start_with_own_all_scope_uses_engine_mint() -> TestResult {
let store = HaematiteStore::create_with_shard_count(
unique_dir("ownall"),
4,
test_node_cache_budget()?,
)?;
assert_eq!(route_start(Some(&store)), RemintOutcome::EngineMint);
Ok(())
}
#[test]
fn start_reminted_id_lands_on_an_owned_shard() -> TestResult {
let store = store_owning("remint", 4, &[1])?;
let RemintOutcome::UseId(workflow_id) = route_start(Some(&store)) else {
return Err(StoreError::Backend(
"subset-owning node must remint, not engine-mint".to_owned(),
));
};
assert!(
store.owns_workflow_shard(&workflow_id),
"reminted id must land on an owned shard"
);
assert_eq!(store.shard_for_workflow(&workflow_id), 1);
Ok(())
}
#[test]
fn steered_start_to_owned_shard_runs_locally() -> TestResult {
let store = HaematiteStore::create_with_shard_count(
unique_dir("steer-local"),
4,
test_node_cache_budget()?,
)?;
let key = "tenant-a/order-1";
let target = store.shard_for_routing_key(key);
let SteerDecision::Local(workflow_id) = route_start_steered(&store, None, key) else {
return Err(StoreError::Backend(
"own-all node must run a steered start locally".to_owned(),
));
};
assert_eq!(
store.shard_for_workflow(&workflow_id),
target,
"the minted id must land on the routing key's shard"
);
Ok(())
}
#[test]
fn steered_start_to_remote_shard_forwards() -> TestResult {
let store = std::sync::Arc::new(store_owning("steer-remote", 4, &[0])?);
let mut key = None;
for index in 0..100_000_u64 {
let candidate = format!("k-{index}");
let shard = store.shard_for_routing_key(&candidate);
if shard != 0 {
key = Some((candidate, shard));
break;
}
}
let Some((key, shard)) = key else {
return Err(StoreError::Backend(
"no off-owner routing key found".to_owned(),
));
};
let grpc_addr = "127.0.0.1:6001"
.parse()
.map_err(|error| StoreError::Backend(format!("bad addr: {error}")))?;
let directory = StaticShardDirectory::new(
std::sync::Arc::clone(&store),
vec![DirectoryPeer {
name: "peer-1".to_owned(),
owned_shards: vec![1, 2, 3],
grpc_addr: Some(grpc_addr),
}],
None,
);
let SteerDecision::Local(workflow_id) =
route_start_steered(store.as_ref(), Some(&directory), &key)
else {
return Err(StoreError::Backend(
"a believed-down owner must route the steered start locally".to_owned(),
));
};
assert_eq!(store.shard_for_workflow(&workflow_id), shard);
Ok(())
}
#[test]
fn steered_start_remote_without_forward_addr_is_not_owner() {
let decision = SteerDecision::NotOwner { shard: 3 };
assert!(matches!(decision, SteerDecision::NotOwner { shard: 3 }));
}
}