use mq_bridge::models::{Endpoint, Middleware};
use std::collections::HashMap;
const REFERENCE: &str = include_str!("../REFERENCE.md");
struct Fence {
line: usize,
info: String,
body: String,
}
fn fenced_blocks() -> Vec<Fence> {
let mut blocks = Vec::new();
let mut open: Option<(usize, String, Vec<&str>)> = None;
for (index, line) in REFERENCE.lines().enumerate() {
match (line.strip_prefix("```"), open.as_mut()) {
(Some(_), Some(_)) => {
let (line_no, info, body) = open.take().expect("checked open above");
blocks.push(Fence {
line: line_no,
info,
body: body.join("\n"),
});
}
(Some(info), None) => open = Some((index + 1, info.trim().to_string(), Vec::new())),
(None, Some((_, _, body))) => body.push(line),
(None, None) => {}
}
}
if let Some((line_no, info, _)) = open {
panic!("REFERENCE.md has an unterminated ```{info} fence opened at line {line_no}");
}
blocks
}
fn tagged(tag: &str) -> Vec<Fence> {
fenced_blocks()
.into_iter()
.filter(|fence| fence.info == tag)
.collect()
}
fn indent(yaml: &str, spaces: usize) -> String {
let pad = " ".repeat(spaces);
yaml.trim()
.lines()
.map(|line| {
if line.trim().is_empty() {
String::new()
} else {
format!("{pad}{line}")
}
})
.collect::<Vec<_>>()
.join("\n")
}
fn middlewares(entries_yaml: &str, line: usize) -> Vec<Middleware> {
let doc = format!(
"middlewares:\n{}\nmemory: {{ topic: \"reference_probe\" }}\n",
indent(entries_yaml, 2)
);
let endpoint: Endpoint = serde_yaml_ng::from_str(&doc).unwrap_or_else(|e| {
panic!("REFERENCE.md middleware block at line {line} does not parse: {e}\n{doc}")
});
endpoint.middlewares
}
fn endpoint(yaml: &str) -> Endpoint {
serde_yaml_ng::from_str(yaml)
.unwrap_or_else(|e| panic!("REFERENCE.md endpoint snippet does not parse: {e}\n{yaml}"))
}
fn endpoints(body: &str, line: usize) -> Vec<Endpoint> {
let mut found = Vec::new();
for chunk in body.split("\n\n") {
if chunk.trim().is_empty() {
continue;
}
let fragment: HashMap<String, Endpoint> =
serde_yaml_ng::from_str(chunk).unwrap_or_else(|e| {
panic!("REFERENCE.md endpoint block at line {line} does not parse: {e}\n{chunk}")
});
for (key, value) in fragment {
assert!(
key == "input" || key == "output",
"REFERENCE.md endpoint block at line {line} uses key '{key}'; \
expected the fragment to sit under 'input' or 'output'"
);
found.push(value);
}
}
assert!(
!found.is_empty(),
"REFERENCE.md endpoint block at line {line} yielded no endpoints"
);
found
}
#[test]
fn documented_middleware_snippets_parse() {
let blocks = tagged("yaml middleware");
assert!(
blocks.len() >= 15,
"expected the reference to tag a `yaml middleware` block per middleware section, found {}",
blocks.len()
);
for fence in blocks {
let parsed = middlewares(&fence.body, fence.line);
assert!(
!parsed.is_empty(),
"REFERENCE.md middleware block at line {} yielded no middleware",
fence.line
);
}
}
#[test]
fn documented_structural_endpoint_snippets_parse() {
let blocks = tagged("yaml endpoint");
assert!(
blocks.len() >= 6,
"expected the reference to tag a `yaml endpoint` block per structural endpoint section, found {}",
blocks.len()
);
for fence in blocks {
endpoints(&fence.body, fence.line);
}
}
#[test]
fn documented_route_snippets_parse() {
let blocks = tagged("yaml route");
assert!(
blocks.len() >= 5,
"expected the reference to tag its complete route examples, found {}",
blocks.len()
);
for fence in blocks {
let routes: HashMap<String, mq_bridge::models::Route> =
serde_yaml_ng::from_str(&fence.body).unwrap_or_else(|e| {
panic!(
"REFERENCE.md route block at line {} does not parse: {e}\n{}",
fence.line, fence.body
)
});
assert!(
!routes.is_empty(),
"REFERENCE.md route block at line {} yielded no routes",
fence.line
);
}
}
#[test]
fn documented_null_endpoint_spelling_parses() {
assert_eq!(endpoint("null").endpoint_type.name(), "null");
assert_eq!(endpoint("null: null").endpoint_type.name(), "null");
let rejected: Result<Endpoint, _> = serde_yaml_ng::from_str("null: {}");
assert!(
rejected.is_err(),
"`null: {{}}` must stay documented as invalid"
);
}
#[test]
fn omitted_output_defaults_to_null_as_documented() {
let route: mq_bridge::models::Route = serde_yaml_ng::from_str(
r#"
input:
memory: { topic: "in" }
"#,
)
.unwrap();
assert_eq!(route.output.endpoint_type.name(), "null");
}
#[tokio::test]
async fn wrong_side_middleware_matches_documented_behaviour() {
use mq_bridge::models::{EndpointType, WeakJoinMiddleware};
async fn publisher_for(middlewares: Vec<Middleware>) -> anyhow::Result<()> {
let mut output = Endpoint::new_memory("reference_docs_wrong_side", 10);
output.middlewares = middlewares;
let inner = match &output.endpoint_type {
EndpointType::Memory(cfg) => {
mq_bridge::endpoints::memory::MemoryPublisher::new_async(cfg)
.await
.unwrap()
}
_ => unreachable!(),
};
mq_bridge::middleware::apply_middlewares_to_publisher(
Box::new(inner),
&output,
"reference_docs_route",
)
.await
.map(|_| ())
}
let weak_join = publisher_for(vec![Middleware::WeakJoin(WeakJoinMiddleware {
group_by: "correlation_id".to_string(),
expected_count: 2,
timeout_ms: 1000,
branch_by: None,
required: Vec::new(),
on_timeout: Default::default(),
})])
.await;
assert!(
weak_join.is_err(),
"weak_join on an output must fail at startup, as documented"
);
let mut input = Endpoint::new_memory("reference_docs_wrong_side_in", 10);
input.middlewares = vec![Middleware::Retry(Default::default())];
let consumer = match &input.endpoint_type {
EndpointType::Memory(cfg) => mq_bridge::endpoints::memory::MemoryConsumer::new_async(cfg)
.await
.unwrap(),
_ => unreachable!(),
};
assert!(
mq_bridge::middleware::apply_middlewares_to_consumer(
Box::new(consumer),
&input,
"reference_docs_route",
)
.await
.is_ok(),
"retry on an input must warn and be skipped, not fail"
);
}
#[tokio::test]
async fn publisher_middleware_wraps_last_entry_outermost() {
use mq_bridge::models::{
DeadLetterQueueMiddleware, EndpointType, FaultMode, RandomPanicMiddleware,
};
use mq_bridge::CanonicalMessage;
let always_fail = RandomPanicMiddleware {
mode: FaultMode::JsonFormatError,
enabled: true,
..Default::default()
};
let dlq_endpoint = Endpoint::new_memory("reference_docs_dlq", 10);
let mut output = Endpoint::new_memory("reference_docs_main", 10);
output.middlewares = vec![
Middleware::RandomPanic(always_fail),
Middleware::Dlq(Box::new(DeadLetterQueueMiddleware {
endpoint: dlq_endpoint.clone(),
})),
];
let inner = match &output.endpoint_type {
EndpointType::Memory(cfg) => mq_bridge::endpoints::memory::MemoryPublisher::new_async(cfg)
.await
.unwrap(),
_ => unreachable!(),
};
let publisher = mq_bridge::middleware::apply_middlewares_to_publisher(
Box::new(inner),
&output,
"reference_docs_route",
)
.await
.unwrap();
publisher
.send(CanonicalMessage::from("payload"))
.await
.expect("dlq listed last should capture the failure");
assert_eq!(
dlq_endpoint.channel().unwrap().drain_messages().len(),
1,
"the failed message should have been dead-lettered"
);
}