use crate::{
AdmissionController, AdmissionError, AdmissionPartition, AdmissionReservation,
AuthorizationFlowQuotaKey, DEFAULT_CANCELLATION_REASON_MAX_BYTES, DEFAULT_CURSOR_MAX_BYTES,
DEFAULT_JSON_RPC_MAX_BODY_BYTES, DEFAULT_METADATA_MAX_BYTES, DEFAULT_METADATA_MAX_ENTRIES,
DEFAULT_URI_MAX_BYTES, PreAuthSourceBucketKey, ProtocolLimit, ProtocolLimits,
ProtocolLimitsError, QuotaPartitionKey, SealedAdmissionKeyError,
};
pub(crate) fn documented_limits() -> ProtocolLimits {
ProtocolLimits::try_new(
DEFAULT_JSON_RPC_MAX_BODY_BYTES,
DEFAULT_METADATA_MAX_ENTRIES,
DEFAULT_METADATA_MAX_BYTES,
DEFAULT_URI_MAX_BYTES,
DEFAULT_CANCELLATION_REASON_MAX_BYTES,
DEFAULT_CURSOR_MAX_BYTES,
)
.expect("the documented defaults must admit")
}
pub(crate) const BOUND_ROW_IDS: [&str; 6] = [
"LIMIT-A-01.01",
"LIMIT-A-01.02",
"LIMIT-A-01.03",
"LIMIT-A-01.04",
"LIMIT-A-01.05",
"LIMIT-A-01.06",
];
pub(crate) fn bound_rows() -> [(&'static str, ProtocolLimit, usize, usize); 6] {
[
(
BOUND_ROW_IDS[0],
ProtocolLimit::JsonRpcBodyBytes,
DEFAULT_JSON_RPC_MAX_BODY_BYTES,
crate::HARD_JSON_RPC_MAX_BODY_BYTES,
),
(
BOUND_ROW_IDS[1],
ProtocolLimit::MetadataEntries,
DEFAULT_METADATA_MAX_ENTRIES as usize,
crate::HARD_METADATA_MAX_ENTRIES as usize,
),
(
BOUND_ROW_IDS[2],
ProtocolLimit::MetadataBytes,
DEFAULT_METADATA_MAX_BYTES,
crate::HARD_METADATA_MAX_BYTES,
),
(
BOUND_ROW_IDS[3],
ProtocolLimit::UriBytes,
DEFAULT_URI_MAX_BYTES,
crate::HARD_URI_MAX_BYTES,
),
(
BOUND_ROW_IDS[4],
ProtocolLimit::CancellationReasonBytes,
DEFAULT_CANCELLATION_REASON_MAX_BYTES,
crate::HARD_CANCELLATION_REASON_MAX_BYTES,
),
(
BOUND_ROW_IDS[5],
ProtocolLimit::CursorBytes,
DEFAULT_CURSOR_MAX_BYTES,
crate::HARD_CURSOR_MAX_BYTES,
),
]
}
pub(crate) fn build_with_override(
limit: ProtocolLimit,
value: usize,
) -> Result<ProtocolLimits, ProtocolLimitsError> {
let pick = |row: ProtocolLimit, documented: usize| -> usize {
if row == limit { value } else { documented }
};
let entries = pick(
ProtocolLimit::MetadataEntries,
DEFAULT_METADATA_MAX_ENTRIES as usize,
);
let entries =
u16::try_from(entries).map_err(|_| ProtocolLimitsError::ExceedsHardCeiling { limit })?;
ProtocolLimits::builder()
.json_rpc_max_body_bytes(pick(
ProtocolLimit::JsonRpcBodyBytes,
DEFAULT_JSON_RPC_MAX_BODY_BYTES,
))
.metadata_max_entries(entries)
.metadata_max_bytes(pick(
ProtocolLimit::MetadataBytes,
DEFAULT_METADATA_MAX_BYTES,
))
.uri_max_bytes(pick(ProtocolLimit::UriBytes, DEFAULT_URI_MAX_BYTES))
.cancellation_reason_max_bytes(pick(
ProtocolLimit::CancellationReasonBytes,
DEFAULT_CANCELLATION_REASON_MAX_BYTES,
))
.cursor_max_bytes(pick(ProtocolLimit::CursorBytes, DEFAULT_CURSOR_MAX_BYTES))
.build()
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct BoundObservation {
pub(crate) id: &'static str,
pub(crate) configured: usize,
pub(crate) ceiling: usize,
pub(crate) generation: u64,
pub(crate) refused_above_ceiling: ProtocolLimitsError,
}
pub(crate) fn run_bound_rows() -> Vec<BoundObservation> {
let accepted = documented_limits();
accepted.validate().expect("the defaults remain valid");
let snapshot = accepted.snapshot();
let mut rows = Vec::new();
for (id, limit, default, ceiling) in bound_rows() {
let at_ceiling =
build_with_override(limit, ceiling).expect("a value at the hard ceiling must admit");
assert_eq!(
at_ceiling
.configured_units(limit)
.expect("a countable row reports its units"),
ceiling,
"{id}: the accepted configuration must carry the ceiling it was given"
);
let refused_above_ceiling = build_with_override(limit, ceiling + 1)
.expect_err("a value above the hard ceiling must refuse");
rows.push(BoundObservation {
id,
configured: snapshot
.configured_units(limit)
.expect("a countable row reports its units"),
ceiling: ProtocolLimits::hard_ceiling(limit)
.expect("a countable row declares a ceiling"),
generation: snapshot.generation(),
refused_above_ceiling,
});
assert_eq!(
rows.last().expect("row just pushed").configured,
default,
"{id}: the configured default must match the documented constant"
);
}
rows
}
pub(crate) fn bounds_receipt(rows: &[BoundObservation]) -> Vec<String> {
let mut receipt = vec![format!(
"LIMIT01-A-BOUNDS-v1 rows={} generation={}",
rows.len(),
rows.first().map_or(0, |row| row.generation)
)];
for row in rows {
receipt.push(format!(
"{} configured={} ceiling={} generation={} refused={:?}",
row.id, row.configured, row.ceiling, row.generation, row.refused_above_ceiling
));
}
receipt
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct PartitionObservation {
pub(crate) id: &'static str,
pub(crate) discriminant: &'static str,
pub(crate) generation: u64,
pub(crate) non_authorizing: bool,
pub(crate) result: Result<(), SealedAdmissionKeyError>,
}
fn refuses_request_identifier(raw: &str) -> bool {
matches!(
AdmissionPartition::try_from_request_identifier(raw),
Err(SealedAdmissionKeyError::RequestSuppliedIdentifier)
)
}
pub(crate) fn run_partition_rows() -> Vec<PartitionObservation> {
let accepted = documented_limits();
let snapshot = accepted.snapshot();
let mut rows = Vec::new();
let pre_auth = AdmissionPartition::pre_auth(
PreAuthSourceBucketKey::from_listener_and_source("mcp.example.test", "tcp:203.0.113.8")
.expect("a transport-observed source is admitted"),
);
assert!(pre_auth.is_pre_auth(), "LIMIT-A-02.01 must remain pre-auth");
assert!(
!pre_auth.is_verified(),
"LIMIT-A-02.01 must not be verified"
);
rows.push(PartitionObservation {
id: "LIMIT-A-02.01",
discriminant: "PreAuth",
generation: snapshot.generation(),
non_authorizing: refuses_request_identifier("tcp:203.0.113.8"),
result: Ok(()),
});
let verified = AdmissionPartition::verified(
QuotaPartitionKey::from_verified_security_facts(
"static-token",
1,
"https://issuer.example.test",
"https://mcp.example.test/mcp",
"tenant-a",
"subject-a",
)
.expect("verified security facts mint a partition key"),
);
assert!(verified.is_verified(), "LIMIT-A-02.02 must be verified");
rows.push(PartitionObservation {
id: "LIMIT-A-02.02",
discriminant: "Verified",
generation: snapshot.generation(),
non_authorizing: refuses_request_identifier("subject-a"),
result: Ok(()),
});
let flow = AdmissionPartition::authorization_flow(
AuthorizationFlowQuotaKey::from_configured_flow(
"https://issuer.example.test",
"https://mcp.example.test/mcp",
"registered-client-1",
"loopback",
"oauth-authorization-code",
)
.expect("a configured flow is admitted"),
);
assert!(
!flow.is_verified() && !flow.is_pre_auth(),
"LIMIT-A-02.03 is neither verified nor pre-auth"
);
rows.push(PartitionObservation {
id: "LIMIT-A-02.03",
discriminant: "AuthorizationFlow",
generation: snapshot.generation(),
non_authorizing: refuses_request_identifier("registered-client-1"),
result: Ok(()),
});
let key_refusal = QuotaPartitionKey::try_from_request_identifier("raw-request-id")
.expect_err("a request-supplied identifier must refuse");
rows.push(PartitionObservation {
id: "LIMIT-A-02.04",
discriminant: "refused",
generation: snapshot.generation(),
non_authorizing: refuses_request_identifier("raw-request-id"),
result: Err(key_refusal),
});
let partition_refusal = AdmissionPartition::try_from_request_identifier("raw-request-id")
.expect_err("a request-supplied identifier must refuse");
rows.push(PartitionObservation {
id: "LIMIT-A-02.05",
discriminant: "refused",
generation: snapshot.generation(),
non_authorizing: refuses_request_identifier("raw-request-id"),
result: Err(partition_refusal),
});
assert_eq!(
snapshot,
accepted.snapshot(),
"the request snapshot must be immutable across its lifecycle"
);
rows
}
pub(crate) fn partitions_receipt(rows: &[PartitionObservation]) -> Vec<String> {
let mut receipt = vec![format!("LIMIT01-A-PARTITIONS-v1 rows={}", rows.len())];
for row in rows {
receipt.push(format!(
"{} discriminant={} generation={} non_authorizing={} result={:?}",
row.id, row.discriminant, row.generation, row.non_authorizing, row.result
));
}
receipt
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ArithmeticObservation {
pub(crate) id: &'static str,
pub(crate) requested: usize,
pub(crate) available: usize,
pub(crate) result: Result<usize, ProtocolLimitsError>,
pub(crate) retained: usize,
}
pub(crate) fn run_arithmetic_rows() -> Vec<ArithmeticObservation> {
let limits = documented_limits();
let limit = ProtocolLimit::MetadataEntries;
let available = limits
.configured_units(limit)
.expect("a countable row reports its units");
let mut retained = 0_usize;
let mut rows = Vec::new();
for (id, current, additional) in [
("LIMIT-A-03.01", 0_usize, available - 1),
("LIMIT-A-03.02", 0_usize, available),
("LIMIT-A-03.03", available, 1_usize),
("LIMIT-A-03.04", usize::MAX, 1_usize),
] {
let result = limits.try_charge(limit, current, additional);
if let Ok(admitted) = result {
retained = admitted;
}
rows.push(ArithmeticObservation {
id,
requested: additional,
available,
result,
retained,
});
}
rows
}
pub(crate) fn arithmetic_receipt(rows: &[ArithmeticObservation]) -> Vec<String> {
let mut receipt = vec![format!("LIMIT01-A-ARITHMETIC-v1 rows={}", rows.len())];
for row in rows {
receipt.push(format!(
"{} requested={} available={} result={:?} retained={}",
row.id, row.requested, row.available, row.result, row.retained
));
}
receipt
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) struct Counters {
pub(crate) global_in_use: usize,
pub(crate) partition_in_use: usize,
pub(crate) committed_work: usize,
pub(crate) release_count: usize,
pub(crate) live: usize,
}
#[rustfmt::skip] pub(crate) fn observe(controller: &AdmissionController, partition: &AdmissionPartition) -> Counters {
Counters {
global_in_use: controller.global_in_use(),
partition_in_use: controller.partition_in_use(partition),
committed_work: controller.committed_work(),
release_count: controller.release_count(),
live: controller.live_reservation_count(),
}
}
pub(crate) fn partition_for(source: &str) -> AdmissionPartition {
AdmissionPartition::pre_auth(
PreAuthSourceBucketKey::from_listener_and_source("mcp.example.test", source)
.expect("a transport-observed source is admitted"),
)
}
pub(crate) const B_GLOBAL_CAPACITY: usize = 6;
pub(crate) const B_PARTITION_CAPACITY: usize = 3;
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct ReserveObservation {
pub(crate) id: &'static str,
pub(crate) ceiling: &'static str,
pub(crate) requested: usize,
pub(crate) state: &'static str,
pub(crate) diagnostic: Option<AdmissionError>,
pub(crate) after: Counters,
}
pub(crate) fn run_reserve_rows() -> Vec<ReserveObservation> {
let snapshot = documented_limits();
let subject = partition_for("tcp:203.0.113.8");
let mut rows = Vec::new();
for (id, requested) in [
("LIMIT-B-01.01", B_PARTITION_CAPACITY - 1),
("LIMIT-B-01.02", B_PARTITION_CAPACITY),
("LIMIT-B-01.03", B_PARTITION_CAPACITY + 1),
] {
let controller = AdmissionController::with_capacities(
snapshot.snapshot(),
B_GLOBAL_CAPACITY,
B_PARTITION_CAPACITY,
)
.expect("declared capacities are positive");
let attempt = controller.reserve(subject.clone(), requested);
let (state, diagnostic, held) = match attempt {
Ok(reservation) => ("held", None, Some(reservation)),
Err(error) => ("none", Some(error), None),
};
rows.push(ReserveObservation {
id,
ceiling: "partition",
requested,
state,
diagnostic,
after: observe(&controller, &subject),
});
drop(held);
}
for (id, requested) in [
("LIMIT-B-01.04", B_GLOBAL_CAPACITY - 5),
("LIMIT-B-01.05", B_GLOBAL_CAPACITY - 4),
("LIMIT-B-01.06", B_GLOBAL_CAPACITY - 3),
] {
let controller = AdmissionController::with_capacities(
snapshot.snapshot(),
B_GLOBAL_CAPACITY,
B_PARTITION_CAPACITY,
)
.expect("declared capacities are positive");
let filler_left = partition_for("tcp:203.0.113.20");
let filler_right = partition_for("tcp:203.0.113.21");
let left_hold = controller
.reserve(filler_left, 2)
.expect("global prefill admits");
let right_hold = controller
.reserve(filler_right, 2)
.expect("global prefill admits");
assert_eq!(
controller.global_in_use(),
4,
"{id}: the prefill must leave exactly two global units"
);
let attempt = controller.reserve(subject.clone(), requested);
let (state, diagnostic, held) = match attempt {
Ok(reservation) => ("held", None, Some(reservation)),
Err(error) => ("none", Some(error), None),
};
rows.push(ReserveObservation {
id,
ceiling: "global",
requested,
state,
diagnostic,
after: observe(&controller, &subject),
});
drop(held);
drop(left_hold);
drop(right_hold);
}
rows
}
pub(crate) fn reserve_receipt(rows: &[ReserveObservation]) -> Vec<String> {
let mut receipt = vec![format!(
"LIMIT01-B-RESERVE-v1 rows={} global={B_GLOBAL_CAPACITY} partition={B_PARTITION_CAPACITY}",
rows.len()
)];
for row in rows {
receipt.push(format!(
"{} ceiling={} requested={} state={} global={} partition={} committed={} released={} diagnostic={:?}",
row.id,
row.ceiling,
row.requested,
row.state,
row.after.global_in_use,
row.after.partition_in_use,
row.after.committed_work,
row.after.release_count,
row.diagnostic,
));
}
receipt
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Terminal {
CommitSuccess,
CommitReject,
ExplicitRelease,
Cancellation,
Deadline,
Drop,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct LifecycleObservation {
pub(crate) id: &'static str,
pub(crate) terminal: Terminal,
pub(crate) units: usize,
pub(crate) state: &'static str,
pub(crate) after: Counters,
pub(crate) repeated_commit: Option<Result<(), AdmissionError>>,
pub(crate) repeated_release: Option<Result<(), AdmissionError>>,
}
pub(crate) fn run_lifecycle_rows() -> Vec<LifecycleObservation> {
const UNITS: usize = 2;
let snapshot = documented_limits();
let subject = partition_for("tcp:203.0.113.8");
let mut rows = Vec::new();
for (id, terminal) in [
("LIMIT-B-02.01", Terminal::CommitSuccess),
("LIMIT-B-02.02", Terminal::CommitReject),
("LIMIT-B-02.03", Terminal::ExplicitRelease),
("LIMIT-B-02.04", Terminal::Cancellation),
("LIMIT-B-02.05", Terminal::Deadline),
("LIMIT-B-02.06", Terminal::Drop),
] {
let controller = AdmissionController::with_capacities(
snapshot.snapshot(),
B_GLOBAL_CAPACITY,
B_PARTITION_CAPACITY,
)
.expect("declared capacities are positive");
let mut settled: Option<AdmissionReservation> = None;
let state = match terminal {
Terminal::CommitSuccess => {
let mut reservation = controller
.reserve(subject.clone(), UNITS)
.expect("the commit-success row admits");
reservation.commit().expect("commit transfers the charge");
assert_eq!(
controller.committed_work(),
UNITS,
"{id}: commit must transfer exactly the reserved units"
);
reservation.release().expect("committed work releases once");
settled = Some(reservation);
"committed"
}
Terminal::CommitReject => {
let mut reservation = controller
.reserve(subject.clone(), UNITS)
.expect("the commit-reject row admits");
reservation.release().expect("release before commit");
assert_eq!(
reservation.commit().expect_err("a settled commit rejects"),
AdmissionError::AlreadySettled,
"{id}: committing a settled reservation must reject"
);
settled = Some(reservation);
"released"
}
Terminal::ExplicitRelease => {
let mut reservation = controller
.reserve(subject.clone(), UNITS)
.expect("the explicit-release row admits");
reservation.release().expect("explicit release discharges");
settled = Some(reservation);
"released"
}
Terminal::Cancellation => {
let mut reservation = controller
.reserve(subject.clone(), UNITS)
.expect("the cancellation row admits");
assert_eq!(
controller.global_in_use(),
UNITS,
"{id}: the cancellation row must actually hold a charge first"
);
reservation
.cancel()
.expect("cancel discharges a held reservation");
settled = Some(reservation);
"cancelled"
}
Terminal::Deadline => {
let mut reservation = controller
.reserve_with_deadline(subject.clone(), UNITS, std::time::Instant::now())
.expect("the deadline row admits and holds occupancy");
assert_eq!(
reservation.commit().expect_err("an expired commit rejects"),
AdmissionError::DeadlineExceeded,
"{id}: an expired reservation must refuse its commit"
);
assert_eq!(
controller.global_in_use(),
UNITS,
"{id}: a refused commit must not release the charge"
);
reservation
.release()
.expect("release after a refused commit");
settled = Some(reservation);
"deadline-exceeded"
}
Terminal::Drop => {
{
let _dropped = controller
.reserve(subject.clone(), UNITS)
.expect("the drop row admits");
assert_eq!(
controller.global_in_use(),
UNITS,
"{id}: the drop row must hold a charge before the scope ends"
);
}
"dropped"
}
};
let after = observe(&controller, &subject);
let (repeated_commit, repeated_release) = match settled.as_mut() {
Some(reservation) => (Some(reservation.commit()), Some(reservation.release())),
None => (None, None),
};
assert_eq!(
observe(&controller, &subject),
after,
"{id}: re-applying a terminal event must not move any counter"
);
rows.push(LifecycleObservation {
id,
terminal,
units: UNITS,
state,
after,
repeated_commit,
repeated_release,
});
drop(settled);
}
rows
}
pub(crate) fn lifecycle_receipt(rows: &[LifecycleObservation]) -> Vec<String> {
let mut receipt = vec![format!("LIMIT01-B-LIFECYCLE-v1 rows={}", rows.len())];
for row in rows {
receipt.push(format!(
"{} terminal={:?} units={} state={} global={} partition={} committed={} released={} repeat_commit={:?} repeat_release={:?}",
row.id,
row.terminal,
row.units,
row.state,
row.after.global_in_use,
row.after.partition_in_use,
row.after.committed_work,
row.after.release_count,
row.repeated_commit,
row.repeated_release,
));
}
receipt
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct FairnessObservation {
pub(crate) id: &'static str,
pub(crate) partition: &'static str,
pub(crate) requested: usize,
pub(crate) outcome: &'static str,
pub(crate) diagnostic: Option<AdmissionError>,
pub(crate) global_in_use: usize,
pub(crate) left_in_use: usize,
pub(crate) right_in_use: usize,
pub(crate) release_count: usize,
}
pub(crate) fn run_fairness_rows() -> Vec<FairnessObservation> {
const GLOBAL: usize = 2;
const PARTITION: usize = 2;
let snapshot = documented_limits();
let controller = AdmissionController::with_capacities(snapshot.snapshot(), GLOBAL, PARTITION)
.expect("declared capacities are positive");
let left = partition_for("tcp:203.0.113.10");
let right = partition_for("tcp:203.0.113.11");
let mut rows = Vec::new();
let record = |id: &'static str,
partition: &'static str,
requested: usize,
outcome: &'static str,
diagnostic: Option<AdmissionError>,
controller: &AdmissionController| FairnessObservation {
id,
partition,
requested,
outcome,
diagnostic,
global_in_use: controller.global_in_use(),
left_in_use: controller.partition_in_use(&left),
right_in_use: controller.partition_in_use(&right),
release_count: controller.release_count(),
};
let mut left_hold = controller
.reserve(left.clone(), GLOBAL)
.expect("left saturates the global row");
rows.push(record(
"LIMIT-B-03.01",
"left",
GLOBAL,
"admitted",
None,
&controller,
));
let refusal = controller
.reserve(right.clone(), 1)
.expect_err("a saturated global row refuses the peer partition");
rows.push(record(
"LIMIT-B-03.02",
"right",
1,
"refused",
Some(refusal),
&controller,
));
left_hold.release().expect("left releases its charge");
let mut right_hold = controller
.reserve(right.clone(), 1)
.expect("the freed global units admit the eligible peer");
rows.push(record(
"LIMIT-B-03.03",
"right",
1,
"admitted",
None,
&controller,
));
let over_partition = controller
.reserve(right.clone(), PARTITION)
.expect_err("a charge above the remaining partition ceiling refuses");
rows.push(record(
"LIMIT-B-03.04",
"right",
PARTITION,
"refused",
Some(over_partition),
&controller,
));
right_hold.release().expect("right releases its charge");
rows
}
pub(crate) fn fairness_receipt(rows: &[FairnessObservation]) -> Vec<String> {
let mut receipt = vec![format!("LIMIT01-B-FAIRNESS-v1 rows={}", rows.len())];
for row in rows {
receipt.push(format!(
"{} partition={} requested={} outcome={} global={} left={} right={} released={} diagnostic={:?}",
row.id,
row.partition,
row.requested,
row.outcome,
row.global_in_use,
row.left_in_use,
row.right_in_use,
row.release_count,
row.diagnostic,
));
}
receipt
}