use std::sync::Arc;
use parking_lot::Mutex;
use theway_core::agent::session::session::Session;
use theway_core::types::{
AfterToolCallHook, AgentMessage, BeforeToolCallHook, StreamFn, ThinkingLevel,
};
use theway_core::{
Agent, AgentOptions, AgentRunError, AgentState, AgentTool, LoopEvent, SessionError,
};
use theway_llm_provider::{Message as PiMessage, Model};
use crate::trigger_engine::event::{TriggerEvent, TriggerListener};
use crate::trigger_engine::runtime::TriggerRuntimeSnapshot;
use crate::trigger_engine::types::Trigger;
use super::RunningTriggerHandle;
use super::promotion::{
PROMOTION_BODY_CAP_BYTES, apply_promotion, compute_sub_agent_outcome, ensure_trigger_prefix,
truncate_on_char_boundary,
};
use super::types::{
BeforeTriggerActionContext, BeforeTriggerActionHook, RunningTriggerState, TriggerAction,
TriggerDelivery,
};
use super::utils::{emit_from_listeners, preview_for_banner};
#[allow(clippy::too_many_arguments)]
pub(super) async fn run_trigger_action(
trigger: Trigger,
trace_id: String,
source_label: String,
event_label: String,
listeners: Arc<Mutex<Vec<TriggerListener>>>,
parent_session: Session,
parent_agent: Arc<Agent>,
running_registry: Arc<Mutex<std::collections::HashMap<String, RunningTriggerHandle>>>,
action_hook: Option<BeforeTriggerActionHook>,
runtime_snapshot: TriggerRuntimeSnapshot,
parent_model: Option<Model>,
parent_system_prompt: String,
parent_tools: Vec<Arc<dyn AgentTool>>,
parent_thinking: Option<ThinkingLevel>,
stream_fn: Option<StreamFn>,
before_tool_call: Option<BeforeToolCallHook>,
after_tool_call: Option<AfterToolCallHook>,
) {
let cancel = tokio_util::sync::CancellationToken::new();
let action = match action_hook {
Some(hook) => {
let ctx = BeforeTriggerActionContext {
trigger: trigger.clone(),
runtime: runtime_snapshot,
};
hook(ctx, cancel.clone()).await
}
None => TriggerAction::default_for(&trigger),
};
if action.delivery == TriggerDelivery::InjectSummary {
let summary = trigger.payload_summary.clone();
emit_from_listeners(
&listeners,
TriggerEvent::TriggerExecutionStarted {
trace_id: trace_id.clone(),
source_label: source_label.clone(),
event_label: event_label.clone(),
prompt_preview: preview_for_banner(
summary.as_deref().unwrap_or("(no summary)"),
80,
),
},
);
let result_data = serde_json::json!({
"trace_id": trace_id,
"branch_id": serde_json::Value::Null,
"success": true,
"summary": summary,
"message_count": 0,
"cost_usd": 0.0,
"reason": serde_json::Value::Null,
"details": serde_json::Value::Null,
"delivery": "inject_summary",
});
if let Err(e) = parent_session
.append_custom("trigger_result", Some(result_data))
.await
{
emit_from_listeners(
&listeners,
TriggerEvent::PersistenceError {
context: "trigger_result".into(),
message: format!("trigger_result (inject) append failed: {:?}", e.code),
},
);
}
emit_from_listeners(
&listeners,
TriggerEvent::TriggerCompleted {
trace_id: trace_id.clone(),
summary: summary.clone(),
cost_usd: Some(0.0),
details: serde_json::Value::Null,
},
);
apply_promotion(
&listeners,
&parent_session,
&parent_agent,
&trace_id,
&trigger,
true,
&summary,
0,
None,
&action.promote,
action.promote_requires_approval,
&serde_json::Value::Null,
)
.await;
return;
}
if action.delivery == TriggerDelivery::InjectAndRun {
let (body, _truncated) =
truncate_on_char_boundary(action.prompt.clone(), PROMOTION_BODY_CAP_BYTES);
let (body, prefix_injected) = ensure_trigger_prefix(body, &trace_id);
emit_from_listeners(
&listeners,
TriggerEvent::TriggerExecutionStarted {
trace_id: trace_id.clone(),
source_label: source_label.clone(),
event_label: event_label.clone(),
prompt_preview: preview_for_banner(&body, 80),
},
);
let user_message = AgentMessage::Llm(PiMessage::User(theway_llm_provider::UserMessage {
role: theway_llm_provider::UserRole::User,
content: theway_llm_provider::UserContent::Text(body.clone()),
timestamp: chrono::Utc::now().timestamp_millis(),
}));
let queued_for_followup = parent_agent.is_streaming();
if queued_for_followup {
parent_agent.enqueue_follow_up(user_message);
} else if let Err(e) = parent_session.append_message(user_message.clone()).await {
emit_from_listeners(
&listeners,
TriggerEvent::PersistenceError {
context: "trigger_inject_and_run".into(),
message: format!("inject_and_run append failed: {:?}", e.code),
},
);
} else {
parent_agent.state().messages.push(user_message);
}
let result_data = serde_json::json!({
"trace_id": trace_id,
"branch_id": serde_json::Value::Null,
"success": true,
"summary": body,
"message_count": 0,
"cost_usd": 0.0,
"reason": serde_json::Value::Null,
"details": serde_json::Value::Null,
"delivery": "inject_and_run",
"prefix_injected": prefix_injected,
"run_dispatch": if queued_for_followup { "follow_up" } else { "main_run_request" },
});
if let Err(e) = parent_session
.append_custom("trigger_result", Some(result_data))
.await
{
emit_from_listeners(
&listeners,
TriggerEvent::PersistenceError {
context: "trigger_result".into(),
message: format!(
"trigger_result (inject_and_run) append failed: {:?}",
e.code
),
},
);
}
emit_from_listeners(
&listeners,
TriggerEvent::TriggerCompleted {
trace_id: trace_id.clone(),
summary: Some(body),
cost_usd: Some(0.0),
details: serde_json::Value::Null,
},
);
if !queued_for_followup {
emit_from_listeners(
&listeners,
TriggerEvent::TriggerRequestsMainRun {
trace_id: trace_id.clone(),
},
);
}
return;
}
let prompt_preview = preview_for_banner(&action.prompt, 80);
let started_at = chrono::Utc::now();
{
let mut reg = running_registry.lock();
reg.insert(
trace_id.clone(),
RunningTriggerHandle {
state: RunningTriggerState {
trace_id: trace_id.clone(),
source_label: source_label.clone(),
event_label: event_label.clone(),
started_at,
prompt_preview: prompt_preview.clone(),
},
cancel: cancel.clone(),
},
);
}
emit_from_listeners(
&listeners,
TriggerEvent::TriggerExecutionStarted {
trace_id: trace_id.clone(),
source_label: source_label.clone(),
event_label: event_label.clone(),
prompt_preview,
},
);
let sub_storage: Arc<dyn theway_core::agent::session::session::SessionStorage> =
Arc::new(theway_core::agent::session::memory_storage::MemorySessionStorage::new());
let sub_session = theway_core::agent::session::session::Session::new(sub_storage);
let mut sub_state = AgentState::default();
sub_state.model = parent_model;
sub_state.thinking_level = parent_thinking;
sub_state.tools = parent_tools;
sub_state.system_prompt = parent_system_prompt;
let sub_agent = Agent::new(AgentOptions {
initial_state: Some(sub_state),
stream_fn,
before_tool_call,
after_tool_call,
observer: parent_agent.runtime_observer(),
observation_context: parent_agent.observation_context(),
observation_parent: parent_agent.active_run_operation(),
..Default::default()
});
let persist_errors: Arc<Mutex<Vec<SessionError>>> = Arc::new(Mutex::new(Vec::new()));
let persist_session = sub_session.clone();
let persist_errors_listener = persist_errors.clone();
let _persist_unsub = sub_agent.subscribe(Arc::new(move |event, _cancel| {
let session = persist_session.clone();
let sink = persist_errors_listener.clone();
Box::pin(async move {
if let LoopEvent::MessageEnd { message } = event {
if let Err(e) = session.append_message(message).await {
sink.lock().push(e);
}
}
})
}));
let user_message = AgentMessage::Llm(PiMessage::User(theway_llm_provider::UserMessage {
role: theway_llm_provider::UserRole::User,
content: theway_llm_provider::UserContent::Text(action.prompt.clone()),
timestamp: chrono::Utc::now().timestamp_millis(),
}));
let run_outcome: Result<(), AgentRunError> = tokio::select! {
biased;
_ = cancel.cancelled() => {
sub_agent.abort();
Err(AgentRunError::Other("aborted".into()))
}
res = sub_agent.prompt(user_message) => res,
};
let (success, summary, message_count) = compute_sub_agent_outcome(&sub_agent, &run_outcome);
let failure_reason: Option<String> = if success {
None
} else {
Some(match &run_outcome {
Err(AgentRunError::Other(msg)) if msg == "aborted" => "aborted".to_string(),
Err(e) => format!("{e}"),
Ok(_) => "unknown failure".to_string(),
})
};
let details_for_promotion: serde_json::Value = serde_json::Value::Null;
let result_data = serde_json::json!({
"trace_id": trace_id,
"branch_id": serde_json::Value::Null,
"success": success,
"summary": summary,
"message_count": message_count,
"cost_usd": serde_json::Value::Null,
"reason": failure_reason,
"details": details_for_promotion,
});
let audit_write_result = parent_session
.append_custom("trigger_result", Some(result_data))
.await;
if let Err(e) = audit_write_result {
emit_from_listeners(
&listeners,
TriggerEvent::PersistenceError {
context: "trigger_result".into(),
message: format!("trigger_result append failed: {:?}", e.code),
},
);
}
for e in persist_errors.lock().iter() {
emit_from_listeners(
&listeners,
TriggerEvent::PersistenceError {
context: "trigger_result".into(),
message: format!("sub-agent session append failed: {:?}", e.code),
},
);
}
if success {
emit_from_listeners(
&listeners,
TriggerEvent::TriggerCompleted {
trace_id: trace_id.clone(),
summary: summary.clone(),
cost_usd: None,
details: details_for_promotion.clone(),
},
);
} else {
emit_from_listeners(
&listeners,
TriggerEvent::TriggerFailed {
trace_id: trace_id.clone(),
reason: failure_reason
.clone()
.unwrap_or_else(|| "unknown failure".to_string()),
},
);
}
apply_promotion(
&listeners,
&parent_session,
&parent_agent,
&trace_id,
&trigger,
success,
&summary,
message_count,
failure_reason.as_deref(),
&action.promote,
action.promote_requires_approval,
&details_for_promotion,
)
.await;
running_registry.lock().remove(&trace_id);
}
#[cfg(test)]
tests_bridge_macro::tests_bridge!("trigger_engine/execution/action");