use std::sync::Arc;
use anyhow::{Context, Result};
use async_trait::async_trait;
use aw_event_bridge::{AgentDispatchInvoker, InvokeOutcome, run_bridge, run_bridge_jetstream};
use serde_json::{Value, json};
use crate::dispatch_ledger::{DispatchLedger, NoopDispatchLedger};
use crate::tenant::TenantContext;
use crate::{AgentInput, AgentRuntime};
pub struct RuntimeAgentDispatchInvoker {
runtime: Arc<AgentRuntime>,
ledger: Arc<dyn DispatchLedger>,
}
impl RuntimeAgentDispatchInvoker {
#[must_use]
pub fn new(runtime: Arc<AgentRuntime>) -> Self {
Self {
runtime,
ledger: Arc::new(NoopDispatchLedger),
}
}
#[must_use]
pub fn with_ledger(runtime: Arc<AgentRuntime>, ledger: Arc<dyn DispatchLedger>) -> Self {
Self { runtime, ledger }
}
}
fn extract_user_text(input: &Value) -> String {
input
.get("user_text")
.or_else(|| input.get("text"))
.and_then(Value::as_str)
.map(str::to_string)
.or_else(|| input.as_str().map(str::to_string))
.unwrap_or_default()
}
fn resolve_session_id(input: &Value, idempotency_key: Option<&str>) -> String {
input
.get("session_id")
.and_then(Value::as_str)
.map(str::to_string)
.or_else(|| idempotency_key.map(str::to_string))
.filter(|hint| !hint.is_empty())
.unwrap_or_else(|| "agentic-dispatch".to_string())
}
#[async_trait]
impl AgentDispatchInvoker for RuntimeAgentDispatchInvoker {
async fn invoke(
&self,
tenant: &str,
env: &str,
target: &str,
_operation: &str,
input: Value,
idempotency_key: Option<&str>,
) -> Result<InvokeOutcome> {
if let Some(key) = idempotency_key {
match self.ledger.get(key).await {
Ok(Some(cached)) => {
tracing::debug!(key, "dispatch ledger hit; returning cached output");
return Ok(InvokeOutcome {
ok: true,
output: cached,
events: vec![],
});
}
Ok(None) => {}
Err(e) => {
tracing::warn!(key, error = %e, "dispatch ledger get failed; proceeding without cache");
}
}
}
let user_text = extract_user_text(&input);
let session_id = resolve_session_id(&input, idempotency_key);
let tenant_ctx = TenantContext::new(tenant, env);
let output = self
.runtime
.step(
tenant_ctx,
&session_id,
target,
AgentInput { text: user_text },
)
.await
.with_context(|| format!("agentic step failed for agent '{target}'"))?;
let outcome_output = json!({
"reply": output.reply,
"trail": output.trail,
"terminated_by": output.terminated_by,
});
if let Some(key) = idempotency_key
&& let Err(e) = self.ledger.record(key, outcome_output.clone()).await
{
tracing::warn!(key, error = %e, "dispatch ledger record failed; redelivery will re-run step");
}
Ok(InvokeOutcome {
ok: true,
output: outcome_output,
events: vec![],
})
}
}
#[must_use]
pub fn use_jetstream(get_env: impl Fn(&str) -> Option<String>) -> bool {
match get_env("GREENTIC_AW_JETSTREAM") {
Some(v) => !matches!(
v.trim().to_ascii_lowercase().as_str(),
"0" | "false" | "no" | "off"
),
None => true,
}
}
#[must_use]
pub fn warm_targets(get_env: impl Fn(&str) -> Option<String>) -> Vec<String> {
get_env("GREENTIC_AW_WARM_PACKS")
.map(|raw| {
raw.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.map(str::to_string)
.collect()
})
.unwrap_or_default()
}
pub fn warm_on_start(get_env: impl Fn(&str) -> Option<String>) {
let targets = warm_targets(get_env);
if targets.is_empty() {
tracing::debug!("aw serve: no warm targets (GREENTIC_AW_WARM_PACKS unset)");
} else {
tracing::info!(
count = targets.len(),
?targets,
"aw serve: warm targets configured"
);
}
}
pub async fn serve_with_ledger(
nats_url: &str,
runtime: Arc<AgentRuntime>,
ledger: Arc<dyn DispatchLedger>,
) -> Result<()> {
warm_on_start(|k| std::env::var(k).ok());
let client = async_nats::connect(nats_url)
.await
.with_context(|| format!("connecting to NATS at {nats_url}"))?;
tracing::info!(
nats_url,
subject = aw_event_bridge::request_topic(aw_event_bridge::RUNTIME_NAME),
"aw event bridge connected; serving agentic dispatch"
);
let invoker = Arc::new(RuntimeAgentDispatchInvoker::with_ledger(runtime, ledger));
if use_jetstream(|k| std::env::var(k).ok()) {
tracing::info!(nats_url, "aw serve: JetStream durable consumer");
run_bridge_jetstream(client, invoker).await
} else {
tracing::info!(nats_url, "aw serve: core-NATS consumer (legacy)");
run_bridge(client, invoker).await
}
}
pub async fn serve(nats_url: &str, runtime: Arc<AgentRuntime>) -> Result<()> {
serve_with_ledger(nats_url, runtime, Arc::new(NoopDispatchLedger)).await
}
#[cfg(feature = "test-mock")]
#[must_use]
pub fn build_test_mock_runtime(agent_id: &str, reply: &str) -> Arc<AgentRuntime> {
use crate::cost::MockTokenMeter;
use crate::llm::LlmResponse;
use crate::mock::{
MockAgentStateStore, MockConfigProvider, MockLlmBackend, MockTelemetry, NoopToolLedger,
};
use crate::tools::ToolLedger;
use crate::{AgentConfig, AgentLimits, LlmProviderRef};
let scripted = (0..64)
.map(|_| {
Ok(LlmResponse {
content: Some(reply.to_string()),
tool_calls: vec![],
tokens_in: 1,
tokens_out: 1,
})
})
.collect();
let llm = Arc::new(MockLlmBackend::new(scripted));
let store = Arc::new(MockAgentStateStore::new());
let telemetry = Arc::new(MockTelemetry::new());
let config_provider = MockConfigProvider::new();
let agent_config = AgentConfig {
agent_id: agent_id.to_string(),
system_prompt: "test-mock agent".to_string(),
tools: vec![],
guardrails: vec![],
llm: LlmProviderRef {
provider: "mock".to_string(),
model: "mock".to_string(),
credential_ref: None,
},
limits: AgentLimits::default(),
memory: None,
knowledge: None,
};
for (tenant, env) in [
("default", "default"),
("acme", "prod"),
("t", "e"),
("sorx", "default"),
] {
config_provider.insert(
&TenantContext::new(tenant, env),
agent_id,
agent_config.clone(),
);
}
let config_provider = Arc::new(config_provider);
let token_meter = Arc::new(MockTokenMeter::new(0));
let ledger: Arc<dyn ToolLedger> = Arc::new(NoopToolLedger);
let ext_runtime = Arc::new(crate::test_support::extension_runtime());
Arc::new(AgentRuntime::new(
config_provider,
store,
ext_runtime,
llm,
telemetry,
token_meter,
ledger,
None,
))
}
#[cfg(test)]
mod warm_targets_tests {
use super::warm_targets;
#[test]
fn warm_targets_parses_csv_and_empty() {
assert_eq!(warm_targets(|_| None), Vec::<String>::new());
assert_eq!(
warm_targets(|k| (k == "GREENTIC_AW_WARM_PACKS").then(|| "a, b ,c".to_string())),
vec!["a", "b", "c"]
);
assert_eq!(
warm_targets(|k| (k == "GREENTIC_AW_WARM_PACKS").then(|| "".to_string())),
Vec::<String>::new()
);
}
}
#[cfg(test)]
mod env_gate_tests {
use super::use_jetstream;
#[test]
fn jetstream_default_on_unless_disabled() {
assert!(use_jetstream(|_| None)); assert!(!use_jetstream(
|k| (k == "GREENTIC_AW_JETSTREAM").then(|| "0".to_string())
));
assert!(!use_jetstream(
|k| (k == "GREENTIC_AW_JETSTREAM").then(|| "off".to_string())
));
assert!(use_jetstream(
|k| (k == "GREENTIC_AW_JETSTREAM").then(|| "on".to_string())
));
}
}
#[cfg(all(test, feature = "test-mock"))]
mod tests {
use super::*;
#[test]
fn extract_user_text_accepts_node_and_raw_shapes() {
assert_eq!(extract_user_text(&json!({"user_text": "hi"})), "hi");
assert_eq!(extract_user_text(&json!({"text": "yo"})), "yo");
assert_eq!(extract_user_text(&json!("bare")), "bare");
assert_eq!(extract_user_text(&json!({"other": 1})), "");
}
#[test]
fn resolve_session_id_prefers_explicit_then_idempotency() {
assert_eq!(
resolve_session_id(&json!({"session_id": "s1"}), Some("corr")),
"s1"
);
assert_eq!(resolve_session_id(&json!({}), Some("corr")), "corr");
assert_eq!(resolve_session_id(&json!({}), Some("")), "agentic-dispatch");
assert_eq!(resolve_session_id(&json!({}), None), "agentic-dispatch");
}
#[tokio::test]
#[allow(clippy::expect_used)] async fn mock_invoker_returns_reply_output() {
let runtime = build_test_mock_runtime("greeter", "pong");
let invoker = RuntimeAgentDispatchInvoker::new(runtime);
let outcome = invoker
.invoke(
"acme",
"prod",
"greeter",
"",
json!({"user_text": "ping"}),
Some("sess-1::pack=p::flow=f"),
)
.await
.expect("mock invoke succeeds");
assert!(outcome.ok);
assert_eq!(outcome.output["reply"], json!("pong"));
assert_eq!(outcome.output["terminated_by"], json!("final_reply"));
}
#[tokio::test]
#[allow(clippy::expect_used)]
async fn redelivery_returns_cached_without_rerunning_step() {
use crate::dispatch_ledger::InMemoryDispatchLedger;
let ledger = Arc::new(InMemoryDispatchLedger::with(
"k1",
json!({"reply": "CACHED", "trail": [], "terminated_by": "final_reply"}),
));
let invoker = RuntimeAgentDispatchInvoker::with_ledger(
build_test_mock_runtime("greeter", "pong"),
ledger,
);
let out = invoker
.invoke(
"acme",
"prod",
"greeter",
"",
json!({"user_text": "hi"}),
Some("k1"),
)
.await
.expect("cached invoke succeeds");
assert!(out.ok);
assert_eq!(
out.output["reply"],
json!("CACHED"),
"expected cached sentinel, got runtime reply"
);
}
#[tokio::test]
#[allow(clippy::expect_used)]
async fn miss_runs_step_and_records() {
use crate::dispatch_ledger::InMemoryDispatchLedger;
let ledger = Arc::new(InMemoryDispatchLedger::default());
let invoker = RuntimeAgentDispatchInvoker::with_ledger(
build_test_mock_runtime("greeter", "pong"),
ledger.clone(),
);
let out = invoker
.invoke(
"acme",
"prod",
"greeter",
"",
json!({"user_text": "hi"}),
Some("k2"),
)
.await
.expect("fresh invoke succeeds");
assert!(out.ok);
assert_eq!(out.output["reply"], json!("pong"), "step ran and replied");
let stored = ledger.stored("k2");
assert!(stored.is_some(), "result was recorded in the ledger");
assert_eq!(
stored.expect("stored entry present")["reply"],
json!("pong"),
"recorded value matches step output"
);
}
#[tokio::test]
#[allow(clippy::expect_used)]
async fn no_idempotency_key_runs_step_without_ledger() {
use crate::dispatch_ledger::InMemoryDispatchLedger;
let ledger = Arc::new(InMemoryDispatchLedger::default());
let invoker = RuntimeAgentDispatchInvoker::with_ledger(
build_test_mock_runtime("greeter", "pong"),
ledger.clone(),
);
let out = invoker
.invoke(
"acme",
"prod",
"greeter",
"",
json!({"user_text": "hi"}),
None,
)
.await
.expect("no-key invoke succeeds");
assert!(out.ok);
assert_eq!(out.output["reply"], json!("pong"));
assert!(ledger.stored("k-absent").is_none());
}
}