Skip to main content

cortexkit_bus_naming/
names.rs

1use std::error::Error;
2use std::fmt;
3use std::num::NonZeroU32;
4
5use crate::token::{validate_account_token, validate_token, NamingError, TokenKind};
6
7/// All six stream names owned by one account.
8#[derive(Debug, Clone, PartialEq, Eq)]
9pub struct StreamNames {
10    pub room: String,
11    pub wake: String,
12    pub peer: String,
13    pub effect: String,
14    pub effect_dead: String,
15    /// Module events (`ck.{acct}.event.>`), read by the flow engine.
16    pub event: String,
17}
18
19impl StreamNames {
20    pub fn all(&self) -> [&str; 6] {
21        [
22            &self.room,
23            &self.wake,
24            &self.peer,
25            &self.effect,
26            &self.effect_dead,
27            &self.event,
28        ]
29    }
30
31    /// The streams that hold per-agent durables (`c_{agent_id}`): WAKE, PEER, EFFECT.
32    ///
33    /// Grants that let a credential read any agent's durable use a whole-token `*`
34    /// in the consumer position on exactly these streams (NATS wildcards cannot
35    /// match a `c_` prefix), which admits every consumer on them. So no non-agent
36    /// durable may ever be created on these streams; ck-bus asserts that against
37    /// this list. ROOM, EFFECT_DEAD and EVENT are not agent streams: EVENT and
38    /// ROOM hold module durables (`m_{module_id}`), never a `c_` durable.
39    pub fn agent_streams(&self) -> [&str; 3] {
40        [&self.wake, &self.peer, &self.effect]
41    }
42}
43
44/// The census KV bucket and its JetStream backing stream.
45#[derive(Debug, Clone, PartialEq, Eq)]
46pub struct BucketNames {
47    pub census: String,
48    pub census_stream: String,
49}
50
51impl BucketNames {
52    pub fn all(&self) -> [&str; 2] {
53        [&self.census, &self.census_stream]
54    }
55}
56
57/// Account-derived names. Derivation never normalizes input.
58#[derive(Debug, Clone, PartialEq, Eq)]
59pub struct AccountNames {
60    account: String,
61    account_upper: String,
62    streams: StreamNames,
63    buckets: BucketNames,
64}
65
66impl AccountNames {
67    pub fn derive(account: &str) -> Result<Self, NamingError> {
68        validate_account_token(account)?;
69        let account_upper = account.to_ascii_uppercase();
70        Ok(Self {
71            account: account.to_owned(),
72            streams: StreamNames {
73                room: format!("CK_{account_upper}_ROOM"),
74                wake: format!("CK_{account_upper}_WAKE"),
75                peer: format!("CK_{account_upper}_PEER"),
76                effect: format!("CK_{account_upper}_EFFECT"),
77                effect_dead: format!("CK_{account_upper}_EFFECT_DEAD"),
78                event: format!("CK_{account_upper}_EVENT"),
79            },
80            buckets: BucketNames {
81                census: format!("CK_{account_upper}_CENSUS"),
82                census_stream: format!("KV_CK_{account_upper}_CENSUS"),
83            },
84            account_upper,
85        })
86    }
87
88    pub fn from_roster_host_id(roster_host_id: &str) -> Result<Self, NamingError> {
89        validate_token(TokenKind::RosterHostId, roster_host_id)?;
90        Self::derive(&format!("box_{roster_host_id}"))
91    }
92
93    pub fn account(&self) -> &str {
94        &self.account
95    }
96
97    pub fn account_upper(&self) -> &str {
98        &self.account_upper
99    }
100
101    pub fn streams(&self) -> &StreamNames {
102        &self.streams
103    }
104
105    pub fn buckets(&self) -> &BucketNames {
106        &self.buckets
107    }
108
109    pub fn room_post(&self, room_id: &str) -> Result<String, NamingError> {
110        validate_token(TokenKind::RoomId, room_id)?;
111        Ok(format!("ck.{}.room.{room_id}.post", self.account))
112    }
113
114    pub fn room_subscription(&self, room_id: &str) -> Result<String, NamingError> {
115        validate_token(TokenKind::RoomId, room_id)?;
116        Ok(format!("ck.{}.room.{room_id}.>", self.account))
117    }
118
119    pub fn room_binding(&self) -> String {
120        format!("ck.{}.room.*.post", self.account)
121    }
122
123    pub fn wake_fire(&self, agent_id: &str) -> Result<String, NamingError> {
124        validate_token(TokenKind::AgentId, agent_id)?;
125        Ok(format!("ck.{}.wake.{agent_id}.fire", self.account))
126    }
127
128    pub fn wake_binding(&self) -> String {
129        format!("ck.{}.wake.*.fire", self.account)
130    }
131
132    pub fn peer_delivery(&self, agent_id: &str, session_id: &str) -> Result<String, NamingError> {
133        validate_token(TokenKind::AgentId, agent_id)?;
134        validate_token(TokenKind::SessionId, session_id)?;
135        Ok(format!(
136            "ck.{}.peer.{agent_id}.{session_id}.deliver",
137            self.account
138        ))
139    }
140
141    pub fn peer_filter(&self, agent_id: &str) -> Result<String, NamingError> {
142        validate_token(TokenKind::AgentId, agent_id)?;
143        Ok(format!("ck.{}.peer.{agent_id}.*.deliver", self.account))
144    }
145
146    pub fn peer_subscription(&self, agent_id: &str) -> Result<String, NamingError> {
147        validate_token(TokenKind::AgentId, agent_id)?;
148        Ok(format!("ck.{}.peer.{agent_id}.>", self.account))
149    }
150
151    pub fn peer_binding(&self) -> String {
152        format!("ck.{}.peer.*.*.deliver", self.account)
153    }
154
155    pub fn effect_intent(&self, agent_id: &str, session_id: &str) -> Result<String, NamingError> {
156        validate_token(TokenKind::AgentId, agent_id)?;
157        validate_token(TokenKind::SessionId, session_id)?;
158        Ok(format!(
159            "ck.{}.effect.{agent_id}.{session_id}.intent",
160            self.account
161        ))
162    }
163
164    pub fn effect_filter(&self, agent_id: &str) -> Result<String, NamingError> {
165        validate_token(TokenKind::AgentId, agent_id)?;
166        Ok(format!("ck.{}.effect.{agent_id}.*.intent", self.account))
167    }
168
169    pub fn effect_publish_grant(&self, agent_id: &str) -> Result<String, NamingError> {
170        validate_token(TokenKind::AgentId, agent_id)?;
171        Ok(format!("ck.{}.effect.{agent_id}.>", self.account))
172    }
173
174    pub fn effect_binding(&self) -> String {
175        format!("ck.{}.effect.*.*.intent", self.account)
176    }
177
178    pub fn effect_dead(&self) -> String {
179        format!("ck.{}.effect.dead", self.account)
180    }
181
182    /// The subject one module publishes one event on:
183    /// `ck.{acct}.event.{module_id}.{event}.v{version}`.
184    ///
185    /// The event name must be a single subject token and the version at least
186    /// 1, so every event subject has exactly six tokens and the version is
187    /// always the last one.
188    pub fn event_subject(
189        &self,
190        module_id: &str,
191        event: &str,
192        version: u32,
193    ) -> Result<String, NamingError> {
194        validate_token(TokenKind::ModuleId, module_id)?;
195        validate_token(TokenKind::EventName, event)?;
196        let version_token = format!("v{version}");
197        validate_token(TokenKind::EventVersion, &version_token)?;
198        Ok(format!(
199            "ck.{}.event.{module_id}.{event}.{version_token}",
200            self.account
201        ))
202    }
203
204    /// The publish grant covering every event subject of one module.
205    pub fn event_publish_grant(&self, module_id: &str) -> Result<String, NamingError> {
206        validate_token(TokenKind::ModuleId, module_id)?;
207        Ok(format!("ck.{}.event.{module_id}.>", self.account))
208    }
209
210    pub fn event_binding(&self) -> String {
211        format!("ck.{}.event.>", self.account)
212    }
213
214    pub fn sentinel_ping(&self) -> String {
215        format!("ck.{}.sentinel.ping", self.account)
216    }
217
218    pub fn consumer_name(agent_id: &str) -> Result<String, NamingError> {
219        validate_token(TokenKind::AgentId, agent_id)?;
220        Ok(format!("c_{agent_id}"))
221    }
222
223    /// The durable a module reads a non-agent stream through, `m_{module_id}`.
224    ///
225    /// The `m_` prefix keeps module durables apart from the agent namespace
226    /// `c_`. They live on the event stream (one per flow-engine module) and on
227    /// the ROOM stream (the delivery authority's one durable, filtered on
228    /// `room_binding`), never on an agent stream.
229    pub fn module_consumer_name(module_id: &str) -> Result<String, NamingError> {
230        validate_token(TokenKind::ModuleId, module_id)?;
231        Ok(format!("m_{module_id}"))
232    }
233
234    /// The census key for one live module inside `buckets().census`.
235    ///
236    /// Issuance writes one record per module under this key (history 1, no TTL)
237    /// so revocation can find the credential a module currently holds. The key
238    /// is the module id itself, refused unless it is a single plain token: a dot
239    /// or wildcard would address a different key or several at once.
240    pub fn census_key(module_id: &str) -> Result<String, NamingError> {
241        validate_token(TokenKind::ModuleId, module_id)?;
242        Ok(module_id.to_owned())
243    }
244
245    /// The KV subject a census write for `module_id` is published on.
246    pub fn census_subject(&self, module_id: &str) -> Result<String, NamingError> {
247        let key = Self::census_key(module_id)?;
248        Ok(format!("$KV.{}.{key}", self.buckets().census))
249    }
250
251    #[deprecated(
252        note = "the foundation amendment keeps no process records in the vault; per-spawn credentials are held by ck-bus in memory"
253    )]
254    pub fn process_record_name(
255        module_id: &str,
256        generation: u64,
257        epoch: u64,
258    ) -> Result<String, NamingError> {
259        validate_token(TokenKind::ModuleId, module_id)?;
260        Ok(format!("nats.{module_id}.g{generation}.e{epoch}"))
261    }
262
263    #[deprecated(
264        note = "vault roots are credential ids from the operator ceremony; use root_credential_id"
265    )]
266    pub fn leaf_record_name(roster_host_id: &str) -> Result<String, NamingError> {
267        validate_token(TokenKind::RosterHostId, roster_host_id)?;
268        Ok(format!("nats.leaf.{roster_host_id}"))
269    }
270
271    #[deprecated(
272        note = "vault roots are credential ids from the operator ceremony; use root_credential_id"
273    )]
274    pub fn operator_record_name(&self) -> String {
275        format!("nats.operator.{}", self.account)
276    }
277
278    #[deprecated(
279        note = "vault roots are credential ids from the operator ceremony; use root_credential_id"
280    )]
281    pub fn account_record_name(&self) -> String {
282        format!("nats.account.{}", self.account)
283    }
284
285    #[deprecated(
286        note = "vault roots are credential ids from the operator ceremony; use root_credential_id"
287    )]
288    pub fn system_account_record_name(&self) -> String {
289        format!("nats.sysaccount.{}", self.account)
290    }
291}
292
293/// Families of root keys the operator creates once in the vault by ceremony
294/// (`ck auth mint-signing-key --id signing:<provider>[:<generation>]`).
295#[derive(Debug, Clone, Copy, PartialEq, Eq)]
296pub enum RootCredentialKind {
297    /// Ed25519 signing keys: account, operator and message-signing roots.
298    Signing,
299    /// Key-encapsulation keys, for sealing federation envelopes.
300    Kem,
301}
302
303impl RootCredentialKind {
304    pub const fn prefix(self) -> &'static str {
305        match self {
306            Self::Signing => "signing",
307            Self::Kem => "kem",
308        }
309    }
310}
311
312/// The vault credential id of a root key: `<kind>:<provider>[:<generation>]`,
313/// for example `signing:ck-bus-account:1` or `signing:msgsig`.
314///
315/// `provider` uses the shared identity lexicon, so it can never carry a colon
316/// and add a segment; a generation, when present, is at least 1 by type.
317pub fn root_credential_id(
318    kind: RootCredentialKind,
319    provider: &str,
320    generation: Option<NonZeroU32>,
321) -> Result<String, NamingError> {
322    validate_token(TokenKind::RootProvider, provider)?;
323    Ok(match generation {
324        Some(generation) => format!("{}:{provider}:{generation}", kind.prefix()),
325        None => format!("{}:{provider}", kind.prefix()),
326    })
327}
328
329/// Namespaces that are deliberately outside the account-token rule.
330#[derive(Debug, Clone, Copy, PartialEq, Eq)]
331pub enum NamingExemption {
332    InboxPrefix,
333    JetStreamPrefix,
334    KeyValuePrefix,
335    SystemPrefix,
336    DurableConsumer,
337}
338
339/// The exemption list is closed and deliberately exposed for exhaustive tests.
340pub const CLOSED_NAMING_EXEMPTIONS: [NamingExemption; 5] = [
341    NamingExemption::InboxPrefix,
342    NamingExemption::JetStreamPrefix,
343    NamingExemption::KeyValuePrefix,
344    NamingExemption::SystemPrefix,
345    NamingExemption::DurableConsumer,
346];
347
348#[derive(Debug, Clone, Copy, PartialEq, Eq)]
349pub enum NamingRule {
350    CkSubject,
351    StreamName,
352    BucketName,
353    Exempt(NamingExemption),
354}
355
356#[derive(Debug, Clone, PartialEq, Eq)]
357pub struct TenancyNameError {
358    name: String,
359    rule: NamingRule,
360}
361
362impl TenancyNameError {
363    pub fn name(&self) -> &str {
364        &self.name
365    }
366
367    pub fn rule(&self) -> NamingRule {
368        self.rule
369    }
370}
371
372impl fmt::Display for TenancyNameError {
373    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
374        write!(
375            f,
376            "name {:?} does not satisfy tenancy naming rule {:?}",
377            self.name, self.rule
378        )
379    }
380}
381
382impl Error for TenancyNameError {}
383
384/// Checks the account-token rule for one emitted name or a named exemption.
385pub fn validate_tenancy_name(
386    account: &AccountNames,
387    name: &str,
388    rule: NamingRule,
389) -> Result<(), TenancyNameError> {
390    let valid = match rule {
391        NamingRule::CkSubject => name.starts_with(&format!("ck.{}.", account.account)),
392        NamingRule::StreamName => name.starts_with(&format!("CK_{}_", account.account_upper)),
393        NamingRule::BucketName => {
394            name.starts_with(&format!("CK_{}_", account.account_upper))
395                || name.starts_with(&format!("KV_CK_{}_", account.account_upper))
396        }
397        NamingRule::Exempt(NamingExemption::InboxPrefix) => name.starts_with("_INBOX."),
398        NamingRule::Exempt(NamingExemption::JetStreamPrefix) => name.starts_with("$JS."),
399        NamingRule::Exempt(NamingExemption::KeyValuePrefix) => name.starts_with("$KV."),
400        NamingRule::Exempt(NamingExemption::SystemPrefix) => name.starts_with("$SYS."),
401        NamingRule::Exempt(NamingExemption::DurableConsumer) => name.starts_with("c_"),
402    };
403    if valid {
404        Ok(())
405    } else {
406        Err(TenancyNameError {
407            name: name.to_owned(),
408            rule,
409        })
410    }
411}