use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::task::{Context, Poll};
use camel_api::{Body, CamelError, Exchange};
use camel_component_api::RuntimeObservability;
use camel_language_minijinja::ResolvedLimits;
use camel_language_minijinja::engine::build_context_bounded;
use tower::Service;
use crate::template_set::SharedTemplates;
#[derive(Clone)]
pub(crate) struct TemplateProducer {
templates: SharedTemplates,
render_limits: ResolvedLimits,
#[allow(dead_code)] rt: Option<Arc<dyn RuntimeObservability>>,
#[allow(dead_code)] route_id: String,
}
impl TemplateProducer {
#[allow(clippy::too_many_arguments)]
pub(crate) fn new(
templates: SharedTemplates,
render_limits: ResolvedLimits,
rt: Option<Arc<dyn RuntimeObservability>>,
route_id: impl Into<String>,
) -> Self {
Self {
templates,
render_limits,
rt,
route_id: route_id.into(),
}
}
}
impl Service<Exchange> for TemplateProducer {
type Response = Exchange;
type Error = CamelError;
type Future = Pin<Box<dyn Future<Output = Result<Exchange, CamelError>> + Send>>;
fn poll_ready(&mut self, _cx: &mut Context<'_>) -> Poll<Result<(), Self::Error>> {
Poll::Ready(Ok(()))
}
fn call(&mut self, mut exchange: Exchange) -> Self::Future {
let producer = self.clone();
Box::pin(async move {
producer.render_into(&mut exchange).await?;
Ok(exchange)
})
}
}
impl TemplateProducer {
#[allow(dead_code)] async fn render_into(&self, exchange: &mut Exchange) -> Result<(), CamelError> {
if matches!(exchange.input.body, Body::Stream(_)) {
return Err(CamelError::ProcessorError(
"template producer cannot render Body::Stream; add `stream_cache` upstream"
.to_string(),
));
}
let context = build_context_bounded(exchange, self.render_limits.max_context_size)
.map_err(|e| CamelError::ProcessorError(e.to_string()))?;
let set = self.templates.load_full();
let rendered = set
.render_entry(context, self.render_limits)
.await
.map_err(|e| CamelError::ProcessorError(e.to_string()))?;
exchange.input.body = Body::from(rendered);
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
use std::sync::Arc;
use arc_swap::ArcSwap;
use camel_api::{Message, Value};
use camel_language_api::MinijinjaLimitsConfig;
use serde_json::json;
use crate::closure::ClosureSnapshot;
use crate::template_set::TemplateSet;
fn make_templates(entry_name: &str, source: &str) -> SharedTemplates {
let snap = ClosureSnapshot::from_single_entry(entry_name, source.as_bytes().to_vec());
let set = TemplateSet::compile(&snap, entry_name, MinijinjaLimitsConfig::default())
.expect("compile");
Arc::new(ArcSwap::from_pointee(set))
}
fn make_exchange_with_body(body: &str) -> Exchange {
let mut msg = Message::new(Body::from(body.to_string()));
msg.headers
.insert("X-Keep".to_string(), Value::String("kept".to_string()));
Exchange::new(msg)
}
fn body_string(exchange: &Exchange) -> Option<String> {
match &exchange.input.body {
Body::Text(s) => Some(s.clone()),
Body::Bytes(b) => Some(String::from_utf8_lossy(b).to_string()),
_ => None,
}
}
#[tokio::test]
async fn producer_replaces_body_on_success() {
let templates = make_templates(
"echo.html",
r#"{% autoescape "none" %}body={{body}}{% endautoescape %}"#,
);
let mut producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let exchange = make_exchange_with_body("hello");
let result = producer.call(exchange).await.expect("call ok");
assert_eq!(
body_string(&result).as_deref(),
Some("body=hello"),
"body must be replaced by rendered output"
);
assert_eq!(
result.input.headers.get("X-Keep").and_then(|v| match v {
Value::String(s) => Some(s.as_str()),
_ => None,
}),
Some("kept"),
"headers must be preserved across render"
);
}
#[tokio::test]
async fn producer_leaves_body_on_render_error() {
let templates = make_templates(
"boom.html",
r#"{% autoescape "none" %}{{undefined_var}}{% endautoescape %}"#,
);
let producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let original_body = "original-payload".to_string();
let mut exchange = Exchange::new(Message::new(Body::from(original_body.clone())));
exchange
.input
.headers
.insert("X-Keep".to_string(), Value::String("kept".to_string()));
let result = producer.render_into(&mut exchange).await;
assert!(
matches!(result, Err(CamelError::ProcessorError(_))),
"strict-undefined is a data-plane render error and must surface as \
CamelError::ProcessorError (NOT CamelError::TemplateReload, which is \
reserved for the control-plane startup-build and hot-reload paths in \
lifecycle.rs / reload.rs), got: {result:?}"
);
match &exchange.input.body {
Body::Text(s) => {
assert_eq!(s, &original_body, "body must be unchanged on render error")
}
other => panic!("expected unchanged Text body, got: {other:?}"),
}
assert_eq!(
exchange.input.headers.get("X-Keep").and_then(|v| match v {
Value::String(s) => Some(s.as_str()),
_ => None,
}),
Some("kept"),
);
}
#[tokio::test]
async fn producer_rejects_oversize_context() {
let small_limits = ResolvedLimits {
max_context_size: 16,
..ResolvedLimits::default()
};
let templates = make_templates(
"echo.html",
r#"{% autoescape "none" %}{{body}}{% endautoescape %}"#,
);
let producer = TemplateProducer::new(templates, small_limits, None, "test-route");
let original_body = "x".repeat(64);
let mut exchange = Exchange::new(Message::new(Body::from(original_body.clone())));
let result = producer.render_into(&mut exchange).await;
assert!(
matches!(result, Err(CamelError::ProcessorError(_))),
"S9 context-overflow must map to CamelError::ProcessorError, got: {result:?}"
);
match &exchange.input.body {
Body::Text(s) => assert_eq!(
s, &original_body,
"body must be byte-identical after S9 rejection"
),
other => panic!("expected unchanged Text body, got: {other:?}"),
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn producer_handles_concurrent_renders() {
let templates = make_templates(
"echo.html",
r#"{% autoescape "none" %}user={{body}}{% endautoescape %}"#,
);
let producer = TemplateProducer::new(
templates,
ResolvedLimits::default(),
None,
"test-concurrent-route",
);
const N: usize = 20;
let mut handles = Vec::with_capacity(N);
for i in 0..N {
let mut producer = producer.clone();
let body = format!("user-{i}");
handles.push(tokio::spawn(async move {
let exchange = make_exchange_with_body(&body);
let result = producer.call(exchange).await.expect("call ok");
let rendered = body_string(&result).expect("rendered body");
let expected = format!("user=user-{i}");
assert_eq!(
rendered, expected,
"task {i} observed wrong body — concurrent render raced"
);
expected
}));
}
let outputs = futures::future::join_all(handles).await;
let mut seen: Vec<String> = Vec::with_capacity(N);
for (i, join) in outputs.into_iter().enumerate() {
seen.push(
join.unwrap_or_else(|e| panic!("task {i} panicked during concurrent render: {e}")),
);
}
let mut expected: Vec<String> = (0..N).map(|i| format!("user=user-{i}")).collect();
seen.sort();
expected.sort();
assert_eq!(
seen, expected,
"concurrent renders must not drop, duplicate, or cross bodies"
);
}
#[tokio::test]
async fn producer_ignores_override_header() {
let templates = make_templates(
"op-entry",
r#"{% autoescape "none" %}{{body}}{% endautoescape %}"#,
);
let mut producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let mut exchange = make_exchange_with_body("op-output");
exchange.input.headers.insert(
"X-Template-File".to_string(),
Value::String("/etc/passwd".to_string()),
);
exchange.input.headers.insert(
"CamelTemplateRoot".to_string(),
Value::String("evil-root".to_string()),
);
let result = producer.call(exchange).await.expect("call ok");
assert_eq!(
body_string(&result).as_deref(),
Some("op-output"),
"producer must render the operator-configured entry, \
ignoring X-Template-File / CamelTemplateRoot headers"
);
let keep: HashMap<_, _> = result
.input
.headers
.iter()
.map(|(k, v)| (k.as_str(), v.clone()))
.collect();
assert!(keep.contains_key("X-Template-File"));
}
#[tokio::test]
async fn producer_fails_closed_on_output_cap() {
const OUTPUT_CAP: usize = 1024;
let tight_limits = ResolvedLimits {
max_output_size: OUTPUT_CAP,
..ResolvedLimits::default()
};
let source = r#"{% autoescape "none" %}{% for i in range(0, 10000) %}x{% endfor %}{% endautoescape %}"#;
let templates = make_templates("flood.html", source);
let producer = TemplateProducer::new(templates, tight_limits, None, "test-output-cap");
let original_body = "original-payload".to_string();
let mut exchange = Exchange::new(Message::new(Body::from(original_body.clone())));
let result = producer.render_into(&mut exchange).await;
let err_msg = match &result {
Err(CamelError::ProcessorError(msg)) => msg.clone(),
other => panic!(
"output-cap overflow must surface as CamelError::ProcessorError, got: {other:?}"
),
};
assert!(
err_msg.contains("render:"),
"output-cap failure must originate from the render path \
(render_entry prefix); got: {err_msg}"
);
match &exchange.input.body {
Body::Text(s) => assert_eq!(
s, &original_body,
"body must be byte-identical after S6 output-cap rejection"
),
other => panic!("expected unchanged Text body, got: {other:?}"),
}
}
#[tokio::test]
async fn producer_exposes_structured_json_body_fields() {
let templates = make_templates(
"dom.html",
r#"{% autoescape "none" %}{{ body.first }} {{ body.last }}{% endautoescape %}"#,
);
let producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let mut exchange = Exchange::new(Message::new(Body::Json(json!({
"first": "Claus",
"last": "Ibsen"
}))));
producer
.render_into(&mut exchange)
.await
.expect("structured JSON body must render");
assert_eq!(
body_string(&exchange).as_deref(),
Some("Claus Ibsen"),
"Body::Json fields must be addressable as body.<field> in the template"
);
}
#[tokio::test]
async fn producer_renders_static_template_with_empty_body() {
let templates = make_templates(
"static.html",
r#"{% autoescape "none" %}static-only, no body ref{% endautoescape %}"#,
);
let producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let mut exchange = Exchange::new(Message::new(Body::Empty));
producer
.render_into(&mut exchange)
.await
.expect("a template that never touches body must render with Body::Empty");
assert_eq!(
body_string(&exchange).as_deref(),
Some("static-only, no body ref"),
"static text must render verbatim regardless of an empty body"
);
}
#[tokio::test]
async fn producer_renders_empty_string_for_referenced_empty_body() {
let templates = make_templates(
"echo.html",
r#"{% autoescape "none" %}{{ body }}{% endautoescape %}"#,
);
let producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let mut exchange = Exchange::new(Message::new(Body::Empty));
producer.render_into(&mut exchange).await.expect(
"referencing an empty-string body does NOT trip strict-undefined; \
the variable is defined",
);
assert_eq!(
body_string(&exchange).as_deref(),
Some(""),
"Body::Empty renders an empty string (rc-wnqj reversal), \
NOT the literal 'none' and NOT a ProcessorError"
);
}
#[tokio::test]
async fn producer_renders_bytes_body_lossy_utf8() {
let templates = make_templates(
"echo.html",
r#"{% autoescape "none" %}{{ body }}{% endautoescape %}"#,
);
let producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let mut exchange = Exchange::new(Message::new(Body::Bytes(
"héllo".as_bytes().to_vec().into(),
)));
producer
.render_into(&mut exchange)
.await
.expect("valid UTF-8 bytes body must render");
assert_eq!(
body_string(&exchange).as_deref(),
Some("héllo"),
"valid UTF-8 bytes must render verbatim"
);
let mut exchange = Exchange::new(Message::new(Body::Bytes(
vec![0xff, 0xfe, 0x00].into(),
)));
producer
.render_into(&mut exchange)
.await
.expect("invalid UTF-8 bytes must render via lossy replacement, not error");
let rendered = body_string(&exchange).expect("rendered body must be text");
assert!(
rendered.contains('\u{FFFD}'),
"invalid UTF-8 bytes must surface as U+FFFD replacement chars (lossy contract); \
got {rendered:?}"
);
}
#[tokio::test]
async fn producer_exposes_exchange_property_key() {
let templates = make_templates(
"prop.html",
r#"{% autoescape "none" %}{{ exchangeProperty.item }}{% endautoescape %}"#,
);
let producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let mut exchange = Exchange::new(Message::new(Body::Empty));
exchange.set_property("item", Value::String("7".into()));
producer
.render_into(&mut exchange)
.await
.expect("exchangeProperty.<key> must resolve from exchange properties");
assert_eq!(
body_string(&exchange).as_deref(),
Some("7"),
"exchange properties must be addressable as exchangeProperty.<name>"
);
}
#[tokio::test]
async fn producer_rejects_unknown_function_call() {
let templates = make_templates(
"evil.html",
r#"{% autoescape "none" %}{{ evil_global_function() }}{% endautoescape %}"#,
);
let producer =
TemplateProducer::new(templates, ResolvedLimits::default(), None, "test-route");
let original_body = "original-payload".to_string();
let mut exchange = Exchange::new(Message::new(Body::from(original_body.clone())));
let result = producer.render_into(&mut exchange).await;
assert!(
matches!(result, Err(CamelError::ProcessorError(_))),
"unknown function call must surface as CamelError::ProcessorError \
(data-plane render error), got: {result:?}"
);
match &exchange.input.body {
Body::Text(s) => assert_eq!(
s, &original_body,
"body must be byte-identical after S8 unknown-function rejection"
),
other => panic!("expected unchanged Text body, got: {other:?}"),
}
}
}