use std::collections::BTreeSet;
use async_trait::async_trait;
use chrono::{DateTime, SecondsFormat, Utc};
use serde::{Deserialize, Serialize};
use crate::StoreError;
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct NamespaceRecord {
pub name: String,
pub created_at: DateTime<Utc>,
pub last_seen: DateTime<Utc>,
pub origin: NamespaceOrigin,
pub config: NamespaceConfig,
pub placement: NamespacePlacement,
pub state: NamespaceState,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum NamespaceOrigin {
WorkerMint,
StartMint,
Explicit,
InferredFromState,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum NamespaceState {
Active,
Deprecated,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum MintOutcome {
Created,
AlreadyExisted,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct NamespaceConfig {
pub kind: Option<String>,
pub max_in_flight_activities: Option<u32>,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub enum NamespacePlacement {
#[default]
Unplaced,
Prefer {
nodes: BTreeSet<String>,
},
Pinned {
nodes: BTreeSet<String>,
},
}
impl NamespaceRecord {
#[must_use]
pub fn new_minted(name: &str, origin: NamespaceOrigin, now: DateTime<Utc>) -> Self {
Self {
name: name.to_owned(),
created_at: now,
last_seen: now,
origin,
config: NamespaceConfig::default(),
placement: NamespacePlacement::default(),
state: NamespaceState::Active,
}
}
pub fn bump_last_seen(&mut self, now: DateTime<Utc>) {
self.last_seen = now;
}
pub fn encode(&self) -> Result<Vec<u8>, StoreError> {
let stored = StoredNamespace {
name: self.name.clone(),
created_at: encode_instant(self.created_at),
last_seen: encode_instant(self.last_seen),
origin: self.origin,
kind: self.config.kind.clone(),
max_in_flight_activities: self.config.max_in_flight_activities,
placement: self.placement.clone(),
state: self.state,
};
serde_json::to_vec(&stored).map_err(|error| StoreError::Serialization(error.to_string()))
}
pub fn decode(bytes: &[u8]) -> Result<Self, StoreError> {
let stored: StoredNamespace = serde_json::from_slice(bytes)
.map_err(|error| StoreError::Serialization(error.to_string()))?;
Ok(Self {
name: stored.name,
created_at: decode_instant(&stored.created_at)?,
last_seen: decode_instant(&stored.last_seen)?,
origin: stored.origin,
config: NamespaceConfig {
kind: stored.kind,
max_in_flight_activities: stored.max_in_flight_activities,
},
placement: stored.placement,
state: stored.state,
})
}
}
#[derive(Serialize, Deserialize)]
struct StoredNamespace {
name: String,
created_at: String,
last_seen: String,
origin: NamespaceOrigin,
kind: Option<String>,
#[serde(default)]
max_in_flight_activities: Option<u32>,
placement: NamespacePlacement,
state: NamespaceState,
}
const PLACEMENT_UNPLACED_TAG: &str = "unplaced";
const PLACEMENT_PREFER_TAG: &str = "prefer";
const PLACEMENT_PINNED_TAG: &str = "pinned";
impl Serialize for NamespacePlacement {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
use serde::ser::SerializeMap;
match self {
Self::Unplaced => serializer.serialize_str(PLACEMENT_UNPLACED_TAG),
Self::Prefer { nodes } => {
let mut map = serializer.serialize_map(Some(1))?;
map.serialize_entry(PLACEMENT_PREFER_TAG, nodes)?;
map.end()
}
Self::Pinned { nodes } => {
let mut map = serializer.serialize_map(Some(1))?;
map.serialize_entry(PLACEMENT_PINNED_TAG, nodes)?;
map.end()
}
}
}
}
impl<'de> Deserialize<'de> for NamespacePlacement {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: serde::Deserializer<'de>,
{
deserializer.deserialize_any(PlacementVisitor)
}
}
struct PlacementVisitor;
impl<'de> serde::de::Visitor<'de> for PlacementVisitor {
type Value = NamespacePlacement;
fn expecting(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
formatter.write_str("\"unplaced\" or a single-key {prefer|pinned: [labels]} map")
}
fn visit_str<E>(self, value: &str) -> Result<Self::Value, E>
where
E: serde::de::Error,
{
match value {
PLACEMENT_UNPLACED_TAG => Ok(NamespacePlacement::Unplaced),
other => Err(E::custom(format!("unknown namespace placement: {other}"))),
}
}
fn visit_map<A>(self, mut map: A) -> Result<Self::Value, A::Error>
where
A: serde::de::MapAccess<'de>,
{
let Some(tag) = map.next_key::<String>()? else {
return Err(serde::de::Error::custom(
"empty namespace placement map: expected one of {prefer|pinned: [labels]}",
));
};
let placement = match tag.as_str() {
PLACEMENT_PREFER_TAG => NamespacePlacement::Prefer {
nodes: map.next_value()?,
},
PLACEMENT_PINNED_TAG => NamespacePlacement::Pinned {
nodes: map.next_value()?,
},
other => {
return Err(serde::de::Error::custom(format!(
"unknown namespace placement: {other}"
)));
}
};
if let Some(extra) = map.next_key::<String>()? {
return Err(serde::de::Error::custom(format!(
"unexpected extra namespace placement key: {extra}"
)));
}
Ok(placement)
}
}
fn encode_instant(instant: DateTime<Utc>) -> String {
instant.to_rfc3339_opts(SecondsFormat::Nanos, true)
}
fn decode_instant(value: &str) -> Result<DateTime<Utc>, StoreError> {
DateTime::parse_from_rfc3339(value)
.map(|date_time| date_time.with_timezone(&Utc))
.map_err(|error| StoreError::Serialization(error.to_string()))
}
#[async_trait]
pub trait NamespaceStore: Send + Sync + 'static {
async fn register_namespace(
&self,
name: &str,
origin: NamespaceOrigin,
) -> Result<MintOutcome, StoreError>;
async fn put_namespace(&self, record: NamespaceRecord) -> Result<MintOutcome, StoreError>;
async fn list_namespaces(&self) -> Result<Vec<NamespaceRecord>, StoreError>;
async fn get_namespace(&self, name: &str) -> Result<Option<NamespaceRecord>, StoreError>;
async fn set_namespace_placement(
&self,
name: &str,
placement: NamespacePlacement,
) -> Result<Option<()>, StoreError>;
async fn deprecate_namespace(&self, name: &str) -> Result<(), StoreError>;
}
#[cfg(test)]
mod tests {
#![allow(clippy::expect_used)]
use super::{
MintOutcome, NamespaceConfig, NamespaceOrigin, NamespacePlacement, NamespaceRecord,
NamespaceState, StoredNamespace,
};
use chrono::{Duration, TimeZone, Utc};
use std::collections::BTreeSet;
fn node_set(labels: &[&str]) -> BTreeSet<String> {
labels.iter().map(|label| (*label).to_owned()).collect()
}
fn fixed_now() -> chrono::DateTime<Utc> {
match Utc.with_ymd_and_hms(2026, 6, 30, 12, 0, 0).single() {
Some(instant) => instant,
None => Utc::now(),
}
}
#[test]
fn new_minted_sets_created_equal_to_last_seen_and_origin() {
let now = fixed_now();
let record = NamespaceRecord::new_minted("orders", NamespaceOrigin::WorkerMint, now);
assert_eq!(record.name, "orders");
assert_eq!(record.created_at, now);
assert_eq!(record.last_seen, now);
assert_eq!(record.created_at, record.last_seen);
assert_eq!(record.origin, NamespaceOrigin::WorkerMint);
assert_eq!(record.state, NamespaceState::Active);
assert_eq!(record.config, NamespaceConfig::default());
assert_eq!(record.placement, NamespacePlacement::Unplaced);
assert_eq!(record.config.kind, None);
assert_eq!(record.config.max_in_flight_activities, None);
}
#[test]
fn bump_last_seen_advances_only_last_seen() {
let now = fixed_now();
let mut record = NamespaceRecord::new_minted("orders", NamespaceOrigin::Explicit, now);
let later = now + Duration::seconds(42);
record.bump_last_seen(later);
assert_eq!(record.last_seen, later);
assert_eq!(record.created_at, now);
assert_eq!(record.origin, NamespaceOrigin::Explicit);
assert_eq!(record.state, NamespaceState::Active);
}
#[test]
fn encode_decode_round_trips() {
let now = fixed_now();
let mut record =
NamespaceRecord::new_minted("billing", NamespaceOrigin::InferredFromState, now);
record.bump_last_seen(now + Duration::seconds(5));
record.state = NamespaceState::Deprecated;
let bytes = record.encode().expect("encode");
let decoded = NamespaceRecord::decode(&bytes).expect("decode");
assert_eq!(record, decoded);
}
#[test]
fn encode_decode_preserves_reserved_kind_discriminator() {
let now = fixed_now();
let mut record = NamespaceRecord::new_minted("tenant-a", NamespaceOrigin::Explicit, now);
record.config.kind = Some("tenant".to_owned());
let bytes = record.encode().expect("encode");
let decoded = NamespaceRecord::decode(&bytes).expect("decode");
assert_eq!(decoded.config.kind.as_deref(), Some("tenant"));
assert_eq!(record, decoded);
}
#[test]
fn decode_rejects_malformed_bytes() {
let err = NamespaceRecord::decode(b"not json").expect_err("must reject");
assert!(matches!(err, crate::StoreError::Serialization(_)));
}
#[test]
fn enum_derives_are_copy_clone_eq() {
let outcome = MintOutcome::Created;
let copied = outcome;
assert_eq!(outcome, copied);
assert_eq!(copied, MintOutcome::Created);
assert_ne!(MintOutcome::Created, MintOutcome::AlreadyExisted);
let origin = NamespaceOrigin::WorkerMint;
let origin_copy = origin;
assert_eq!(origin, origin_copy);
let state = NamespaceState::Active;
let state_copy = state;
assert_eq!(state, state_copy);
assert_ne!(NamespaceState::Active, NamespaceState::Deprecated);
}
#[test]
fn placement_variants_round_trip_through_record() {
let now = fixed_now();
for placement in [
NamespacePlacement::Unplaced,
NamespacePlacement::Prefer {
nodes: node_set(&["az-a", "az-b"]),
},
NamespacePlacement::Pinned {
nodes: node_set(&["gpu-pool"]),
},
] {
let mut record = NamespaceRecord::new_minted("placed", NamespaceOrigin::Explicit, now);
record.placement = placement.clone();
let bytes = record.encode().expect("encode");
let decoded = NamespaceRecord::decode(&bytes).expect("decode");
assert_eq!(decoded.placement, placement);
assert_eq!(record, decoded);
}
}
#[test]
fn unplaced_encodes_as_bare_string_tag() {
let placement = NamespacePlacement::Unplaced;
let json = serde_json::to_string(&placement).expect("serialize");
assert_eq!(json, "\"unplaced\"");
}
#[test]
fn prefer_and_pinned_encode_as_single_key_maps() {
let prefer = NamespacePlacement::Prefer {
nodes: node_set(&["b", "a"]),
};
let pinned = NamespacePlacement::Pinned {
nodes: node_set(&["only"]),
};
assert_eq!(
serde_json::to_string(&prefer).expect("serialize prefer"),
r#"{"prefer":["a","b"]}"#
);
assert_eq!(
serde_json::to_string(&pinned).expect("serialize pinned"),
r#"{"pinned":["only"]}"#
);
}
#[test]
fn old_unplaced_record_without_quota_field_decodes() {
let old = StoredNamespace {
name: "legacy".to_owned(),
created_at: super::encode_instant(fixed_now()),
last_seen: super::encode_instant(fixed_now()),
origin: NamespaceOrigin::WorkerMint,
kind: None,
max_in_flight_activities: None,
placement: NamespacePlacement::Unplaced,
state: NamespaceState::Active,
};
let old_json = format!(
r#"{{"name":"legacy","created_at":"{}","last_seen":"{}","origin":"WorkerMint","kind":null,"placement":"unplaced","state":"Active"}}"#,
old.created_at, old.last_seen
);
let decoded = NamespaceRecord::decode(old_json.as_bytes()).expect("decode old bytes");
assert_eq!(decoded.placement, NamespacePlacement::Unplaced);
assert_eq!(decoded.config.max_in_flight_activities, None);
assert_eq!(decoded.config.kind, None);
assert_eq!(decoded.name, "legacy");
assert_eq!(decoded.origin, NamespaceOrigin::WorkerMint);
assert_eq!(decoded.state, NamespaceState::Active);
}
#[test]
fn max_in_flight_activities_round_trips_present_and_absent() {
let now = fixed_now();
let mut with_quota = NamespaceRecord::new_minted("capped", NamespaceOrigin::Explicit, now);
with_quota.config.max_in_flight_activities = Some(256);
let decoded =
NamespaceRecord::decode(&with_quota.encode().expect("encode")).expect("decode");
assert_eq!(decoded.config.max_in_flight_activities, Some(256));
assert_eq!(with_quota, decoded);
let without = NamespaceRecord::new_minted("uncapped", NamespaceOrigin::Explicit, now);
let decoded = NamespaceRecord::decode(&without.encode().expect("encode")).expect("decode");
assert_eq!(decoded.config.max_in_flight_activities, None);
assert_eq!(without, decoded);
}
#[test]
fn unknown_placement_tag_is_rejected() {
let bad_string: Result<NamespacePlacement, _> = serde_json::from_str("\"elsewhere\"");
assert!(bad_string.is_err());
let bad_map: Result<NamespacePlacement, _> = serde_json::from_str(r#"{"banish":["x"]}"#);
assert!(bad_map.is_err());
}
}