orion-server 1.11.0

Turn business logic into live REST/Kafka services, declared as JSON
//! The one step between "the guards admitted this" and "shape a response".
//!
//! `channel::guards::apply_guards` is the single enforcement point for
//! ingress, and everything after it used to be copied per transport: the sync
//! HTTP route, the async trace queue, the Kafka consume loop and the
//! in-process `channel_call` handler each built a `Message`, loaded the engine
//! snapshot, called [`run_for_channel`], re-derived the "engine returned Ok
//! but the message carries task errors" rule, and mapped a timeout — **four
//! different ways**: an `OrionError::Timeout`, a synthesised
//! `DataflowError::Timeout`, a `fail("timeout", …)`, and another
//! `DataflowError::Timeout` with different wording. Nothing tied the four
//! together, so the next change to dataflow-rs's error semantics was four
//! edits and a test suite that could not tell if one was missed.
//!
//! They had already drifted. Kafka alone omitted the rollout bucket (correctly
//! — a record has no sticky caller identity), and `channel_call` never applied
//! the `has_errors` rule at all, so a target channel whose workflow failed its
//! tasks reported success to its caller.
//!
//! [`execute_admitted`] is that step, once. What stays with each transport is
//! what genuinely differs: persistence, response shaping, and what a
//! [`RunOutcome`] *means* to it — an HTTP caller gets the errors in its
//! envelope, the async path routes them to the DLQ, Kafka rewinds or
//! dead-letters.

use std::sync::Arc;
use std::time::{Duration, Instant};

use serde_json::Value;

use super::runner::{TraceCapture, run_for_channel};

/// What one admitted execution did.
///
/// The four outcomes are exhaustive over what the engine can report, which is
/// the point: a new dataflow-rs error shape lands here, once, and every
/// transport's `match` fails to compile until it is handled.
#[derive(Debug)]
pub enum RunOutcome {
    /// Every workflow that matched ran, and none reported a task error.
    Ok,
    /// The engine returned `Ok(())` and the message carries task errors — the
    /// v3 contract. Carries the joined `code: message` summary, which is what
    /// all three of the transports that look at it were building by hand.
    WorkflowErrors(String),
    /// The deadline expired, with the budget in milliseconds.
    Timeout(u64),
    /// The engine itself failed the call.
    EngineError(dataflow_rs::DataflowError),
}

impl RunOutcome {
    /// The `status` label for `orion_messages_total`.
    ///
    /// One derivation, so the counter means the same thing on every transport.
    /// It did not: the sync route counted [`Self::WorkflowErrors`] as `ok`
    /// while the Kafka and async paths counted it as `error`, so a channel
    /// whose workflow failed every synchronous request reported a 100% success
    /// rate. `error` is the answer the metric's own documentation implies —
    /// "messages processed, by outcome" — and a workflow that failed its tasks
    /// did not achieve its outcome.
    ///
    /// The label is derived here; *emitting* it stays with the transport,
    /// because whether an execution should be counted at all is knowledge only
    /// the transport has. Kafka suppresses the counter on an in-place retry
    /// (K10 — one poison record otherwise inflated the error rate once a
    /// minute forever), and `channel_call` emits nothing because the request
    /// that triggered it was already counted at its own ingress.
    pub fn status_label(&self) -> &'static str {
        match self {
            Self::Ok => "ok",
            Self::WorkflowErrors(_) | Self::EngineError(_) => "error",
            Self::Timeout(_) => "timeout",
        }
    }

    /// Whether the run produced the result the caller asked for.
    pub fn is_ok(&self) -> bool {
        matches!(self, Self::Ok)
    }
}

/// Everything about one execution that is not the channel or the payload.
#[derive(Default)]
pub struct ExecOpts<'a> {
    /// Deadline for the whole call. `None` runs untimed.
    pub timeout_ms: Option<u64>,
    /// Per-task trace capture, when the channel opted in via
    /// `config.tracing.task_details`. Also decides the message's
    /// `capture_changes`: only a traced run pays for the per-write value copies.
    pub capture: Option<TraceCapture>,
    /// The rollout bucket to stamp on the message, from
    /// [`crate::engine::utils::rollout_bucket_for_identity`].
    ///
    /// `None` means "admitted by every workflow, rollout or not", and is
    /// deliberate on two transports: a Kafka record has no sticky caller
    /// identity and no forwarded address, so a random bucket per record would
    /// split one topic's traffic across canary versions non-deterministically;
    /// an in-process `channel_call` is already running inside its caller's
    /// bucket.
    pub routing_bucket: Option<u8>,
    /// The per-request profile scope, when one is in flight.
    pub profile: Option<&'a Arc<super::profile::ProfileCollector>>,
}

/// One execution and everything a transport needs from it.
pub struct Execution {
    /// The message as the workflows left it — the source of the response body,
    /// the persisted result, and the error list.
    pub message: dataflow_rs::Message,
    /// The per-task trace, when [`ExecOpts::capture`] asked for one.
    pub task_trace: Option<dataflow_rs::ExecutionTrace>,
    pub outcome: RunOutcome,
    /// Wall time of the engine call alone, for the latency histogram.
    pub duration: Duration,
}

/// Run `channel`'s workflows over `data`, after the guards have admitted it.
///
/// `engine` is the one belonging to the generation that admitted the request
/// (`RuntimeGeneration::engine`) — passed in rather than loaded here, so a
/// transport cannot admit against one generation and execute on the next. It
/// used to take the engine *handle* and load it at this line, which is exactly
/// where that mismatch entered.
///
/// Owns the message build, the deadline arm and the `has_errors` rule.
/// Everything downstream — persisting a trace, shaping a response, committing
/// an offset, emitting the counter whose label [`RunOutcome::status_label`]
/// derives — is the caller's.
pub async fn execute_admitted(
    engine: &Arc<dataflow_rs::Engine>,
    channel: &str,
    data: &Value,
    metadata: &Value,
    opts: ExecOpts<'_>,
) -> Execution {
    // Per-write capture is on only for a traced run. With it on, every write
    // deep-copies its old and new value into the audit trail, and those copies
    // stay on the message until the run returns — so a looping workflow holds
    // sweeps × writes × value size (#350; about 65 bytes per number written).
    // Nothing on this path reads `AuditTrail::changes`; the trace's per-step
    // diff does, and `TraceOptions { changes: true }` only reports what was
    // captured, it never turns capture on. The audit entries themselves are
    // recorded either way.
    let mut builder = dataflow_rs::Message::builder()
        .payload_json(data)
        .metadata_json(metadata)
        .capture_changes(opts.capture.is_some());
    if let Some(bucket) = opts.routing_bucket {
        builder = builder.routing_bucket(bucket);
    }
    let mut message = builder.build();

    let started = Instant::now();
    let call = run_for_channel(
        engine,
        channel,
        &mut message,
        opts.timeout_ms,
        opts.profile,
        opts.capture,
    )
    .await;
    let duration = started.elapsed();

    let (outcome, task_trace) = match call {
        Err(ms) => (RunOutcome::Timeout(ms), None),
        Ok((Err(e), trace)) => (RunOutcome::EngineError(e), trace),
        Ok((Ok(()), trace)) => {
            // The v3 contract: `process_message_for_channel` answers `Ok(())`
            // even when individual workflows failed — the failures are pushed
            // into `message.errors()`. Derived here so a transport cannot
            // forget to look, which is what `channel_call` was doing.
            if message.has_errors() {
                (RunOutcome::WorkflowErrors(error_summary(&message)), trace)
            } else {
                (RunOutcome::Ok, trace)
            }
        }
    };

    Execution {
        message,
        task_trace,
        outcome,
        duration,
    }
}

/// The `code: message; code: message` summary three transports were each
/// building from `message.errors()`.
fn error_summary(message: &dataflow_rs::Message) -> String {
    message
        .errors()
        .iter()
        .map(|e| format!("{}: {}", e.code, e.message))
        .collect::<Vec<_>>()
        .join("; ")
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn the_status_label_is_one_derivation() {
        assert_eq!(RunOutcome::Ok.status_label(), "ok");
        assert_eq!(RunOutcome::Timeout(50).status_label(), "timeout");
        // The unification: workflow errors are an `error` on every transport,
        // where the sync route used to call them `ok`.
        assert_eq!(
            RunOutcome::WorkflowErrors("boom".into()).status_label(),
            "error"
        );
        assert_eq!(
            RunOutcome::EngineError(dataflow_rs::DataflowError::Unknown("x".into())).status_label(),
            "error"
        );
    }

    /// #350: a looping workflow run without a trace keeps one audit entry per
    /// task per sweep but no value copies, so its memory does not grow with the
    /// sweeps. A traced run still captures them — its per-step diff is built
    /// from them.
    #[tokio::test]
    async fn only_a_traced_run_captures_per_write_changes() {
        let workflow = dataflow_rs::Workflow::from_json(
            r#"{"id":"w","name":"w","channel":"c","priority":0,"condition":true,
                "loop":{"counter":"i","max":3},
                "tasks":[{"id":"m","name":"m","function":{"name":"map","input":
                  {"mappings":[{"path":"temp_data.seen","logic":{"var":"temp_data.i"}}]}}}]}"#,
        )
        .expect("workflow parses");
        let engine = Arc::new(
            dataflow_rs::Engine::new(vec![workflow], std::collections::HashMap::new())
                .expect("engine builds"),
        );
        let data = serde_json::json!({"xs": [1, 2, 3]});
        let metadata = serde_json::json!({});

        let untraced = execute_admitted(&engine, "c", &data, &metadata, ExecOpts::default()).await;
        assert!(untraced.outcome.is_ok());
        let trail = untraced.message.audit_trail();
        assert_eq!(trail.len(), 3, "one entry per sweep is still recorded");
        assert!(
            trail.iter().all(|entry| entry.changes.is_empty()),
            "an untraced run holds no value copies"
        );

        let traced = execute_admitted(
            &engine,
            "c",
            &data,
            &metadata,
            ExecOpts {
                capture: Some(TraceCapture {
                    max_snapshot_bytes: 0,
                }),
                ..Default::default()
            },
        )
        .await;
        assert!(traced.outcome.is_ok());
        assert!(
            traced
                .message
                .audit_trail()
                .iter()
                .all(|entry| entry.changes.len() == 1),
            "a traced run captures each write"
        );
        let trace = traced.task_trace.expect("a traced run returns its trace");
        assert!(
            trace
                .steps
                .iter()
                .all(|step| step.changes.as_ref().is_some_and(|c| c.len() == 1)),
            "the trace's per-step diff is still populated"
        );
    }
}