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 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#[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 pub system_credential: &'static str,
71 pub flow_engine_module: &'static str,
73 pub foreign_module: &'static str,
76 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
94const GOLDEN_SESSION: &str = "session_gold";
96
97const 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
196pub 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
224pub 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
279fn 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 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 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 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
422pub 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
460pub 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 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 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 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 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 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 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
805pub 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}