#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
pub enum ConsumerAckKind {
TableRow,
DeliveryAck,
OffsetCommit,
StreamAck,
HttpResponse,
InProcess,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
pub enum KnativeIntegrationKind {
Native,
CustomBridge,
None,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, serde::Serialize, serde::Deserialize)]
pub struct TransportCapabilities {
pub durable_receive: bool,
pub publish_confirm: bool,
pub platform_managed_retry: bool,
pub consumer_ack: ConsumerAckKind,
pub knative_integration: KnativeIntegrationKind,
}
impl TransportCapabilities {
pub const fn postgres() -> Self {
Self {
durable_receive: true,
publish_confirm: true,
platform_managed_retry: false,
consumer_ack: ConsumerAckKind::TableRow,
knative_integration: KnativeIntegrationKind::CustomBridge,
}
}
pub const fn rabbitmq() -> Self {
Self {
durable_receive: true,
publish_confirm: true,
platform_managed_retry: false,
consumer_ack: ConsumerAckKind::DeliveryAck,
knative_integration: KnativeIntegrationKind::Native,
}
}
pub const fn kafka() -> Self {
Self {
durable_receive: true,
publish_confirm: true,
platform_managed_retry: false,
consumer_ack: ConsumerAckKind::OffsetCommit,
knative_integration: KnativeIntegrationKind::Native,
}
}
pub const fn nats_jetstream() -> Self {
Self {
durable_receive: true,
publish_confirm: true,
platform_managed_retry: false,
consumer_ack: ConsumerAckKind::StreamAck,
knative_integration: KnativeIntegrationKind::CustomBridge,
}
}
pub const fn knative() -> Self {
Self {
durable_receive: false,
publish_confirm: true,
platform_managed_retry: true,
consumer_ack: ConsumerAckKind::HttpResponse,
knative_integration: KnativeIntegrationKind::Native,
}
}
pub const fn in_memory() -> Self {
Self {
durable_receive: false,
publish_confirm: false,
platform_managed_retry: false,
consumer_ack: ConsumerAckKind::InProcess,
knative_integration: KnativeIntegrationKind::None,
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn postgres_is_durable_with_custom_knative_bridge() {
let caps = TransportCapabilities::postgres();
assert!(caps.durable_receive);
assert!(caps.publish_confirm);
assert!(!caps.platform_managed_retry);
assert_eq!(caps.consumer_ack, ConsumerAckKind::TableRow);
assert_eq!(
caps.knative_integration,
KnativeIntegrationKind::CustomBridge
);
}
#[test]
fn knative_delegates_retry_to_the_platform() {
let caps = TransportCapabilities::knative();
assert!(caps.platform_managed_retry);
assert!(!caps.durable_receive);
assert_eq!(caps.consumer_ack, ConsumerAckKind::HttpResponse);
assert_eq!(caps.knative_integration, KnativeIntegrationKind::Native);
}
#[test]
fn in_memory_is_the_only_profile_without_durable_receive_or_publish_confirm() {
let caps = TransportCapabilities::in_memory();
assert!(!caps.durable_receive);
assert!(!caps.publish_confirm);
assert_eq!(caps.consumer_ack, ConsumerAckKind::InProcess);
assert_eq!(caps.knative_integration, KnativeIntegrationKind::None);
}
#[test]
fn direct_transports_own_their_retry() {
for caps in [
TransportCapabilities::postgres(),
TransportCapabilities::rabbitmq(),
TransportCapabilities::kafka(),
TransportCapabilities::nats_jetstream(),
] {
assert!(
!caps.platform_managed_retry,
"direct transports own retry; only Knative is platform-managed"
);
assert!(caps.durable_receive, "direct transports receive durably");
}
}
#[test]
fn each_known_transport_has_a_distinct_ack_kind() {
let acks = [
TransportCapabilities::postgres().consumer_ack,
TransportCapabilities::rabbitmq().consumer_ack,
TransportCapabilities::kafka().consumer_ack,
TransportCapabilities::nats_jetstream().consumer_ack,
TransportCapabilities::knative().consumer_ack,
TransportCapabilities::in_memory().consumer_ack,
];
for (i, a) in acks.iter().enumerate() {
for b in &acks[i + 1..] {
assert_ne!(a, b, "ack kinds should be distinct across transports");
}
}
}
#[test]
fn knative_integration_matches_the_evaluation_matrix() {
use KnativeIntegrationKind::*;
assert_eq!(
TransportCapabilities::rabbitmq().knative_integration,
Native
);
assert_eq!(TransportCapabilities::kafka().knative_integration, Native);
assert_eq!(TransportCapabilities::knative().knative_integration, Native);
assert_eq!(
TransportCapabilities::postgres().knative_integration,
CustomBridge
);
assert_eq!(
TransportCapabilities::nats_jetstream().knative_integration,
CustomBridge
);
assert_eq!(TransportCapabilities::in_memory().knative_integration, None);
}
#[test]
fn capabilities_round_trip_through_serde() {
let caps = TransportCapabilities::kafka();
let json = serde_json::to_string(&caps).unwrap();
let restored: TransportCapabilities = serde_json::from_str(&json).unwrap();
assert_eq!(caps, restored);
}
}