Skip to main content

cortexkit_bus_naming/
grants.rs

1use std::collections::BTreeSet;
2use std::error::Error;
3use std::fmt;
4
5use crate::names::AccountNames;
6use crate::token::{validate_token, NamingError, TokenKind};
7
8#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
9pub enum Principal {
10    Participant,
11    DeliveryAuthority,
12    Bus,
13    System,
14    /// The module that runs flows (basal): it reads the event stream.
15    FlowEngine,
16}
17
18impl Principal {
19    pub const fn as_str(self) -> &'static str {
20        match self {
21            Self::Participant => "participant",
22            Self::DeliveryAuthority => "delivery-authority",
23            Self::Bus => "bus",
24            Self::System => "system",
25            Self::FlowEngine => "flow-engine",
26        }
27    }
28}
29
30#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
31pub enum Operation {
32    Publish,
33    Subscribe,
34}
35
36impl Operation {
37    pub const fn as_str(self) -> &'static str {
38        match self {
39            Self::Publish => "publish",
40            Self::Subscribe => "subscribe",
41        }
42    }
43}
44
45/// One permission entry. The model intentionally has no deny variant.
46#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
47pub struct AllowEntry {
48    pub principal: Principal,
49    pub operation: Operation,
50    pub subject: String,
51}
52
53#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
54pub struct RefusedEntry {
55    pub principal: Principal,
56    pub operation: Operation,
57    pub subject: String,
58    pub reason: &'static str,
59}
60
61#[derive(Debug, Clone, Copy, PartialEq, Eq)]
62pub struct GoldenFixture {
63    pub account: &'static str,
64    pub module_id: &'static str,
65    pub bound_agent: &'static str,
66    pub foreign_agent: &'static str,
67    pub bound_room: &'static str,
68    pub unbound_room: &'static str,
69    /// Inbox component of the system user in the golden.
70    pub system_credential: &'static str,
71    /// Module id, and inbox component, of the flow engine in the golden.
72    pub flow_engine_module: &'static str,
73    /// A module that is neither the participant nor the flow engine; the
74    /// golden refuses its event subjects and its durable to both.
75    pub foreign_module: &'static str,
76    /// Module id, and inbox component, of the delivery authority in the
77    /// golden; its ROOM durable is `m_` followed by this id.
78    pub delivery_authority_module: &'static str,
79}
80
81pub const PINNED_GOLDEN_FIXTURE: GoldenFixture = GoldenFixture {
82    account: "box_goldenfixture",
83    module_id: "ckbus",
84    bound_agent: "agent_gold_a",
85    foreign_agent: "agent_gold_b",
86    bound_room: "room_gold_bound",
87    unbound_room: "room_gold_unbound",
88    system_credential: "cksys",
89    flow_engine_module: "basal",
90    foreign_module: "other",
91    delivery_authority_module: "prefrontal-core",
92};
93
94/// The session token in the golden's refused peer and effect subjects.
95const GOLDEN_SESSION: &str = "session_gold";
96
97/// The event name and version in the golden's refused event subjects.
98const GOLDEN_EVENT: &str = "pull_request_review";
99const GOLDEN_EVENT_VERSION: u32 = 1;
100
101#[derive(Debug, Clone, PartialEq, Eq)]
102pub struct PermissionDocument {
103    fixture: GoldenFixture,
104    allows: Vec<AllowEntry>,
105    refused: Vec<RefusedEntry>,
106}
107
108impl PermissionDocument {
109    pub fn allows(&self) -> &[AllowEntry] {
110        &self.allows
111    }
112
113    pub fn refused(&self) -> &[RefusedEntry] {
114        &self.refused
115    }
116
117    pub fn render(&self) -> String {
118        let mut output = String::from(
119            "# Generated by cortexkit-bus-naming; edit the generator, not this file.\n",
120        );
121        for (name, value) in [
122            ("account", self.fixture.account),
123            ("module", self.fixture.module_id),
124            ("bound-agent", self.fixture.bound_agent),
125            ("foreign-agent", self.fixture.foreign_agent),
126            ("bound-room", self.fixture.bound_room),
127            ("unbound-room", self.fixture.unbound_room),
128            ("flow-engine-module", self.fixture.flow_engine_module),
129            ("foreign-module", self.fixture.foreign_module),
130            (
131                "delivery-authority-module",
132                self.fixture.delivery_authority_module,
133            ),
134        ] {
135            output.push_str(&format!("fixture {name} {value}\n"));
136        }
137        output.push('\n');
138        for entry in &self.allows {
139            output.push_str(&format!(
140                "allow {} {} {}\n",
141                entry.principal.as_str(),
142                entry.operation.as_str(),
143                entry.subject
144            ));
145        }
146        output.push('\n');
147        for entry in &self.refused {
148            output.push_str(&format!(
149                "expect-refused {} {} {} {}\n",
150                entry.principal.as_str(),
151                entry.operation.as_str(),
152                entry.subject,
153                entry.reason
154            ));
155        }
156        output
157    }
158}
159
160#[derive(Debug, Clone, PartialEq, Eq)]
161pub enum GrantError {
162    Naming(NamingError),
163    InvalidPermissionFile { line: usize, reason: String },
164    MissingStreamExpansion { stream: String },
165}
166
167impl fmt::Display for GrantError {
168    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
169        match self {
170            Self::Naming(error) => error.fmt(f),
171            Self::InvalidPermissionFile { line, reason } => {
172                write!(f, "invalid permission file at line {line}: {reason}")
173            }
174            Self::MissingStreamExpansion { stream } => {
175                write!(f, "permission file does not expand literal stream {stream}")
176            }
177        }
178    }
179}
180
181impl Error for GrantError {
182    fn source(&self) -> Option<&(dyn Error + 'static)> {
183        match self {
184            Self::Naming(error) => Some(error),
185            _ => None,
186        }
187    }
188}
189
190impl From<NamingError> for GrantError {
191    fn from(value: NamingError) -> Self {
192        Self::Naming(value)
193    }
194}
195
196/// The grant for every host and module except the delivery authority.
197///
198/// Credentials are account-scoped and name no agents: a participant may pull,
199/// inspect and acknowledge any agent's durable on the agent streams, and which
200/// durable it actually reads is decided by prefrontal-core, which owns agent
201/// residence. It holds no workload publish at all: hosts and modules consume
202/// wakes, peer deliveries and effect intents but never produce them.
203///
204/// It may publish its own module events, on `ck.{acct}.event.{module_id}.>`
205/// and on no other module's event subjects.
206pub fn participant_permissions(
207    account: &AccountNames,
208    credential_public: &str,
209    module_id: &str,
210    bound_rooms: &[&str],
211) -> Result<Vec<AllowEntry>, GrantError> {
212    let mut entries = BTreeSet::new();
213    add_consumer_permissions(
214        &mut entries,
215        Principal::Participant,
216        account,
217        credential_public,
218        module_id,
219        bound_rooms,
220    )?;
221    Ok(entries.into_iter().collect())
222}
223
224/// The grant for the delivery authority, prefrontal-core alone.
225///
226/// Everything a participant holds, plus publish on each agent stream's own
227/// binding: prefrontal-core is the only producer of wakes, peer deliveries and
228/// effect intents, so it publishes every agent's wake fires and peer deliveries.
229/// Keeping workload publish here, and out of every
230/// host credential and out of ck-bus's own, means the module that decides a
231/// delivery is the only one that can make one.
232///
233/// Rooms follow the same shape. prefrontal-core is the only producer and the
234/// only consumer of room posts (hosts never consume rooms), so it publishes on
235/// the ROOM binding `ck.{acct}.room.*.post` rather than per bound room, which
236/// would re-mint its credential on every room create. It reads every room
237/// through one module durable, `m_{module_id}` on the ROOM stream, filtered on
238/// that binding and created by ck-bus: pull, ack and consumer info are granted
239/// on that one durable by name, never with a wildcard consumer token, so it can
240/// read no other consumer on ROOM, and it holds no consumer create or delete
241/// there.
242pub fn delivery_authority_permissions(
243    account: &AccountNames,
244    credential_public: &str,
245    module_id: &str,
246    bound_rooms: &[&str],
247) -> Result<Vec<AllowEntry>, GrantError> {
248    let mut entries = BTreeSet::new();
249    add_consumer_permissions(
250        &mut entries,
251        Principal::DeliveryAuthority,
252        account,
253        credential_public,
254        module_id,
255        bound_rooms,
256    )?;
257    let room_stream = &account.streams().room;
258    let room_durable = AccountNames::module_consumer_name(module_id)?;
259    for subject in [
260        account.wake_binding(),
261        account.peer_binding(),
262        account.effect_binding(),
263        account.room_binding(),
264        format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{room_durable}"),
265        format!("$JS.API.CONSUMER.INFO.{room_stream}.{room_durable}"),
266        format!("$JS.ACK.{room_stream}.{room_durable}.>"),
267        format!("$JS.API.STREAM.INFO.{room_stream}"),
268    ] {
269        add(
270            &mut entries,
271            Principal::DeliveryAuthority,
272            Operation::Publish,
273            subject,
274        );
275    }
276    Ok(entries.into_iter().collect())
277}
278
279/// The permissions participants and the delivery authority share: their own
280/// inbox, the dead-letter record, census reads, bound rooms, publishing their
281/// own module events, and reading any agent's durable on the agent streams.
282fn add_consumer_permissions(
283    entries: &mut BTreeSet<AllowEntry>,
284    principal: Principal,
285    account: &AccountNames,
286    credential_public: &str,
287    module_id: &str,
288    bound_rooms: &[&str],
289) -> Result<(), GrantError> {
290    validate_token(TokenKind::CredentialPublic, credential_public)?;
291
292    add(
293        entries,
294        principal,
295        Operation::Subscribe,
296        format!("_INBOX.{credential_public}.>"),
297    );
298    add(
299        entries,
300        principal,
301        Operation::Publish,
302        account.effect_dead(),
303    );
304    add(
305        entries,
306        principal,
307        Operation::Publish,
308        account.event_publish_grant(module_id)?,
309    );
310
311    // A whole-token `*` in the consumer position: NATS wildcards cannot match the
312    // `c_` prefix, so this admits every consumer on the stream. That is safe only
313    // because no non-agent durable exists on the agent streams (see
314    // `StreamNames::agent_streams`).
315    for stream in account.streams().agent_streams() {
316        for subject in [
317            format!("$JS.ACK.{stream}.*.>"),
318            format!("$JS.API.CONSUMER.MSG.NEXT.{stream}.*"),
319            format!("$JS.API.CONSUMER.INFO.{stream}.*"),
320            format!("$JS.API.STREAM.INFO.{stream}"),
321        ] {
322            add(entries, principal, Operation::Publish, subject);
323        }
324    }
325
326    for room_id in bound_rooms {
327        validate_token(TokenKind::RoomId, room_id)?;
328        add(
329            entries,
330            principal,
331            Operation::Publish,
332            account.room_post(room_id)?,
333        );
334        add(
335            entries,
336            principal,
337            Operation::Subscribe,
338            account.room_subscription(room_id)?,
339        );
340    }
341
342    add_census_read_permissions(entries, principal, account);
343    Ok(())
344}
345
346pub fn bus_permissions(
347    account: &AccountNames,
348    credential_public: &str,
349) -> Result<Vec<AllowEntry>, GrantError> {
350    validate_token(TokenKind::CredentialPublic, credential_public)?;
351    let mut entries = BTreeSet::new();
352    let inbox = format!("_INBOX.{credential_public}.>");
353    for operation in [Operation::Publish, Operation::Subscribe] {
354        add(&mut entries, Principal::Bus, operation, inbox.clone());
355        add(
356            &mut entries,
357            Principal::Bus,
358            operation,
359            account.sentinel_ping(),
360        );
361        add(
362            &mut entries,
363            Principal::Bus,
364            operation,
365            format!("$KV.{}.>", account.buckets().census),
366        );
367    }
368    add_census_read_permissions(&mut entries, Principal::Bus, account);
369
370    let census_stream = &account.buckets().census_stream;
371    for subject in [
372        format!("$JS.API.STREAM.CREATE.{census_stream}"),
373        format!("$JS.API.STREAM.UPDATE.{census_stream}"),
374    ] {
375        add(&mut entries, Principal::Bus, Operation::Publish, subject);
376    }
377
378    // Stream and consumer management on every stream, the event stream
379    // included: the consumer-create rows are how ck-bus creates the module
380    // durables `m_{module_id}` on the event and ROOM streams, as it creates the
381    // agents' `c_` durables on the agent streams.
382    for stream in account.streams().all() {
383        for subject in [
384            format!("$JS.API.STREAM.CREATE.{stream}"),
385            format!("$JS.API.STREAM.UPDATE.{stream}"),
386            format!("$JS.API.STREAM.INFO.{stream}"),
387            format!("$JS.API.STREAM.DELETE.{stream}"),
388            format!("$JS.API.CONSUMER.CREATE.{stream}.>"),
389            format!("$JS.API.CONSUMER.DURABLE.CREATE.{stream}.>"),
390            format!("$JS.API.CONSUMER.INFO.{stream}.>"),
391            format!("$JS.API.CONSUMER.DELETE.{stream}.>"),
392        ] {
393            add(&mut entries, Principal::Bus, Operation::Publish, subject);
394        }
395    }
396
397    let dead_stream = &account.streams().effect_dead;
398    for subject in [
399        format!("$JS.ACK.{dead_stream}.c_ckbus_dead.>"),
400        format!("$JS.API.CONSUMER.MSG.NEXT.{dead_stream}.c_ckbus_dead"),
401        format!("$JS.API.CONSUMER.INFO.{dead_stream}.c_ckbus_dead"),
402    ] {
403        add(&mut entries, Principal::Bus, Operation::Publish, subject);
404    }
405
406    // Agent durable lifecycle: listing the durables, and purging an agent's
407    // subject when its durable is deleted (a consumer delete alone leaves its
408    // messages in the stream for a later bind to redeliver). The purge filter
409    // travels in the request body, which permissions cannot see, so ck-bus's
410    // code restricts each purge to one agent's filter.
411    for stream in account.streams().agent_streams() {
412        for subject in [
413            format!("$JS.API.CONSUMER.NAMES.{stream}"),
414            format!("$JS.API.STREAM.PURGE.{stream}"),
415        ] {
416            add(&mut entries, Principal::Bus, Operation::Publish, subject);
417        }
418    }
419    Ok(entries.into_iter().collect())
420}
421
422/// The system user's grant. It publishes claims updates, per-account claims
423/// lookups (`$SYS.REQ.ACCOUNT.<account>.CLAIMS.LOOKUP`, which the revocation
424/// check reads back) and the claims list (`$SYS.REQ.CLAIMS.LIST`, which finds
425/// an existing account by name when its id was not recorded, so a lost state
426/// file adopts the account instead of creating a second one), kicks, and
427/// watches connects; replies to its requests arrive on its own
428/// credential-scoped inbox.
429pub fn system_permissions(credential_public: &str) -> Result<Vec<AllowEntry>, GrantError> {
430    validate_token(TokenKind::CredentialPublic, credential_public)?;
431    let mut entries = BTreeSet::new();
432    for subject in [
433        "$SYS.REQ.CLAIMS.UPDATE",
434        "$SYS.REQ.CLAIMS.LIST",
435        "$SYS.REQ.ACCOUNT.*.CLAIMS.LOOKUP",
436        "$SYS.REQ.SERVER.*.KICK",
437    ] {
438        add(
439            &mut entries,
440            Principal::System,
441            Operation::Publish,
442            subject.to_owned(),
443        );
444    }
445    for subject in [
446        "$SYS.ACCOUNT.*.CONNECT".to_owned(),
447        "$SYS.ACCOUNT.*.DISCONNECT".to_owned(),
448        format!("_INBOX.{credential_public}.>"),
449    ] {
450        add(
451            &mut entries,
452            Principal::System,
453            Operation::Subscribe,
454            subject,
455        );
456    }
457    Ok(entries.into_iter().collect())
458}
459
460/// The grant for the flow engine (basal), issued by its attested module id.
461///
462/// It reads module events and produces nothing: its own inbox, census reads,
463/// the dead-letter record, pull and ack on its own durable `m_{module_id}` on
464/// the event stream, and ephemeral ordered consumers on that stream for dry-run
465/// replay. It publishes on no workload subject and on no event subject.
466///
467/// The ephemeral consumers rest on three rows. Consumer create is granted only
468/// in its unnamed form, the bare `$JS.API.CONSUMER.CREATE.<stream>` with no
469/// trailing token: a named or durable create carries the consumer name as a
470/// further token and matches nothing here, so the flow engine cannot create a
471/// durable, its own included (ck-bus creates that one). Consumer info on any
472/// consumer of the stream lets it inspect the ephemeral consumers, whose names
473/// the server picks. Publish on the stream's flow-control subjects is needed
474/// because an ordered consumer that cannot answer flow control stalls without
475/// an error. There is deliberately no consumer delete on the event stream, in
476/// any form: the flow engine creates its ephemeral consumers with a short
477/// inactivity threshold, so the server removes an abandoned one itself, and
478/// without delete the flow engine cannot remove another module's durable.
479pub fn flow_engine_permissions(
480    account: &AccountNames,
481    credential_public: &str,
482    module_id: &str,
483) -> Result<Vec<AllowEntry>, GrantError> {
484    validate_token(TokenKind::CredentialPublic, credential_public)?;
485    let durable = AccountNames::module_consumer_name(module_id)?;
486    let stream = &account.streams().event;
487    let mut entries = BTreeSet::new();
488
489    add(
490        &mut entries,
491        Principal::FlowEngine,
492        Operation::Subscribe,
493        format!("_INBOX.{credential_public}.>"),
494    );
495    add_census_read_permissions(&mut entries, Principal::FlowEngine, account);
496    for subject in [
497        account.effect_dead(),
498        // The flow engine's own durable, `m_{module_id}`.
499        format!("$JS.API.CONSUMER.MSG.NEXT.{stream}.{durable}"),
500        format!("$JS.API.CONSUMER.INFO.{stream}.{durable}"),
501        format!("$JS.ACK.{stream}.{durable}.>"),
502        format!("$JS.API.STREAM.INFO.{stream}"),
503        // The flow engine's ephemeral ordered consumers, for dry-run replay.
504        format!("$JS.API.CONSUMER.CREATE.{stream}"),
505        format!("$JS.API.CONSUMER.INFO.{stream}.*"),
506        format!("$JS.FC.{stream}.>"),
507    ] {
508        add(
509            &mut entries,
510            Principal::FlowEngine,
511            Operation::Publish,
512            subject,
513        );
514    }
515    Ok(entries.into_iter().collect())
516}
517
518pub fn generate_permission_golden(
519    fixture: GoldenFixture,
520) -> Result<PermissionDocument, GrantError> {
521    let account = AccountNames::derive(fixture.account)?;
522    validate_token(TokenKind::ModuleId, fixture.module_id)?;
523    validate_token(TokenKind::AgentId, fixture.foreign_agent)?;
524    validate_token(TokenKind::RoomId, fixture.unbound_room)?;
525    validate_token(TokenKind::ModuleId, fixture.foreign_module)?;
526    validate_token(TokenKind::ModuleId, fixture.delivery_authority_module)?;
527
528    let mut allows = BTreeSet::new();
529    allows.extend(participant_permissions(
530        &account,
531        fixture.module_id,
532        fixture.module_id,
533        &[fixture.bound_room],
534    )?);
535    allows.extend(delivery_authority_permissions(
536        &account,
537        fixture.delivery_authority_module,
538        fixture.delivery_authority_module,
539        &[fixture.bound_room],
540    )?);
541    allows.extend(bus_permissions(&account, fixture.module_id)?);
542    allows.extend(system_permissions(fixture.system_credential)?);
543    allows.extend(flow_engine_permissions(
544        &account,
545        fixture.flow_engine_module,
546        fixture.flow_engine_module,
547    )?);
548
549    let mut refused = BTreeSet::new();
550    // Credentials name no agents, so no participant may produce a workload
551    // message for any agent, its own or another's.
552    for agent_id in [fixture.bound_agent, fixture.foreign_agent] {
553        for subject in [
554            account.wake_fire(agent_id)?,
555            account.peer_delivery(agent_id, GOLDEN_SESSION)?,
556            account.effect_intent(agent_id, GOLDEN_SESSION)?,
557        ] {
558            refused.insert(RefusedEntry {
559                principal: Principal::Participant,
560                operation: Operation::Publish,
561                subject,
562                reason: "workload-publish",
563            });
564        }
565    }
566    for (operation, subject, reason) in [
567        (
568            Operation::Publish,
569            format!("$KV.{}.forbidden", account.buckets().census),
570            "census-write",
571        ),
572        (
573            Operation::Publish,
574            account.room_post(fixture.unbound_room)?,
575            "unbound-room",
576        ),
577        (
578            Operation::Publish,
579            "$SYS.REQ.CLAIMS.UPDATE".to_owned(),
580            "system-subject",
581        ),
582        (
583            Operation::Publish,
584            account.sentinel_ping(),
585            "sentinel-subject",
586        ),
587        (
588            Operation::Publish,
589            account.event_subject(fixture.foreign_module, GOLDEN_EVENT, GOLDEN_EVENT_VERSION)?,
590            "foreign-module-event",
591        ),
592    ] {
593        refused.insert(RefusedEntry {
594            principal: Principal::Participant,
595            operation,
596            subject,
597            reason,
598        });
599    }
600    for principal in [Principal::Participant, Principal::DeliveryAuthority] {
601        for stream in [
602            &account.streams().room,
603            &account.streams().wake,
604            &account.streams().peer,
605            &account.streams().effect,
606        ] {
607            refused.insert(RefusedEntry {
608                principal,
609                operation: Operation::Publish,
610                subject: format!("$JS.API.CONSUMER.CREATE.{stream}.c_{}", fixture.bound_agent),
611                reason: "consumer-create",
612            });
613        }
614    }
615    // The delivery authority publishes on the whole ROOM binding, so an
616    // unbound room's post is not refused to it the way it is to a participant.
617    for (subject, reason) in [
618        (
619            format!("$KV.{}.forbidden", account.buckets().census),
620            "census-write",
621        ),
622        ("$SYS.REQ.CLAIMS.UPDATE".to_owned(), "system-subject"),
623    ] {
624        refused.insert(RefusedEntry {
625            principal: Principal::DeliveryAuthority,
626            operation: Operation::Publish,
627            subject,
628            reason,
629        });
630    }
631    for principal in [Principal::Bus, Principal::System] {
632        for subject in [
633            account.room_post(fixture.bound_room)?,
634            account.wake_fire(fixture.bound_agent)?,
635            account.peer_delivery(fixture.bound_agent, GOLDEN_SESSION)?,
636            account.effect_intent(fixture.bound_agent, GOLDEN_SESSION)?,
637        ] {
638            refused.insert(RefusedEntry {
639                principal,
640                operation: Operation::Publish,
641                subject,
642                reason: "workload-publish",
643            });
644        }
645    }
646
647    // The ROOM module durable is the delivery authority's alone. A participant
648    // cannot read it; the delivery authority cannot read another module's
649    // durable on ROOM, nor create or delete any ROOM consumer (ck-bus creates
650    // its durable); the flow engine does not touch ROOM at all.
651    let room_stream = &account.streams().room;
652    let room_durable = AccountNames::module_consumer_name(fixture.delivery_authority_module)?;
653    let foreign_room_durable = AccountNames::module_consumer_name(fixture.foreign_module)?;
654    let participant_room_durable = AccountNames::module_consumer_name(fixture.module_id)?;
655    let room_binding = account.room_binding();
656    for (principal, subject, reason) in [
657        (
658            Principal::Participant,
659            format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{room_durable}"),
660            "room-durable",
661        ),
662        (
663            Principal::Participant,
664            format!("$JS.ACK.{room_stream}.{room_durable}.>"),
665            "room-durable",
666        ),
667        (
668            Principal::DeliveryAuthority,
669            format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{foreign_room_durable}"),
670            "foreign-room-durable",
671        ),
672        (
673            Principal::DeliveryAuthority,
674            format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{participant_room_durable}"),
675            "foreign-room-durable",
676        ),
677        (
678            Principal::DeliveryAuthority,
679            format!("$JS.ACK.{room_stream}.{foreign_room_durable}.>"),
680            "foreign-room-durable",
681        ),
682        (
683            Principal::DeliveryAuthority,
684            format!("$JS.API.CONSUMER.CREATE.{room_stream}"),
685            "consumer-create",
686        ),
687        (
688            Principal::DeliveryAuthority,
689            format!("$JS.API.CONSUMER.CREATE.{room_stream}.{room_durable}.{room_binding}"),
690            "named-consumer-create",
691        ),
692        (
693            Principal::DeliveryAuthority,
694            format!("$JS.API.CONSUMER.DURABLE.CREATE.{room_stream}.{room_durable}"),
695            "durable-create",
696        ),
697        (
698            Principal::DeliveryAuthority,
699            format!("$JS.API.CONSUMER.DELETE.{room_stream}.{room_durable}"),
700            "consumer-delete",
701        ),
702        (
703            Principal::DeliveryAuthority,
704            format!("$JS.API.CONSUMER.DELETE.{room_stream}.{foreign_room_durable}"),
705            "consumer-delete",
706        ),
707        (
708            Principal::FlowEngine,
709            format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{room_durable}"),
710            "room-stream",
711        ),
712        (
713            Principal::FlowEngine,
714            format!("$JS.ACK.{room_stream}.{room_durable}.>"),
715            "room-stream",
716        ),
717        (
718            Principal::FlowEngine,
719            format!("$JS.API.CONSUMER.INFO.{room_stream}.{room_durable}"),
720            "room-stream",
721        ),
722        (
723            Principal::FlowEngine,
724            format!("$JS.API.CONSUMER.CREATE.{room_stream}"),
725            "room-stream",
726        ),
727        (
728            Principal::FlowEngine,
729            format!("$JS.API.STREAM.INFO.{room_stream}"),
730            "room-stream",
731        ),
732        (
733            Principal::FlowEngine,
734            account.room_post(fixture.bound_room)?,
735            "room-stream",
736        ),
737    ] {
738        refused.insert(RefusedEntry {
739            principal,
740            operation: Operation::Publish,
741            subject,
742            reason,
743        });
744    }
745    refused.insert(RefusedEntry {
746        principal: Principal::FlowEngine,
747        operation: Operation::Subscribe,
748        subject: account.room_subscription(fixture.bound_room)?,
749        reason: "room-stream",
750    });
751
752    // The flow engine cannot create a named or durable consumer on the event
753    // stream, cannot delete any consumer there, cannot read an agent's
754    // durable, and publishes no workload message.
755    let event_stream = &account.streams().event;
756    let own_durable = AccountNames::module_consumer_name(fixture.flow_engine_module)?;
757    let foreign_durable = AccountNames::module_consumer_name(fixture.foreign_module)?;
758    for (subject, reason) in [
759        (
760            format!("$JS.API.CONSUMER.DURABLE.CREATE.{event_stream}.{foreign_durable}"),
761            "durable-create",
762        ),
763        (
764            format!(
765                "$JS.API.CONSUMER.CREATE.{event_stream}.{own_durable}.{}",
766                account.event_binding()
767            ),
768            "named-consumer-create",
769        ),
770        (
771            format!("$JS.API.CONSUMER.DELETE.{event_stream}.{foreign_durable}"),
772            "consumer-delete",
773        ),
774        (
775            format!("$JS.API.CONSUMER.DELETE.{event_stream}.{own_durable}"),
776            "consumer-delete",
777        ),
778        (
779            format!(
780                "$JS.API.CONSUMER.MSG.NEXT.{}.{}",
781                account.streams().wake,
782                AccountNames::consumer_name(fixture.bound_agent)?
783            ),
784            "agent-durable",
785        ),
786        (account.wake_fire(fixture.bound_agent)?, "workload-publish"),
787    ] {
788        refused.insert(RefusedEntry {
789            principal: Principal::FlowEngine,
790            operation: Operation::Publish,
791            subject,
792            reason,
793        });
794    }
795
796    let document = PermissionDocument {
797        fixture,
798        allows: allows.into_iter().collect(),
799        refused: refused.into_iter().collect(),
800    };
801    validate_permission_file(&document.render(), &account)?;
802    Ok(document)
803}
804
805/// Validates the checked-in line format and its allow-only, whole-token contract.
806pub fn validate_permission_file(contents: &str, account: &AccountNames) -> Result<(), GrantError> {
807    let expected_streams = account.streams().all();
808    let mut expanded_streams = BTreeSet::new();
809    let mut seen_entries = BTreeSet::new();
810
811    for (zero_index, raw_line) in contents.lines().enumerate() {
812        let line_number = zero_index + 1;
813        let line = raw_line.trim();
814        if line.is_empty() || line.starts_with('#') {
815            continue;
816        }
817        let fields = line.split_whitespace().collect::<Vec<_>>();
818        if fields
819            .first()
820            .is_some_and(|field| field.eq_ignore_ascii_case("deny"))
821        {
822            return Err(file_error(
823                line_number,
824                "deny entries are forbidden; absence from allow is the denial",
825            ));
826        }
827
828        let subject = match fields.as_slice() {
829            ["fixture", _, _] => continue,
830            ["allow", principal, operation, subject] => {
831                validate_principal(line_number, principal)?;
832                validate_operation(line_number, operation)?;
833                if !seen_entries.insert(line.to_owned()) {
834                    return Err(file_error(line_number, "duplicate permission entry"));
835                }
836                for stream in expected_streams {
837                    if subject.split('.').any(|token| token == stream) {
838                        expanded_streams.insert(stream.to_owned());
839                    }
840                }
841                *subject
842            }
843            ["expect-refused", principal, operation, subject, _] => {
844                validate_principal(line_number, principal)?;
845                validate_operation(line_number, operation)?;
846                if !seen_entries.insert(line.to_owned()) {
847                    return Err(file_error(line_number, "duplicate expectation entry"));
848                }
849                *subject
850            }
851            _ => {
852                return Err(file_error(
853                    line_number,
854                    "expected fixture, allow, or expect-refused entry",
855                ));
856            }
857        };
858        validate_subject(line_number, subject)?;
859    }
860
861    for stream in expected_streams {
862        if !expanded_streams.contains(stream) {
863            return Err(GrantError::MissingStreamExpansion {
864                stream: stream.to_owned(),
865            });
866        }
867    }
868    Ok(())
869}
870
871fn add(
872    entries: &mut BTreeSet<AllowEntry>,
873    principal: Principal,
874    operation: Operation,
875    subject: String,
876) {
877    entries.insert(AllowEntry {
878        principal,
879        operation,
880        subject,
881    });
882}
883
884fn add_census_read_permissions(
885    entries: &mut BTreeSet<AllowEntry>,
886    principal: Principal,
887    account: &AccountNames,
888) {
889    add(
890        entries,
891        principal,
892        Operation::Subscribe,
893        format!("$KV.{}.>", account.buckets().census),
894    );
895    let stream = &account.buckets().census_stream;
896    for subject in [
897        format!("$JS.API.DIRECT.GET.{stream}.>"),
898        format!("$JS.API.STREAM.MSG.GET.{stream}"),
899        format!("$JS.API.STREAM.INFO.{stream}"),
900        format!("$JS.API.CONSUMER.CREATE.{stream}"),
901        format!("$JS.API.CONSUMER.CREATE.{stream}.>"),
902        format!("$JS.API.CONSUMER.DURABLE.CREATE.{stream}.>"),
903        format!("$JS.API.CONSUMER.DELETE.{stream}.>"),
904    ] {
905        add(entries, principal, Operation::Publish, subject);
906    }
907}
908
909fn validate_principal(line: usize, principal: &str) -> Result<(), GrantError> {
910    if matches!(
911        principal,
912        "participant" | "delivery-authority" | "bus" | "system" | "flow-engine"
913    ) {
914        Ok(())
915    } else {
916        Err(file_error(line, "unknown principal"))
917    }
918}
919
920fn validate_operation(line: usize, operation: &str) -> Result<(), GrantError> {
921    if matches!(operation, "publish" | "subscribe") {
922        Ok(())
923    } else {
924        Err(file_error(line, "unknown permission operation"))
925    }
926}
927
928fn validate_subject(line: usize, subject: &str) -> Result<(), GrantError> {
929    if subject.contains('{') || subject.contains('}') {
930        return Err(file_error(line, "template placeholders must be expanded"));
931    }
932    let tokens = subject.split('.').collect::<Vec<_>>();
933    if tokens.iter().any(|token| token.is_empty()) {
934        return Err(file_error(line, "subject contains an empty token"));
935    }
936    for (index, token) in tokens.iter().enumerate() {
937        if token.bytes().any(|byte| byte.is_ascii_whitespace()) {
938            return Err(file_error(line, "subject tokens cannot contain whitespace"));
939        }
940        if (token.contains('*') && *token != "*") || (token.contains('>') && *token != ">") {
941            return Err(file_error(
942                line,
943                "wildcards must occupy a whole subject token",
944            ));
945        }
946        if *token == ">" && index + 1 != tokens.len() {
947            return Err(file_error(line, "the > wildcard must be the final token"));
948        }
949    }
950    Ok(())
951}
952
953fn file_error(line: usize, reason: impl Into<String>) -> GrantError {
954    GrantError::InvalidPermissionFile {
955        line,
956        reason: reason.into(),
957    }
958}