use std::collections::BTreeSet;
use thiserror::Error;
use crate::execution::ExecutionLabel;
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub enum BrokerFamily {
Connector,
Peer,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ReturnedContent {
Untrusted,
Quarantined,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BrokerCorrelation {
pub tenant: String,
pub conversation: String,
pub execution: String,
pub attempt: String,
pub step: String,
pub tool_call: String,
pub audience: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BrokerCapability {
pub correlation: BrokerCorrelation,
pub admitted_label: ExecutionLabel,
pub family: BrokerFamily,
pub target: String,
pub method: String,
pub resource: String,
pub max_request_bytes: usize,
pub expires_at_millis: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BrokerRequest {
pub correlation: BrokerCorrelation,
pub label: ExecutionLabel,
pub family: BrokerFamily,
pub target: String,
pub method: String,
pub resource: String,
pub encoded_request_bytes: usize,
pub deadline_at_millis: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TrustedTarget {
pub target: String,
pub family: BrokerFamily,
pub method: String,
pub resource: String,
pub read_only: bool,
pub classification: ReturnedContent,
}
pub trait TrustedRegistry {
fn resolve(&self, family: BrokerFamily, target: &str) -> Option<TrustedTarget>;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct BrokerLimits {
pub max_calls: usize,
pub max_pending: usize,
pub session_deadline_at_millis: u64,
pub broker_deadline_at_millis: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BrokerAdmission {
pub target: TrustedTarget,
pub deadline_at_millis: u64,
pub correlation: BrokerCorrelation,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)]
pub enum BrokerRefusal {
#[error("broker capability binding mismatch")]
BindingMismatch,
#[error("Execution label is not admitted by this broker session")]
UnadmittedLabel,
#[error("broker capability expired")]
Expired,
#[error("trusted registry has no matching target")]
UnknownTarget,
#[error("broker request exceeds its exact target authority")]
ScopeMismatch,
#[error("broker request exceeds encoded body limit")]
RequestTooLarge,
#[error("broker session capacity exhausted")]
CapacityExceeded,
#[error("broker deadline exhausted")]
DeadlineExceeded,
#[error("broker capability was already used")]
CapabilityReplayed,
}
#[derive(Debug)]
pub struct BrokerSession {
limits: BrokerLimits,
calls: usize,
pending: usize,
used: BTreeSet<String>,
}
impl BrokerSession {
#[must_use]
pub const fn new(limits: BrokerLimits) -> Self {
Self {
limits,
calls: 0,
pending: 0,
used: BTreeSet::new(),
}
}
pub fn admit(
&mut self,
capability: &BrokerCapability,
request: &BrokerRequest,
registry: &impl TrustedRegistry,
now_millis: u64,
) -> Result<BrokerAdmission, BrokerRefusal> {
if capability.correlation != request.correlation
|| capability.family != request.family
|| capability.target != request.target
|| capability.method != request.method
|| capability.resource != request.resource
{
return Err(BrokerRefusal::BindingMismatch);
}
if capability.admitted_label != request.label {
return Err(BrokerRefusal::UnadmittedLabel);
}
if now_millis >= capability.expires_at_millis {
return Err(BrokerRefusal::Expired);
}
if request.encoded_request_bytes > capability.max_request_bytes {
return Err(BrokerRefusal::RequestTooLarge);
}
let deadline_at_millis = capability
.expires_at_millis
.min(self.limits.session_deadline_at_millis)
.min(self.limits.broker_deadline_at_millis)
.min(request.deadline_at_millis);
if now_millis >= deadline_at_millis {
return Err(BrokerRefusal::DeadlineExceeded);
}
let use_key = capability_use_key(capability);
if self.used.contains(&use_key) {
return Err(BrokerRefusal::CapabilityReplayed);
}
if self.calls >= self.limits.max_calls || self.pending >= self.limits.max_pending {
return Err(BrokerRefusal::CapacityExceeded);
}
let target = registry
.resolve(request.family, &request.target)
.filter(|target| target.family == request.family)
.ok_or(BrokerRefusal::UnknownTarget)?;
if target.method != request.method
|| target.resource != request.resource
|| !target.read_only
{
return Err(BrokerRefusal::ScopeMismatch);
}
self.used.insert(use_key);
self.calls += 1;
self.pending += 1;
Ok(BrokerAdmission {
target,
deadline_at_millis,
correlation: capability.correlation.clone(),
})
}
pub const fn complete(&mut self) {
self.pending = self.pending.saturating_sub(1);
}
#[must_use]
pub const fn pending(&self) -> usize {
self.pending
}
}
fn capability_use_key(capability: &BrokerCapability) -> String {
let c = &capability.correlation;
format!(
"{}:{}:{}:{}:{}:{}:{}:{}",
c.tenant,
c.conversation,
c.execution,
c.attempt,
c.step,
c.tool_call,
c.audience,
capability.target,
)
}
#[cfg(test)]
mod tests {
#![allow(clippy::pedantic, clippy::nursery, missing_docs)]
use std::collections::BTreeMap;
use polyc_state::{command::FencingToken, id::AttemptId};
use super::*;
use crate::execution::{ExecutionId, ExecutionIdentity};
#[derive(Default)]
struct Registry(BTreeMap<(BrokerFamily, String), TrustedTarget>);
impl TrustedRegistry for Registry {
fn resolve(&self, family: BrokerFamily, target: &str) -> Option<TrustedTarget> {
self.0.get(&(family, target.to_owned())).cloned()
}
}
fn label() -> ExecutionLabel {
ExecutionLabel::new(
ExecutionIdentity::new(
ExecutionId::new("execution"),
AttemptId::new("attempt"),
FencingToken::new(4),
)
.expect("identity"),
2,
)
}
fn correlation() -> BrokerCorrelation {
BrokerCorrelation {
tenant: "tenant-a".to_owned(),
conversation: "conversation-a".to_owned(),
execution: "execution".to_owned(),
attempt: "attempt".to_owned(),
step: label().step().id().as_str().to_owned(),
tool_call: "call-7".to_owned(),
audience: "connector-read".to_owned(),
}
}
fn capability() -> BrokerCapability {
BrokerCapability {
correlation: correlation(),
admitted_label: label(),
family: BrokerFamily::Connector,
target: "calendar".to_owned(),
method: "GET".to_owned(),
resource: "/v1/events".to_owned(),
max_request_bytes: 128,
expires_at_millis: 1_000,
}
}
fn request() -> BrokerRequest {
BrokerRequest {
correlation: correlation(),
label: label(),
family: BrokerFamily::Connector,
target: "calendar".to_owned(),
method: "GET".to_owned(),
resource: "/v1/events".to_owned(),
encoded_request_bytes: 8,
deadline_at_millis: 900,
}
}
fn registry() -> Registry {
let mut registry = Registry::default();
registry.0.insert(
(BrokerFamily::Connector, "calendar".to_owned()),
TrustedTarget {
target: "calendar".to_owned(),
family: BrokerFamily::Connector,
method: "GET".to_owned(),
resource: "/v1/events".to_owned(),
read_only: true,
classification: ReturnedContent::Untrusted,
},
);
registry
}
fn session() -> BrokerSession {
BrokerSession::new(BrokerLimits {
max_calls: 2,
max_pending: 1,
session_deadline_at_millis: 800,
broker_deadline_at_millis: 700,
})
}
#[test]
fn binds_every_identity_field_and_uses_the_smallest_deadline() {
let mut session = session();
let admission = session
.admit(&capability(), &request(), ®istry(), 100)
.expect("the exact trusted request is admitted");
assert_eq!(admission.deadline_at_millis, 700);
assert_eq!(admission.target.classification, ReturnedContent::Untrusted);
}
#[test]
fn changed_bound_fields_refuse_before_registry_resolution() {
let mut changed = request();
changed.correlation.tenant = "tenant-b".to_owned();
assert_eq!(
session().admit(&capability(), &changed, ®istry(), 100),
Err(BrokerRefusal::BindingMismatch)
);
}
#[test]
fn a_label_outside_the_admitted_session_is_not_freshness_authority() {
let mut changed = request();
changed.label = ExecutionLabel::new(
ExecutionIdentity::new(
ExecutionId::new("execution"),
AttemptId::new("other-attempt"),
FencingToken::new(5),
)
.expect("identity"),
2,
);
assert_eq!(
session().admit(&capability(), &changed, ®istry(), 100),
Err(BrokerRefusal::UnadmittedLabel)
);
}
#[test]
fn a_one_use_capability_cannot_start_a_second_call_after_completion() {
let mut session = session();
session
.admit(&capability(), &request(), ®istry(), 100)
.expect("first call");
session.complete();
assert_eq!(
session.admit(&capability(), &request(), ®istry(), 100),
Err(BrokerRefusal::CapabilityReplayed)
);
}
#[test]
fn unregistered_mutating_or_out_of_scope_targets_fail_closed() {
let mut changed_request = request();
changed_request.method = "POST".to_owned();
assert_eq!(
session().admit(&capability(), &changed_request, ®istry(), 100),
Err(BrokerRefusal::BindingMismatch)
);
let mut mutating = registry();
mutating
.0
.get_mut(&(BrokerFamily::Connector, "calendar".to_owned()))
.expect("target")
.read_only = false;
assert_eq!(
session().admit(&capability(), &request(), &mutating, 100),
Err(BrokerRefusal::ScopeMismatch)
);
}
#[test]
fn bounds_expiry_pending_and_deadline_are_enforced_before_a_dial() {
let mut oversized = request();
oversized.encoded_request_bytes = 129;
assert_eq!(
session().admit(&capability(), &oversized, ®istry(), 100),
Err(BrokerRefusal::RequestTooLarge)
);
assert_eq!(
session().admit(&capability(), &request(), ®istry(), 1_000),
Err(BrokerRefusal::Expired)
);
assert_eq!(
session().admit(&capability(), &request(), ®istry(), 700),
Err(BrokerRefusal::DeadlineExceeded)
);
let mut session = session();
session
.admit(&capability(), &request(), ®istry(), 100)
.expect("first pending call");
let mut second = capability();
second.correlation.tool_call = "call-8".to_owned();
assert_eq!(
session.admit(
&second,
&BrokerRequest {
correlation: second.correlation.clone(),
..request()
},
®istry(),
100
),
Err(BrokerRefusal::CapacityExceeded)
);
assert_eq!(session.pending(), 1);
}
}