use std::collections::BTreeSet;
use std::error::Error;
use std::fmt;
use crate::names::AccountNames;
use crate::token::{validate_token, NamingError, TokenKind};
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum Principal {
Participant,
DeliveryAuthority,
Bus,
System,
FlowEngine,
}
impl Principal {
pub const fn as_str(self) -> &'static str {
match self {
Self::Participant => "participant",
Self::DeliveryAuthority => "delivery-authority",
Self::Bus => "bus",
Self::System => "system",
Self::FlowEngine => "flow-engine",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum Operation {
Publish,
Subscribe,
}
impl Operation {
pub const fn as_str(self) -> &'static str {
match self {
Self::Publish => "publish",
Self::Subscribe => "subscribe",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
pub struct AllowEntry {
pub principal: Principal,
pub operation: Operation,
pub subject: String,
}
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord)]
pub struct RefusedEntry {
pub principal: Principal,
pub operation: Operation,
pub subject: String,
pub reason: &'static str,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct GoldenFixture {
pub account: &'static str,
pub module_id: &'static str,
pub bound_agent: &'static str,
pub foreign_agent: &'static str,
pub bound_room: &'static str,
pub unbound_room: &'static str,
pub system_credential: &'static str,
pub flow_engine_module: &'static str,
pub foreign_module: &'static str,
pub delivery_authority_module: &'static str,
}
pub const PINNED_GOLDEN_FIXTURE: GoldenFixture = GoldenFixture {
account: "box_goldenfixture",
module_id: "ckbus",
bound_agent: "agent_gold_a",
foreign_agent: "agent_gold_b",
bound_room: "room_gold_bound",
unbound_room: "room_gold_unbound",
system_credential: "cksys",
flow_engine_module: "basal",
foreign_module: "other",
delivery_authority_module: "prefrontal-core",
};
const GOLDEN_SESSION: &str = "session_gold";
const GOLDEN_EVENT: &str = "pull_request_review";
const GOLDEN_EVENT_VERSION: u32 = 1;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PermissionDocument {
fixture: GoldenFixture,
allows: Vec<AllowEntry>,
refused: Vec<RefusedEntry>,
}
impl PermissionDocument {
pub fn allows(&self) -> &[AllowEntry] {
&self.allows
}
pub fn refused(&self) -> &[RefusedEntry] {
&self.refused
}
pub fn render(&self) -> String {
let mut output = String::from(
"# Generated by cortexkit-bus-naming; edit the generator, not this file.\n",
);
for (name, value) in [
("account", self.fixture.account),
("module", self.fixture.module_id),
("bound-agent", self.fixture.bound_agent),
("foreign-agent", self.fixture.foreign_agent),
("bound-room", self.fixture.bound_room),
("unbound-room", self.fixture.unbound_room),
("flow-engine-module", self.fixture.flow_engine_module),
("foreign-module", self.fixture.foreign_module),
(
"delivery-authority-module",
self.fixture.delivery_authority_module,
),
] {
output.push_str(&format!("fixture {name} {value}\n"));
}
output.push('\n');
for entry in &self.allows {
output.push_str(&format!(
"allow {} {} {}\n",
entry.principal.as_str(),
entry.operation.as_str(),
entry.subject
));
}
output.push('\n');
for entry in &self.refused {
output.push_str(&format!(
"expect-refused {} {} {} {}\n",
entry.principal.as_str(),
entry.operation.as_str(),
entry.subject,
entry.reason
));
}
output
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum GrantError {
Naming(NamingError),
InvalidPermissionFile { line: usize, reason: String },
MissingStreamExpansion { stream: String },
}
impl fmt::Display for GrantError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Naming(error) => error.fmt(f),
Self::InvalidPermissionFile { line, reason } => {
write!(f, "invalid permission file at line {line}: {reason}")
}
Self::MissingStreamExpansion { stream } => {
write!(f, "permission file does not expand literal stream {stream}")
}
}
}
}
impl Error for GrantError {
fn source(&self) -> Option<&(dyn Error + 'static)> {
match self {
Self::Naming(error) => Some(error),
_ => None,
}
}
}
impl From<NamingError> for GrantError {
fn from(value: NamingError) -> Self {
Self::Naming(value)
}
}
pub fn participant_permissions(
account: &AccountNames,
credential_public: &str,
module_id: &str,
bound_rooms: &[&str],
) -> Result<Vec<AllowEntry>, GrantError> {
let mut entries = BTreeSet::new();
add_consumer_permissions(
&mut entries,
Principal::Participant,
account,
credential_public,
module_id,
bound_rooms,
)?;
Ok(entries.into_iter().collect())
}
pub fn delivery_authority_permissions(
account: &AccountNames,
credential_public: &str,
module_id: &str,
bound_rooms: &[&str],
) -> Result<Vec<AllowEntry>, GrantError> {
let mut entries = BTreeSet::new();
add_consumer_permissions(
&mut entries,
Principal::DeliveryAuthority,
account,
credential_public,
module_id,
bound_rooms,
)?;
let room_stream = &account.streams().room;
let room_durable = AccountNames::module_consumer_name(module_id)?;
for subject in [
account.wake_binding(),
account.peer_binding(),
account.effect_binding(),
account.room_binding(),
format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{room_durable}"),
format!("$JS.API.CONSUMER.INFO.{room_stream}.{room_durable}"),
format!("$JS.ACK.{room_stream}.{room_durable}.>"),
format!("$JS.API.STREAM.INFO.{room_stream}"),
] {
add(
&mut entries,
Principal::DeliveryAuthority,
Operation::Publish,
subject,
);
}
Ok(entries.into_iter().collect())
}
fn add_consumer_permissions(
entries: &mut BTreeSet<AllowEntry>,
principal: Principal,
account: &AccountNames,
credential_public: &str,
module_id: &str,
bound_rooms: &[&str],
) -> Result<(), GrantError> {
validate_token(TokenKind::CredentialPublic, credential_public)?;
add(
entries,
principal,
Operation::Subscribe,
format!("_INBOX.{credential_public}.>"),
);
add(
entries,
principal,
Operation::Publish,
account.effect_dead(),
);
add(
entries,
principal,
Operation::Publish,
account.event_publish_grant(module_id)?,
);
for stream in account.streams().agent_streams() {
for subject in [
format!("$JS.ACK.{stream}.*.>"),
format!("$JS.API.CONSUMER.MSG.NEXT.{stream}.*"),
format!("$JS.API.CONSUMER.INFO.{stream}.*"),
format!("$JS.API.STREAM.INFO.{stream}"),
] {
add(entries, principal, Operation::Publish, subject);
}
}
for room_id in bound_rooms {
validate_token(TokenKind::RoomId, room_id)?;
add(
entries,
principal,
Operation::Publish,
account.room_post(room_id)?,
);
add(
entries,
principal,
Operation::Subscribe,
account.room_subscription(room_id)?,
);
}
add_census_read_permissions(entries, principal, account);
Ok(())
}
pub fn bus_permissions(
account: &AccountNames,
credential_public: &str,
) -> Result<Vec<AllowEntry>, GrantError> {
validate_token(TokenKind::CredentialPublic, credential_public)?;
let mut entries = BTreeSet::new();
let inbox = format!("_INBOX.{credential_public}.>");
for operation in [Operation::Publish, Operation::Subscribe] {
add(&mut entries, Principal::Bus, operation, inbox.clone());
add(
&mut entries,
Principal::Bus,
operation,
account.sentinel_ping(),
);
add(
&mut entries,
Principal::Bus,
operation,
format!("$KV.{}.>", account.buckets().census),
);
}
add_census_read_permissions(&mut entries, Principal::Bus, account);
let census_stream = &account.buckets().census_stream;
for subject in [
format!("$JS.API.STREAM.CREATE.{census_stream}"),
format!("$JS.API.STREAM.UPDATE.{census_stream}"),
] {
add(&mut entries, Principal::Bus, Operation::Publish, subject);
}
for stream in account.streams().all() {
for subject in [
format!("$JS.API.STREAM.CREATE.{stream}"),
format!("$JS.API.STREAM.UPDATE.{stream}"),
format!("$JS.API.STREAM.INFO.{stream}"),
format!("$JS.API.STREAM.DELETE.{stream}"),
format!("$JS.API.CONSUMER.CREATE.{stream}.>"),
format!("$JS.API.CONSUMER.DURABLE.CREATE.{stream}.>"),
format!("$JS.API.CONSUMER.INFO.{stream}.>"),
format!("$JS.API.CONSUMER.DELETE.{stream}.>"),
] {
add(&mut entries, Principal::Bus, Operation::Publish, subject);
}
}
let dead_stream = &account.streams().effect_dead;
for subject in [
format!("$JS.ACK.{dead_stream}.c_ckbus_dead.>"),
format!("$JS.API.CONSUMER.MSG.NEXT.{dead_stream}.c_ckbus_dead"),
format!("$JS.API.CONSUMER.INFO.{dead_stream}.c_ckbus_dead"),
] {
add(&mut entries, Principal::Bus, Operation::Publish, subject);
}
for stream in account.streams().agent_streams() {
for subject in [
format!("$JS.API.CONSUMER.NAMES.{stream}"),
format!("$JS.API.STREAM.PURGE.{stream}"),
] {
add(&mut entries, Principal::Bus, Operation::Publish, subject);
}
}
Ok(entries.into_iter().collect())
}
pub fn system_permissions(credential_public: &str) -> Result<Vec<AllowEntry>, GrantError> {
validate_token(TokenKind::CredentialPublic, credential_public)?;
let mut entries = BTreeSet::new();
for subject in [
"$SYS.REQ.CLAIMS.UPDATE",
"$SYS.REQ.CLAIMS.LIST",
"$SYS.REQ.ACCOUNT.*.CLAIMS.LOOKUP",
"$SYS.REQ.SERVER.*.KICK",
] {
add(
&mut entries,
Principal::System,
Operation::Publish,
subject.to_owned(),
);
}
for subject in [
"$SYS.ACCOUNT.*.CONNECT".to_owned(),
"$SYS.ACCOUNT.*.DISCONNECT".to_owned(),
format!("_INBOX.{credential_public}.>"),
] {
add(
&mut entries,
Principal::System,
Operation::Subscribe,
subject,
);
}
Ok(entries.into_iter().collect())
}
pub fn flow_engine_permissions(
account: &AccountNames,
credential_public: &str,
module_id: &str,
) -> Result<Vec<AllowEntry>, GrantError> {
validate_token(TokenKind::CredentialPublic, credential_public)?;
let durable = AccountNames::module_consumer_name(module_id)?;
let stream = &account.streams().event;
let mut entries = BTreeSet::new();
add(
&mut entries,
Principal::FlowEngine,
Operation::Subscribe,
format!("_INBOX.{credential_public}.>"),
);
add_census_read_permissions(&mut entries, Principal::FlowEngine, account);
for subject in [
account.effect_dead(),
format!("$JS.API.CONSUMER.MSG.NEXT.{stream}.{durable}"),
format!("$JS.API.CONSUMER.INFO.{stream}.{durable}"),
format!("$JS.ACK.{stream}.{durable}.>"),
format!("$JS.API.STREAM.INFO.{stream}"),
format!("$JS.API.CONSUMER.CREATE.{stream}"),
format!("$JS.API.CONSUMER.INFO.{stream}.*"),
format!("$JS.FC.{stream}.>"),
] {
add(
&mut entries,
Principal::FlowEngine,
Operation::Publish,
subject,
);
}
Ok(entries.into_iter().collect())
}
pub fn generate_permission_golden(
fixture: GoldenFixture,
) -> Result<PermissionDocument, GrantError> {
let account = AccountNames::derive(fixture.account)?;
validate_token(TokenKind::ModuleId, fixture.module_id)?;
validate_token(TokenKind::AgentId, fixture.foreign_agent)?;
validate_token(TokenKind::RoomId, fixture.unbound_room)?;
validate_token(TokenKind::ModuleId, fixture.foreign_module)?;
validate_token(TokenKind::ModuleId, fixture.delivery_authority_module)?;
let mut allows = BTreeSet::new();
allows.extend(participant_permissions(
&account,
fixture.module_id,
fixture.module_id,
&[fixture.bound_room],
)?);
allows.extend(delivery_authority_permissions(
&account,
fixture.delivery_authority_module,
fixture.delivery_authority_module,
&[fixture.bound_room],
)?);
allows.extend(bus_permissions(&account, fixture.module_id)?);
allows.extend(system_permissions(fixture.system_credential)?);
allows.extend(flow_engine_permissions(
&account,
fixture.flow_engine_module,
fixture.flow_engine_module,
)?);
let mut refused = BTreeSet::new();
for agent_id in [fixture.bound_agent, fixture.foreign_agent] {
for subject in [
account.wake_fire(agent_id)?,
account.peer_delivery(agent_id, GOLDEN_SESSION)?,
account.effect_intent(agent_id, GOLDEN_SESSION)?,
] {
refused.insert(RefusedEntry {
principal: Principal::Participant,
operation: Operation::Publish,
subject,
reason: "workload-publish",
});
}
}
for (operation, subject, reason) in [
(
Operation::Publish,
format!("$KV.{}.forbidden", account.buckets().census),
"census-write",
),
(
Operation::Publish,
account.room_post(fixture.unbound_room)?,
"unbound-room",
),
(
Operation::Publish,
"$SYS.REQ.CLAIMS.UPDATE".to_owned(),
"system-subject",
),
(
Operation::Publish,
account.sentinel_ping(),
"sentinel-subject",
),
(
Operation::Publish,
account.event_subject(fixture.foreign_module, GOLDEN_EVENT, GOLDEN_EVENT_VERSION)?,
"foreign-module-event",
),
] {
refused.insert(RefusedEntry {
principal: Principal::Participant,
operation,
subject,
reason,
});
}
for principal in [Principal::Participant, Principal::DeliveryAuthority] {
for stream in [
&account.streams().room,
&account.streams().wake,
&account.streams().peer,
&account.streams().effect,
] {
refused.insert(RefusedEntry {
principal,
operation: Operation::Publish,
subject: format!("$JS.API.CONSUMER.CREATE.{stream}.c_{}", fixture.bound_agent),
reason: "consumer-create",
});
}
}
for (subject, reason) in [
(
format!("$KV.{}.forbidden", account.buckets().census),
"census-write",
),
("$SYS.REQ.CLAIMS.UPDATE".to_owned(), "system-subject"),
] {
refused.insert(RefusedEntry {
principal: Principal::DeliveryAuthority,
operation: Operation::Publish,
subject,
reason,
});
}
for principal in [Principal::Bus, Principal::System] {
for subject in [
account.room_post(fixture.bound_room)?,
account.wake_fire(fixture.bound_agent)?,
account.peer_delivery(fixture.bound_agent, GOLDEN_SESSION)?,
account.effect_intent(fixture.bound_agent, GOLDEN_SESSION)?,
] {
refused.insert(RefusedEntry {
principal,
operation: Operation::Publish,
subject,
reason: "workload-publish",
});
}
}
let room_stream = &account.streams().room;
let room_durable = AccountNames::module_consumer_name(fixture.delivery_authority_module)?;
let foreign_room_durable = AccountNames::module_consumer_name(fixture.foreign_module)?;
let participant_room_durable = AccountNames::module_consumer_name(fixture.module_id)?;
let room_binding = account.room_binding();
for (principal, subject, reason) in [
(
Principal::Participant,
format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{room_durable}"),
"room-durable",
),
(
Principal::Participant,
format!("$JS.ACK.{room_stream}.{room_durable}.>"),
"room-durable",
),
(
Principal::DeliveryAuthority,
format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{foreign_room_durable}"),
"foreign-room-durable",
),
(
Principal::DeliveryAuthority,
format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{participant_room_durable}"),
"foreign-room-durable",
),
(
Principal::DeliveryAuthority,
format!("$JS.ACK.{room_stream}.{foreign_room_durable}.>"),
"foreign-room-durable",
),
(
Principal::DeliveryAuthority,
format!("$JS.API.CONSUMER.CREATE.{room_stream}"),
"consumer-create",
),
(
Principal::DeliveryAuthority,
format!("$JS.API.CONSUMER.CREATE.{room_stream}.{room_durable}.{room_binding}"),
"named-consumer-create",
),
(
Principal::DeliveryAuthority,
format!("$JS.API.CONSUMER.DURABLE.CREATE.{room_stream}.{room_durable}"),
"durable-create",
),
(
Principal::DeliveryAuthority,
format!("$JS.API.CONSUMER.DELETE.{room_stream}.{room_durable}"),
"consumer-delete",
),
(
Principal::DeliveryAuthority,
format!("$JS.API.CONSUMER.DELETE.{room_stream}.{foreign_room_durable}"),
"consumer-delete",
),
(
Principal::FlowEngine,
format!("$JS.API.CONSUMER.MSG.NEXT.{room_stream}.{room_durable}"),
"room-stream",
),
(
Principal::FlowEngine,
format!("$JS.ACK.{room_stream}.{room_durable}.>"),
"room-stream",
),
(
Principal::FlowEngine,
format!("$JS.API.CONSUMER.INFO.{room_stream}.{room_durable}"),
"room-stream",
),
(
Principal::FlowEngine,
format!("$JS.API.CONSUMER.CREATE.{room_stream}"),
"room-stream",
),
(
Principal::FlowEngine,
format!("$JS.API.STREAM.INFO.{room_stream}"),
"room-stream",
),
(
Principal::FlowEngine,
account.room_post(fixture.bound_room)?,
"room-stream",
),
] {
refused.insert(RefusedEntry {
principal,
operation: Operation::Publish,
subject,
reason,
});
}
refused.insert(RefusedEntry {
principal: Principal::FlowEngine,
operation: Operation::Subscribe,
subject: account.room_subscription(fixture.bound_room)?,
reason: "room-stream",
});
let event_stream = &account.streams().event;
let own_durable = AccountNames::module_consumer_name(fixture.flow_engine_module)?;
let foreign_durable = AccountNames::module_consumer_name(fixture.foreign_module)?;
for (subject, reason) in [
(
format!("$JS.API.CONSUMER.DURABLE.CREATE.{event_stream}.{foreign_durable}"),
"durable-create",
),
(
format!(
"$JS.API.CONSUMER.CREATE.{event_stream}.{own_durable}.{}",
account.event_binding()
),
"named-consumer-create",
),
(
format!("$JS.API.CONSUMER.DELETE.{event_stream}.{foreign_durable}"),
"consumer-delete",
),
(
format!("$JS.API.CONSUMER.DELETE.{event_stream}.{own_durable}"),
"consumer-delete",
),
(
format!(
"$JS.API.CONSUMER.MSG.NEXT.{}.{}",
account.streams().wake,
AccountNames::consumer_name(fixture.bound_agent)?
),
"agent-durable",
),
(account.wake_fire(fixture.bound_agent)?, "workload-publish"),
] {
refused.insert(RefusedEntry {
principal: Principal::FlowEngine,
operation: Operation::Publish,
subject,
reason,
});
}
let document = PermissionDocument {
fixture,
allows: allows.into_iter().collect(),
refused: refused.into_iter().collect(),
};
validate_permission_file(&document.render(), &account)?;
Ok(document)
}
pub fn validate_permission_file(contents: &str, account: &AccountNames) -> Result<(), GrantError> {
let expected_streams = account.streams().all();
let mut expanded_streams = BTreeSet::new();
let mut seen_entries = BTreeSet::new();
for (zero_index, raw_line) in contents.lines().enumerate() {
let line_number = zero_index + 1;
let line = raw_line.trim();
if line.is_empty() || line.starts_with('#') {
continue;
}
let fields = line.split_whitespace().collect::<Vec<_>>();
if fields
.first()
.is_some_and(|field| field.eq_ignore_ascii_case("deny"))
{
return Err(file_error(
line_number,
"deny entries are forbidden; absence from allow is the denial",
));
}
let subject = match fields.as_slice() {
["fixture", _, _] => continue,
["allow", principal, operation, subject] => {
validate_principal(line_number, principal)?;
validate_operation(line_number, operation)?;
if !seen_entries.insert(line.to_owned()) {
return Err(file_error(line_number, "duplicate permission entry"));
}
for stream in expected_streams {
if subject.split('.').any(|token| token == stream) {
expanded_streams.insert(stream.to_owned());
}
}
*subject
}
["expect-refused", principal, operation, subject, _] => {
validate_principal(line_number, principal)?;
validate_operation(line_number, operation)?;
if !seen_entries.insert(line.to_owned()) {
return Err(file_error(line_number, "duplicate expectation entry"));
}
*subject
}
_ => {
return Err(file_error(
line_number,
"expected fixture, allow, or expect-refused entry",
));
}
};
validate_subject(line_number, subject)?;
}
for stream in expected_streams {
if !expanded_streams.contains(stream) {
return Err(GrantError::MissingStreamExpansion {
stream: stream.to_owned(),
});
}
}
Ok(())
}
fn add(
entries: &mut BTreeSet<AllowEntry>,
principal: Principal,
operation: Operation,
subject: String,
) {
entries.insert(AllowEntry {
principal,
operation,
subject,
});
}
fn add_census_read_permissions(
entries: &mut BTreeSet<AllowEntry>,
principal: Principal,
account: &AccountNames,
) {
add(
entries,
principal,
Operation::Subscribe,
format!("$KV.{}.>", account.buckets().census),
);
let stream = &account.buckets().census_stream;
for subject in [
format!("$JS.API.DIRECT.GET.{stream}.>"),
format!("$JS.API.STREAM.MSG.GET.{stream}"),
format!("$JS.API.STREAM.INFO.{stream}"),
format!("$JS.API.CONSUMER.CREATE.{stream}"),
format!("$JS.API.CONSUMER.CREATE.{stream}.>"),
format!("$JS.API.CONSUMER.DURABLE.CREATE.{stream}.>"),
format!("$JS.API.CONSUMER.DELETE.{stream}.>"),
] {
add(entries, principal, Operation::Publish, subject);
}
}
fn validate_principal(line: usize, principal: &str) -> Result<(), GrantError> {
if matches!(
principal,
"participant" | "delivery-authority" | "bus" | "system" | "flow-engine"
) {
Ok(())
} else {
Err(file_error(line, "unknown principal"))
}
}
fn validate_operation(line: usize, operation: &str) -> Result<(), GrantError> {
if matches!(operation, "publish" | "subscribe") {
Ok(())
} else {
Err(file_error(line, "unknown permission operation"))
}
}
fn validate_subject(line: usize, subject: &str) -> Result<(), GrantError> {
if subject.contains('{') || subject.contains('}') {
return Err(file_error(line, "template placeholders must be expanded"));
}
let tokens = subject.split('.').collect::<Vec<_>>();
if tokens.iter().any(|token| token.is_empty()) {
return Err(file_error(line, "subject contains an empty token"));
}
for (index, token) in tokens.iter().enumerate() {
if token.bytes().any(|byte| byte.is_ascii_whitespace()) {
return Err(file_error(line, "subject tokens cannot contain whitespace"));
}
if (token.contains('*') && *token != "*") || (token.contains('>') && *token != ">") {
return Err(file_error(
line,
"wildcards must occupy a whole subject token",
));
}
if *token == ">" && index + 1 != tokens.len() {
return Err(file_error(line, "the > wildcard must be the final token"));
}
}
Ok(())
}
fn file_error(line: usize, reason: impl Into<String>) -> GrantError {
GrantError::InvalidPermissionFile {
line,
reason: reason.into(),
}
}