diff --git a/AGENTS.md b/AGENTS.md index fe5cfe033..42a7ab8f9 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -79,7 +79,7 @@ Fabro is an AI-powered workflow orchestration platform. Workflows are defined as When working on Rust crates, read the relevant strategy doc **before** making changes: - **`docs-internal/logging-strategy.md`** — read when adding `tracing` calls (`info!`, `debug!`, `warn!`, `error!`), working on error handling paths, or adding new operations that should be observable -- **`docs-internal/events-strategy.md`** — read when adding or modifying `Event` variants, touching `EventEmitter`/`emit()`, changing `progress.jsonl` output, or adding new workflow stage types +- **`docs-internal/events-strategy.md`** — read when adding or modifying `Event` variants, touching `Emitter`/`emit()`, changing `progress.jsonl` output, or adding new workflow stage types - **`files-internal/testing-strategy.md`** — read when adding or reorganizing tests, choosing between unit vs `tests/it`, deciding whether a test belongs in `cmd` vs `workflow` vs `scenario`, or deciding how to structure snapshots and fixtures ## Shell quoting in sandbox code diff --git a/docs-internal/events-strategy.md b/docs-internal/events-strategy.md index ee69a070f..474fe46f3 100644 --- a/docs-internal/events-strategy.md +++ b/docs-internal/events-strategy.md @@ -9,7 +9,7 @@ Detached runs rely on this distinction. If something needs to be visible after r ## Architecture ```text -Engine/Handler -> Event -> EventEmitter::emit() +Engine/Handler -> Event -> Emitter::emit() |- trace(raw event) |- canonicalize -> RunEventEnvelope `- on_event(&RunEventEnvelope) @@ -22,7 +22,7 @@ Engine/Handler -> Event -> EventEmitter::emit() The canonical envelope is built exactly once in `fabro-workflow/src/event.rs`. - `Event` remains the internal typed source of truth. -- `EventEmitter` owns an immutable `run_id` and converts typed events into `RunEventEnvelope`. +- `Emitter` owns an immutable `run_id` and converts typed events into `RunEventEnvelope`. - Every listener receives `&RunEventEnvelope`, not `&Event`. - Bypass paths that cannot go through the emitter must call `canonicalize_event()` once and reuse the same envelope for every sink. @@ -100,7 +100,7 @@ Agent events now use explicit session links: ## Direct-Write Paths -Most events flow through `EventEmitter::emit()`. The remaining direct-write paths must use: +Most events flow through `Emitter::emit()`. The remaining direct-write paths must use: 1. `canonicalize_event(run_id, event)` 2. Serialize and redact once @@ -132,7 +132,7 @@ Update `extract_envelope_fields()`: ### 5. Emit it -Prefer `EventEmitter::emit(&Event::...)`. +Prefer `Emitter::emit(&Event::...)`. Use `canonicalize_event()` only for true bypass paths. diff --git a/lib/crates/fabro-agent/README.md b/lib/crates/fabro-agent/README.md index 7b1442e97..9003b5d09 100644 --- a/lib/crates/fabro-agent/README.md +++ b/lib/crates/fabro-agent/README.md @@ -44,7 +44,7 @@ User Input - **`Sandbox`** (trait) -- Abstracts filesystem, shell, grep, and glob operations. `LocalSandbox` provides a real implementation; the trait enables sandboxing and testing. - **`ToolRegistry`** -- Maps tool names to definitions and async executor functions. Tools are registered per-profile. - **`History`** -- Ordered list of `Turn` variants (`User`, `Assistant`, `ToolResults`, `System`, `Steering`) that converts to LLM messages. -- **`EventEmitter`** -- Broadcasts `SessionEvent`s (tool calls, text, errors, warnings) over a `tokio::sync::broadcast` channel for UI or logging. +- **`Emitter`** -- Broadcasts `SessionEvent`s (tool calls, text, errors, warnings) over a `tokio::sync::broadcast` channel for UI or logging. - **`SubAgentManager`** -- Spawns child `Session`s on background tasks for delegated work, with depth limits. - **`SessionConfig`** -- Tunable parameters: max turns, tool round limits, command timeouts, loop detection, output truncation limits, and user instructions. diff --git a/lib/crates/fabro-agent/src/compaction.rs b/lib/crates/fabro-agent/src/compaction.rs index 14d630aaa..729441a50 100644 --- a/lib/crates/fabro-agent/src/compaction.rs +++ b/lib/crates/fabro-agent/src/compaction.rs @@ -2,7 +2,7 @@ use std::fmt::Write; use crate::agent_profile::AgentProfile; use crate::error::AgentError; -use crate::event::EventEmitter; +use crate::event::Emitter; use crate::file_tracker::FileTracker; use crate::history::History; use crate::truncation; @@ -19,7 +19,7 @@ pub fn check_context_usage( history: &History, provider_profile: &dyn AgentProfile, threshold_percent: usize, - emitter: &EventEmitter, + emitter: &Emitter, session_id: &str, ) -> bool { let estimated_tokens = estimate_token_count(system_prompt, history); @@ -57,7 +57,7 @@ pub async fn compact_context( system_prompt: &str, file_tracker: &FileTracker, preserve_count: usize, - emitter: &EventEmitter, + emitter: &Emitter, session_id: &str, ) -> Result<(), AgentError> { let estimated_tokens = estimate_token_count(system_prompt, history); @@ -251,7 +251,7 @@ pub fn render_turns_for_summary(turns: &[Turn]) -> String { #[cfg(test)] mod tests { use super::*; - use crate::event::EventEmitter; + use crate::event::Emitter; use crate::history::History; use crate::test_support::TestProfile; use crate::tool_registry::ToolRegistry; @@ -331,7 +331,7 @@ mod tests { #[test] fn check_context_usage_below_threshold() { let history = History::default(); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let profile = TestProfile::new(); // Empty history, huge context window => well below threshold let over = check_context_usage("short", &history, &profile, 80, &emitter, "sess"); @@ -346,7 +346,7 @@ mod tests { content: "x".repeat(1000), timestamp: SystemTime::now(), }); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let mut rx = emitter.subscribe(); // TestProfile has context_window=200_000 by default; use a small one let profile = TestProfile::with_context_window(ToolRegistry::new(), 100); diff --git a/lib/crates/fabro-agent/src/event.rs b/lib/crates/fabro-agent/src/event.rs index 1fc120932..2e2cb9d27 100644 --- a/lib/crates/fabro-agent/src/event.rs +++ b/lib/crates/fabro-agent/src/event.rs @@ -3,11 +3,11 @@ use std::time::SystemTime; use tokio::sync::broadcast; #[derive(Clone)] -pub struct EventEmitter { +pub struct Emitter { sender: broadcast::Sender, } -impl EventEmitter { +impl Emitter { #[must_use] pub fn new() -> Self { let (sender, _) = broadcast::channel(1024); @@ -36,7 +36,7 @@ impl EventEmitter { } } -impl Default for EventEmitter { +impl Default for Emitter { fn default() -> Self { Self::new() } @@ -49,7 +49,7 @@ mod tests { #[tokio::test] async fn emit_and_receive_event() { - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let mut receiver = emitter.subscribe(); emitter.emit( @@ -74,7 +74,7 @@ mod tests { #[tokio::test] async fn emit_with_data() { - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let mut receiver = emitter.subscribe(); emitter.emit( @@ -93,7 +93,7 @@ mod tests { #[tokio::test] async fn multiple_subscribers() { - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let mut rx1 = emitter.subscribe(); let mut rx2 = emitter.subscribe(); @@ -111,7 +111,7 @@ mod tests { #[test] fn emit_without_subscribers_does_not_panic() { - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); emitter.emit( "sess-4".into(), AgentEvent::Error { @@ -122,13 +122,13 @@ mod tests { #[test] fn default_creates_emitter() { - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let _rx = emitter.subscribe(); } #[tokio::test] async fn forward_preserves_session_ids() { - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let mut receiver = emitter.subscribe(); emitter.forward(SessionEvent { diff --git a/lib/crates/fabro-agent/src/lib.rs b/lib/crates/fabro-agent/src/lib.rs index ac6ebede1..aa3d5465d 100644 --- a/lib/crates/fabro-agent/src/lib.rs +++ b/lib/crates/fabro-agent/src/lib.rs @@ -31,7 +31,7 @@ pub use config::{SessionOptions, ToolApprovalAdapter, ToolHookCallback, ToolHook #[cfg(feature = "docker")] pub use docker_sandbox::{DockerSandbox, DockerSandboxOptions}; pub use error::{AbortReason, AgentError}; -pub use event::EventEmitter; +pub use event::Emitter; pub use fabro_mcp::config::McpServerSettings; pub use history::History; pub use local_sandbox::LocalSandbox; diff --git a/lib/crates/fabro-agent/src/session.rs b/lib/crates/fabro-agent/src/session.rs index e7e74d73b..440485b87 100644 --- a/lib/crates/fabro-agent/src/session.rs +++ b/lib/crates/fabro-agent/src/session.rs @@ -2,7 +2,7 @@ use crate::agent_profile::AgentProfile; use crate::compaction::{check_context_usage, compact_context}; use crate::config::SessionOptions; use crate::error::{AbortReason, AgentError}; -use crate::event::EventEmitter; +use crate::event::Emitter; use crate::file_tracker::FileTracker; use crate::history::History; use crate::loop_detection::detect_loop; @@ -39,7 +39,7 @@ pub struct Session { id: String, config: SessionOptions, history: History, - event_emitter: EventEmitter, + event_emitter: Emitter, state: SessionState, llm_client: Client, provider_profile: Arc, @@ -70,7 +70,7 @@ impl Session { id: uuid::Uuid::new_v4().to_string(), config, history: History::default(), - event_emitter: EventEmitter::new(), + event_emitter: Emitter::new(), state: SessionState::Idle, llm_client, provider_profile, diff --git a/lib/crates/fabro-agent/src/tool_execution.rs b/lib/crates/fabro-agent/src/tool_execution.rs index 6bf3c41b8..291d89e59 100644 --- a/lib/crates/fabro-agent/src/tool_execution.rs +++ b/lib/crates/fabro-agent/src/tool_execution.rs @@ -1,5 +1,5 @@ use crate::config::{SessionOptions, ToolHookCallback, ToolHookDecision}; -use crate::event::EventEmitter; +use crate::event::Emitter; use crate::sandbox::Sandbox; use crate::tool_registry::{RegisteredTool, ToolContext, ToolRegistry}; use crate::truncation::truncate_tool_output; @@ -21,7 +21,7 @@ pub async fn execute_tool_calls( tool_hooks: Option<&Arc>, cancel_token: &CancellationToken, config: &SessionOptions, - emitter: &EventEmitter, + emitter: &Emitter, session_id: &str, tool_env: Option<&HashMap>, ) -> Vec { @@ -62,7 +62,7 @@ async fn execute_tool_calls_sequential( tool_hooks: Option<&Arc>, cancel_token: &CancellationToken, config: &SessionOptions, - emitter: &EventEmitter, + emitter: &Emitter, session_id: &str, tool_env: Option<&HashMap>, ) -> Vec { @@ -98,7 +98,7 @@ async fn execute_tool_calls_parallel( tool_hooks: Option<&Arc>, cancel_token: &CancellationToken, config: &SessionOptions, - emitter: &EventEmitter, + emitter: &Emitter, session_id: &str, tool_env: Option<&HashMap>, ) -> Vec { @@ -145,7 +145,7 @@ pub async fn execute_and_emit_one_tool( tool_hooks: Option<&Arc>, cancel_token: CancellationToken, config: &SessionOptions, - emitter: &EventEmitter, + emitter: &Emitter, session_id: &str, tool_env: Option<&HashMap>, ) -> ToolResult { @@ -172,7 +172,7 @@ async fn execute_and_emit_one_tool_with_lookup( tool_hooks: Option<&Arc>, cancel_token: CancellationToken, config: &SessionOptions, - emitter: &EventEmitter, + emitter: &Emitter, session_id: &str, tool_env: Option<&HashMap>, ) -> ToolResult { @@ -345,7 +345,7 @@ pub fn validate_tool_args( mod tests { use super::*; use crate::config::{ToolHookCallback, ToolHookDecision}; - use crate::event::EventEmitter; + use crate::event::Emitter; use crate::local_sandbox::LocalSandbox; use crate::read_before_write_sandbox::ReadBeforeWriteSandbox; use crate::test_support::MutableMockSandbox; @@ -460,7 +460,7 @@ mod tests { })); let tc = make_tool_call("echo", "call_1", serde_json::json!({"text": "hello"})); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); let result = execute_and_emit_one_tool( @@ -490,7 +490,7 @@ mod tests { Arc::new(MockHookCallback::new(ToolHookDecision::Proceed)); let tc = make_tool_call("echo", "call_1", serde_json::json!({"text": "hello"})); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); let result = execute_and_emit_one_tool( @@ -520,7 +520,7 @@ mod tests { let hooks: Arc = mock.clone(); let tc = make_tool_call("echo", "call_1", serde_json::json!({"text": "hello"})); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); execute_and_emit_one_tool( @@ -555,7 +555,7 @@ mod tests { let hooks: Arc = mock.clone(); let tc = make_tool_call("fail_tool", "call_1", serde_json::json!({})); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); execute_and_emit_one_tool( @@ -587,7 +587,7 @@ mod tests { registry.register(make_echo_tool()); let tc = make_tool_call("echo", "call_1", serde_json::json!({"text": "hello"})); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); let result = execute_and_emit_one_tool( @@ -627,7 +627,7 @@ mod tests { "call_1", serde_json::json!({"file_path": "a.ts", "content": "new"}), ); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); let result = execute_and_emit_one_tool( @@ -654,7 +654,7 @@ mod tests { registry.register(make_write_file_tool()); let sandbox = make_guarded_sandbox(HashMap::from([("a.ts".into(), "content".into())])); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); // First read the file @@ -706,7 +706,7 @@ mod tests { registry.register(make_write_file_tool()); let sandbox = make_guarded_sandbox(HashMap::from([("a.ts".into(), "content".into())])); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); // Grep matching a.ts @@ -758,7 +758,7 @@ mod tests { "call_1", serde_json::json!({"file_path": "a.ts", "old_string": "content", "new_string": "updated"}), ); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); let result = execute_and_emit_one_tool( @@ -789,7 +789,7 @@ mod tests { "call_1", serde_json::json!({"file_path": "new.ts", "content": "hello"}), ); - let emitter = EventEmitter::new(); + let emitter = Emitter::new(); let config = SessionOptions::default(); let result = execute_and_emit_one_tool( diff --git a/lib/crates/fabro-cli/src/commands/run/detached.rs b/lib/crates/fabro-cli/src/commands/run/detached.rs index c453b2acc..bf1d4b98d 100644 --- a/lib/crates/fabro-cli/src/commands/run/detached.rs +++ b/lib/crates/fabro-cli/src/commands/run/detached.rs @@ -5,7 +5,7 @@ use anyhow::{Result, anyhow}; use fabro_interview::FileInterviewer; use fabro_store::RuntimeState; use fabro_types::RunId; -use fabro_workflow::event::EventEmitter; +use fabro_workflow::event::Emitter; use fabro_workflow::operations::{StartServices, resume as resume_run, start as start_run}; use crate::shared; @@ -50,7 +50,7 @@ pub(crate) async fn execute( let services = StartServices { run_id: run_record.run_id, cancel_token: None, - emitter: Arc::new(EventEmitter::new(run_record.run_id)), + emitter: Arc::new(Emitter::new(run_record.run_id)), interviewer: Arc::new(FileInterviewer::new( runtime_state.interview_request_path(), runtime_state.interview_response_path(), diff --git a/lib/crates/fabro-cli/tests/it/workflow/real_cli.rs b/lib/crates/fabro-cli/tests/it/workflow/real_cli.rs index f00d9a377..dcc906bb4 100644 --- a/lib/crates/fabro-cli/tests/it/workflow/real_cli.rs +++ b/lib/crates/fabro-cli/tests/it/workflow/real_cli.rs @@ -4,7 +4,7 @@ use std::time::Duration; use fabro_graphviz::graph::{AttrValue, Node}; use fabro_llm::provider::Provider; use fabro_workflow::context::Context; -use fabro_workflow::event::EventEmitter; +use fabro_workflow::event::Emitter; use fabro_workflow::handler::agent::{CodergenBackend, CodergenResult}; use fabro_workflow::handler::llm::cli::AgentCliBackend; @@ -24,7 +24,7 @@ async fn run_real_cli_test(provider: Provider, model: &str) { ); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = backend .run( &node, diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 640948c71..632958071 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -47,7 +47,7 @@ use crate::static_files; use crate::web_auth; use fabro_interview::{Answer, Interviewer, QuestionType, WebInterviewer}; use fabro_workflow::context::Context; -use fabro_workflow::event::EventEmitter; +use fabro_workflow::event::Emitter; use fabro_workflow::operations::{self, CreateRunInput, WorkflowInput}; use fabro_workflow::pipeline::Persisted; use fabro_workflow::records::Checkpoint; @@ -650,7 +650,7 @@ async fn execute_run(state: Arc, run_id: RunId) { // Create interviewer and event plumbing (this is the "provisioning" phase) let interviewer = Arc::new(WebInterviewer::new()); let context = Context::new(); - let emitter = EventEmitter::new(run_id); + let emitter = Emitter::new(run_id); if let Some(tx_clone) = event_tx { emitter.on_event(move |event| { let _ = tx_clone.send(event.clone()); diff --git a/lib/crates/fabro-workflow/src/devcontainer_bridge.rs b/lib/crates/fabro-workflow/src/devcontainer_bridge.rs index b7558da27..6ed9ee497 100644 --- a/lib/crates/fabro-workflow/src/devcontainer_bridge.rs +++ b/lib/crates/fabro-workflow/src/devcontainer_bridge.rs @@ -4,7 +4,7 @@ use sha2::{Digest, Sha256}; use fabro_devcontainer::DevcontainerSpec; -use crate::event::{Event, EventEmitter}; +use crate::event::{Emitter, Event}; use fabro_agent::sandbox::Sandbox; use fabro_sandbox::daytona::{DaytonaSnapshotConfig, DockerfileSource}; use futures::future::try_join_all; @@ -32,7 +32,7 @@ pub fn devcontainer_to_snapshot_config(dc: &DevcontainerSpec) -> DaytonaSnapshot /// Follows the same pattern as setup commands in `run.rs`. pub async fn run_devcontainer_lifecycle( sandbox: &dyn Sandbox, - emitter: &EventEmitter, + emitter: &Emitter, phase: &str, commands: &[fabro_devcontainer::Command], timeout_ms: u64, @@ -139,7 +139,7 @@ pub async fn run_devcontainer_lifecycle( async fn run_single_lifecycle_command( sandbox: &dyn Sandbox, - emitter: &EventEmitter, + emitter: &Emitter, phase: &str, command: &str, index: usize, @@ -348,7 +348,7 @@ mod tests { #[tokio::test] async fn shell_command_executed() { let sandbox = TestSandbox::new(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let commands = vec![fabro_devcontainer::Command::Shell("echo hi".to_string())]; run_devcontainer_lifecycle(&sandbox, &emitter, "on_create", &commands, 300_000) .await @@ -361,7 +361,7 @@ mod tests { #[tokio::test] async fn args_command_joins() { let sandbox = TestSandbox::new(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let commands = vec![fabro_devcontainer::Command::Args(vec![ "echo".to_string(), "hi".to_string(), @@ -380,7 +380,7 @@ mod tests { #[tokio::test] async fn emits_started_and_completed_events() { - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = Arc::new(Mutex::new(Vec::::new())); let events_clone = Arc::clone(&events); emitter.on_event(move |event| { @@ -417,7 +417,7 @@ mod tests { #[tokio::test] async fn failed_command_emits_failed_and_returns_error() { - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = Arc::new(Mutex::new(Vec::::new())); let events_clone = Arc::clone(&events); emitter.on_event(move |event| { @@ -438,7 +438,7 @@ mod tests { #[tokio::test] async fn empty_commands_is_noop() { - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = Arc::new(Mutex::new(Vec::new())); let events_clone = Arc::clone(&events); emitter.on_event(move |event| { @@ -454,7 +454,7 @@ mod tests { #[tokio::test] async fn parallel_commands_run() { let sandbox = TestSandbox::new(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let mut map = HashMap::new(); map.insert("install".to_string(), "npm install".to_string()); map.insert("build".to_string(), "npm run build".to_string()); diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index fb0e73dd2..bbb253670 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -1489,7 +1489,7 @@ impl StoreProgressLogger { Self { tx } } - pub fn register(&self, emitter: &EventEmitter) { + pub fn register(&self, emitter: &Emitter) { let tx = self.tx.clone(); emitter.on_event( move |event| match build_redacted_event_payload(event, &event.run_id) { @@ -1532,17 +1532,17 @@ fn epoch_millis() -> i64 { type EventListener = Arc; /// Callback-based event emitter for workflow run events. -pub struct EventEmitter { +pub struct Emitter { run_id: RunId, listeners: std::sync::Mutex>, /// Epoch milliseconds of the last `emit()` or `touch()` call. 0 until first event. last_event_at: AtomicI64, } -impl std::fmt::Debug for EventEmitter { +impl std::fmt::Debug for Emitter { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { let count = self.listeners.lock().map(|l| l.len()).unwrap_or(0); - f.debug_struct("EventEmitter") + f.debug_struct("Emitter") .field("run_id", &self.run_id) .field("listener_count", &count) .field("last_event_at", &self.last_event_at.load(Ordering::Relaxed)) @@ -1550,13 +1550,13 @@ impl std::fmt::Debug for EventEmitter { } } -impl Default for EventEmitter { +impl Default for Emitter { fn default() -> Self { Self::new(RunId::new()) } } -impl EventEmitter { +impl Emitter { #[must_use] pub fn new(run_id: RunId) -> Self { Self { @@ -1642,13 +1642,13 @@ mod tests { #[test] fn event_emitter_new_has_no_listeners() { - let emitter = EventEmitter::new(fixtures::RUN_1); + let emitter = Emitter::new(fixtures::RUN_1); assert_eq!(emitter.listeners.lock().unwrap().len(), 0); } #[test] fn event_emitter_calls_listener_with_envelope() { - let emitter = EventEmitter::new(fixtures::RUN_1); + let emitter = Emitter::new(fixtures::RUN_1); let received = Arc::new(Mutex::new(Vec::new())); let received_clone = Arc::clone(&received); emitter.on_event(move |event| { @@ -1672,7 +1672,7 @@ mod tests { #[test] fn event_emitter_default() { - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); assert_eq!(emitter.listeners.lock().unwrap().len(), 0); } diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs index 03f601db4..4a01b78ed 100644 --- a/lib/crates/fabro-workflow/src/handler/agent.rs +++ b/lib/crates/fabro-workflow/src/handler/agent.rs @@ -10,7 +10,7 @@ use fabro_types::RunId; use crate::context::keys; use crate::context::{Context, WorkflowContext}; use crate::error::FabroError; -use crate::event::{Event, EventEmitter}; +use crate::event::{Emitter, Event}; use crate::outcome::{ FailureCategory, FailureDetail, Outcome, OutcomeExt, StageStatus, StageUsage, }; @@ -42,7 +42,7 @@ pub trait CodergenBackend: Send + Sync { prompt: &str, context: &Context, thread_id: Option<&str>, - emitter: &Arc, + emitter: &Arc, sandbox: &Arc, tool_hooks: Option>, ) -> Result; @@ -393,7 +393,7 @@ impl Handler for AgentHandler { #[cfg(test)] mod tests { use super::*; - use crate::event::EventEmitter; + use crate::event::Emitter; use fabro_graphviz::graph::AttrValue; use fabro_store::{SlateRunStore, SlateStore, StageId}; use fabro_types::fixtures; @@ -422,7 +422,7 @@ mod tests { let store = test_store(); let run_store = store.create_run(&fixtures::RUN_1).await.unwrap(); let services = EngineServices { - emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)), + emitter: Arc::new(crate::event::Emitter::new(fixtures::RUN_1)), run_store: run_store.clone(), ..EngineServices::test_default() }; @@ -591,7 +591,7 @@ mod tests { _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -647,7 +647,7 @@ mod tests { _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -705,7 +705,7 @@ mod tests { _prompt: &str, context: &Context, _thread_id: Option<&str>, - emitter: &Arc, + emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -812,7 +812,7 @@ mod tests { _prompt: &str, _context: &Context, thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -864,7 +864,7 @@ mod tests { _prompt: &str, _context: &Context, thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -911,7 +911,7 @@ mod tests { _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -1054,7 +1054,7 @@ Some text in between. _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -1092,7 +1092,7 @@ Some text in between. prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -1161,7 +1161,7 @@ Some text in between. prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { diff --git a/lib/crates/fabro-workflow/src/handler/command.rs b/lib/crates/fabro-workflow/src/handler/command.rs index af65d97bf..b03438daa 100644 --- a/lib/crates/fabro-workflow/src/handler/command.rs +++ b/lib/crates/fabro-workflow/src/handler/command.rs @@ -195,7 +195,7 @@ mod tests { let store = test_store(); let run_store = store.create_run(&fixtures::RUN_1).await.unwrap(); let services = EngineServices { - emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)), + emitter: Arc::new(crate::event::Emitter::new(fixtures::RUN_1)), run_store: run_store.clone(), ..EngineServices::test_default() }; diff --git a/lib/crates/fabro-workflow/src/handler/fan_in.rs b/lib/crates/fabro-workflow/src/handler/fan_in.rs index 3989aaa4b..ae3560dc2 100644 --- a/lib/crates/fabro-workflow/src/handler/fan_in.rs +++ b/lib/crates/fabro-workflow/src/handler/fan_in.rs @@ -4,7 +4,7 @@ use std::sync::Arc; use crate::context::Context; use crate::context::keys; use crate::error::FabroError; -use crate::event::{Event, EventEmitter}; +use crate::event::{Emitter, Event}; use crate::outcome::{Outcome, OutcomeExt}; use crate::run_dir::visit_from_context; use crate::sandbox_git::git_merge_ff_only; @@ -220,7 +220,7 @@ async fn llm_evaluate( context: &Context, _run_dir: &Path, node_id: &str, - emitter: &Arc, + emitter: &Arc, sandbox: &Arc, ) -> Result { let results_text = @@ -458,7 +458,7 @@ mod tests { _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { diff --git a/lib/crates/fabro-workflow/src/handler/human.rs b/lib/crates/fabro-workflow/src/handler/human.rs index 1cdb3eec8..695c86d3a 100644 --- a/lib/crates/fabro-workflow/src/handler/human.rs +++ b/lib/crates/fabro-workflow/src/handler/human.rs @@ -7,7 +7,7 @@ use async_trait::async_trait; use crate::context::Context; use crate::context::keys; use crate::error::FabroError; -use crate::event::{Event, EventEmitter}; +use crate::event::{Emitter, Event}; use crate::millis_u64; use crate::outcome::{Outcome, OutcomeExt}; use fabro_graphviz::graph::{Graph, Node}; @@ -68,7 +68,7 @@ fn parse_accelerator_key(label: &str) -> String { /// Blocks until a human selects an option derived from outgoing edges. pub struct HumanHandler { interviewer: Arc, - emitter: Option>, + emitter: Option>, } impl HumanHandler { @@ -80,7 +80,7 @@ impl HumanHandler { } #[must_use] - pub fn with_emitter(mut self, emitter: Arc) -> Self { + pub fn with_emitter(mut self, emitter: Arc) -> Self { self.emitter = Some(emitter); self } diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs index 4eae3e871..18b17f5dc 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/api.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs @@ -19,7 +19,7 @@ use super::super::agent::{CodergenBackend, CodergenResult}; use crate::context::keys::Fidelity; use crate::context::{Context, WorkflowContext}; use crate::error::FabroError; -use crate::event::{Event, EventEmitter}; +use crate::event::{Emitter, Event}; use crate::outcome::StageUsage; use crate::outcome::compute_stage_cost; use crate::run_dir::visit_from_context; @@ -88,7 +88,7 @@ fn spawn_event_forwarder( session: &Session, node_id: String, visit: u32, - emitter: Arc, + emitter: Arc, file_tracking: Arc>, ) { let mut rx = session.subscribe(); @@ -410,7 +410,7 @@ impl CodergenBackend for AgentApiBackend { prompt: &str, context: &Context, thread_id: Option<&str>, - emitter: &Arc, + emitter: &Arc, sandbox: &Arc, tool_hooks: Option>, ) -> Result { diff --git a/lib/crates/fabro-workflow/src/handler/llm/cli.rs b/lib/crates/fabro-workflow/src/handler/llm/cli.rs index ead561430..70bd13df7 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/cli.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/cli.rs @@ -10,7 +10,7 @@ use tokio::time::sleep; use super::super::agent::{CodergenBackend, CodergenResult}; use crate::context::Context; use crate::error::FabroError; -use crate::event::{Event, EventEmitter}; +use crate::event::{Emitter, Event}; use crate::outcome::StageUsage; use crate::outcome::compute_stage_cost; use crate::run_dir::visit_from_context; @@ -67,7 +67,7 @@ async fn ensure_cli( cli: AgentCli, provider: Provider, sandbox: &Arc, - emitter: &Arc, + emitter: &Arc, ) -> Result<(), FabroError> { let start = std::time::Instant::now(); let cli_name = cli.name(); @@ -463,7 +463,7 @@ impl CodergenBackend for AgentCliBackend { prompt: &str, _context: &Context, _thread_id: Option<&str>, - emitter: &Arc, + emitter: &Arc, sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -752,7 +752,7 @@ impl CodergenBackend for BackendRouter { prompt: &str, context: &Context, thread_id: Option<&str>, - emitter: &Arc, + emitter: &Arc, sandbox: &Arc, tool_hooks: Option>, ) -> Result { @@ -944,7 +944,7 @@ mod tests { vec![ok_result()], Arc::clone(&commands), )); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = ensure_cli(AgentCli::Claude, Provider::Anthropic, &sandbox, &emitter).await; assert!(result.is_ok()); @@ -965,7 +965,7 @@ mod tests { ], Arc::clone(&commands), )); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = ensure_cli(AgentCli::Claude, Provider::Anthropic, &sandbox, &emitter).await; assert!(result.is_ok()); @@ -985,7 +985,7 @@ mod tests { ], Arc::clone(&commands), )); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = ensure_cli(AgentCli::Claude, Provider::Anthropic, &sandbox, &emitter).await; assert!(result.is_err()); @@ -1190,7 +1190,7 @@ mod tests { _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { diff --git a/lib/crates/fabro-workflow/src/handler/mod.rs b/lib/crates/fabro-workflow/src/handler/mod.rs index 9ca028451..96ea888ca 100644 --- a/lib/crates/fabro-workflow/src/handler/mod.rs +++ b/lib/crates/fabro-workflow/src/handler/mod.rs @@ -28,7 +28,7 @@ use object_store::memory::InMemory; use crate::context::Context; use crate::error::FabroError; -use crate::event::EventEmitter; +use crate::event::Emitter; use crate::outcome::{Outcome, OutcomeExt}; use crate::sandbox_git::GitState; use fabro_graphviz::graph::{Graph, Node, shape_to_handler_type}; @@ -38,7 +38,7 @@ use fabro_interview::Interviewer; /// Shared services available to all handlers during execution. pub struct EngineServices { pub registry: Arc, - pub emitter: Arc, + pub emitter: Arc, pub sandbox: Arc, pub run_store: SlateRunStore, /// Git state for the current run. Set via `set_git_state` at the start of @@ -82,7 +82,7 @@ impl EngineServices { )); Self { registry: Arc::new(HandlerRegistry::new(Box::new(start::StartHandler))), - emitter: Arc::new(EventEmitter::default()), + emitter: Arc::new(Emitter::default()), sandbox: Arc::new(fabro_agent::LocalSandbox::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index 9c410bd91..f4aace1a8 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -627,7 +627,7 @@ mod tests { let store = test_store(); let run_store = store.create_run(&fixtures::RUN_1).await.unwrap(); let services = EngineServices { - emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)), + emitter: Arc::new(crate::event::Emitter::new(fixtures::RUN_1)), run_store: run_store.clone(), ..EngineServices::test_default() }; @@ -679,7 +679,7 @@ mod tests { let store = test_store(); let run_store = store.create_run(&fixtures::RUN_1).await.unwrap(); let services = EngineServices { - emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)), + emitter: Arc::new(crate::event::Emitter::new(fixtures::RUN_1)), run_store: run_store.clone(), ..EngineServices::test_default() }; diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs index 9db3e5dac..c8af00352 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -202,7 +202,7 @@ mod tests { let store = test_store(); let run_store = store.create_run(&fixtures::RUN_1).await.unwrap(); let services = EngineServices { - emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)), + emitter: Arc::new(crate::event::Emitter::new(fixtures::RUN_1)), run_store: run_store.clone(), ..EngineServices::test_default() }; @@ -260,7 +260,7 @@ mod tests { _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -320,7 +320,7 @@ mod tests { _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -377,7 +377,7 @@ mod tests { _prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { diff --git a/lib/crates/fabro-workflow/src/lifecycle/artifact.rs b/lib/crates/fabro-workflow/src/lifecycle/artifact.rs index be3330d69..4661a5a4e 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/artifact.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/artifact.rs @@ -12,7 +12,7 @@ use fabro_core::state::ExecutionState; use crate::artifact::{offload_large_values, sync_artifacts_to_env}; use crate::artifact_snapshot::collect_artifacts; -use crate::event::{Event, EventEmitter, RunNoticeLevel}; +use crate::event::{Emitter, Event, RunNoticeLevel}; use crate::graph::WorkflowGraph; use crate::graph::WorkflowNode; use crate::outcome::StageUsage; @@ -28,7 +28,7 @@ pub(crate) struct ArtifactLifecycle { pub sandbox: Arc, pub run_store: SlateRunStore, pub blob_cache_dir: PathBuf, - pub emitter: Arc, + pub emitter: Arc, pub artifacts_dir: PathBuf, pub artifact_globs: Vec, pub captured_artifact_count: Arc, @@ -42,7 +42,7 @@ impl ArtifactLifecycle { sandbox: Arc, run_store: SlateRunStore, blob_cache_dir: PathBuf, - emitter: Arc, + emitter: Arc, artifacts_dir: PathBuf, artifact_globs: Vec, captured_artifact_count: Arc, diff --git a/lib/crates/fabro-workflow/src/lifecycle/event.rs b/lib/crates/fabro-workflow/src/lifecycle/event.rs index 9714bb0c7..46fc955b7 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/event.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/event.rs @@ -17,7 +17,7 @@ use super::circuit_breaker::CircuitBreakerLifecycle; use super::git::GitCheckpointResult; use crate::context; use crate::error::FabroError; -use crate::event::{Event, EventEmitter}; +use crate::event::{Emitter, Event}; use crate::graph::WorkflowGraph; use crate::graph::WorkflowNode; use crate::outcome::{ @@ -34,7 +34,7 @@ type FailureSignatureSnapshot = ( /// Sub-lifecycle responsible for emitting workflow run events. pub(crate) struct EventLifecycle { - pub emitter: Arc, + pub emitter: Arc, pub graph_name: String, pub run_id: RunId, pub run_start: Mutex, diff --git a/lib/crates/fabro-workflow/src/lifecycle/git.rs b/lib/crates/fabro-workflow/src/lifecycle/git.rs index 492030bbf..a92d6ea84 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/git.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/git.rs @@ -13,7 +13,7 @@ use fabro_core::lifecycle::RunLifecycle; use fabro_core::outcome::NodeResult; use fabro_core::state::ExecutionState; -use crate::event::{Event, EventEmitter, RunNoticeLevel}; +use crate::event::{Emitter, Event, RunNoticeLevel}; use crate::git::MetadataStore; use crate::graph::WorkflowGraph; use crate::graph::WorkflowNode; @@ -63,7 +63,7 @@ pub(crate) struct GitCheckpointResult { /// Sub-lifecycle responsible for git operations (checkpoint commits, pushes, diffs). pub(crate) struct GitLifecycle { pub sandbox: Arc, - pub emitter: Arc, + pub emitter: Arc, pub run_dir: PathBuf, pub run_id: RunId, pub run_store: SlateRunStore, diff --git a/lib/crates/fabro-workflow/src/lifecycle/mod.rs b/lib/crates/fabro-workflow/src/lifecycle/mod.rs index 7ec808f46..5d87b5809 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/mod.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/mod.rs @@ -27,7 +27,7 @@ use fabro_core::state::ExecutionState; use crate::context; use crate::error::{FailureSignature, FailureSignatureExt}; -use crate::event::EventEmitter; +use crate::event::Emitter; use crate::graph::WorkflowGraph; use crate::graph::WorkflowNode; use crate::outcome::{Outcome, StageUsage}; @@ -76,7 +76,7 @@ pub(crate) struct WorkflowLifecycle { impl WorkflowLifecycle { #[allow(clippy::too_many_arguments)] pub(crate) fn new( - emitter: &Arc, + emitter: &Arc, hook_runner: Option>, sandbox: &Arc, graph: Arc, diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs index 0e9ea3563..dc3cabe87 100644 --- a/lib/crates/fabro-workflow/src/operations/start.rs +++ b/lib/crates/fabro-workflow/src/operations/start.rs @@ -15,7 +15,7 @@ use fabro_types::{RunId, Settings}; use crate::context::Context; use crate::error::FabroError; use crate::event::{ - Event, EventBody, EventEmitter, RunNoticeLevel, StoreProgressLogger, append_event, + Emitter, Event, EventBody, RunNoticeLevel, StoreProgressLogger, append_event, event_payload_from_redacted_json, redacted_event_json, to_run_event, }; use crate::git::MetadataStore; @@ -37,7 +37,7 @@ use tokio::runtime::Handle; struct RunSession { cancel_token: Option>, - emitter: Arc, + emitter: Arc, sandbox: SandboxSpec, llm: LlmSpec, interviewer: Arc, @@ -63,7 +63,7 @@ struct RunSession { pub struct StartServices { pub run_id: RunId, pub cancel_token: Option>, - pub emitter: Arc, + pub emitter: Arc, pub interviewer: Arc, pub run_store: SlateRunStore, pub github_app: Option, @@ -811,7 +811,7 @@ mod tests { use super::*; use crate::context::Context; - use crate::event::EventEmitter; + use crate::event::Emitter; use crate::handler::HandlerRegistry; use crate::handler::exit::ExitHandler; use crate::handler::start::StartHandler; @@ -871,7 +871,7 @@ mod tests { async fn test_start_services( store: &SlateStore, _run_dir: &Path, - emitter: Arc, + emitter: Arc, registry: Arc, ) -> StartServices { StartServices { @@ -890,7 +890,7 @@ mod tests { async fn start_captures_checkpoint_git_sha_in_conclusion() { let temp = tempfile::tempdir().unwrap(); let run_dir = temp.path().join("run"); - let emitter = Arc::new(EventEmitter::new(fixtures::RUN_1)); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let injected = Arc::new(AtomicBool::new(false)); @@ -944,7 +944,7 @@ mod tests { async fn start_loads_persisted_from_run_dir() { let temp = tempfile::tempdir().unwrap(); let run_dir = temp.path().join("run"); - let emitter = Arc::new(EventEmitter::new(fixtures::RUN_1)); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; @@ -965,7 +965,7 @@ mod tests { async fn start_invokes_on_node_callback_before_execution() { let temp = tempfile::tempdir().unwrap(); let run_dir = temp.path().join("run"); - let emitter = Arc::new(EventEmitter::new(fixtures::RUN_1)); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let visited = Arc::new(Mutex::new(Vec::new())); @@ -994,7 +994,7 @@ mod tests { async fn start_errors_when_checkpoint_exists() { let temp = tempfile::tempdir().unwrap(); let run_dir = temp.path().join("run"); - let emitter = Arc::new(EventEmitter::new(fixtures::RUN_1)); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; @@ -1063,7 +1063,7 @@ mod tests { async fn resume_errors_when_checkpoint_missing() { let temp = tempfile::tempdir().unwrap(); let run_dir = temp.path().join("run"); - let emitter = Arc::new(EventEmitter::new(fixtures::RUN_1)); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; @@ -1086,7 +1086,7 @@ mod tests { let temp = tempfile::tempdir().unwrap(); let run_dir = temp.path().join("run"); std::fs::create_dir_all(&run_dir).unwrap(); - let emitter = Arc::new(EventEmitter::new(fixtures::RUN_1)); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); let registry = Arc::new(test_registry()); let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &run_dir).await; diff --git a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs index 6344492ca..f9a3a1958 100644 --- a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs +++ b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs @@ -19,7 +19,7 @@ use object_store::memory::InMemory; use super::*; use crate::context::{self, Context}; use crate::error::FabroError; -use crate::event::{EventEmitter, StoreProgressLogger}; +use crate::event::{Emitter, StoreProgressLogger}; use crate::handler::start::StartHandler; use crate::handler::{Handler as HandlerTrait, HandlerRegistry}; use crate::outcome::{Outcome, OutcomeExt, StageStatus}; @@ -76,11 +76,11 @@ fn test_run_id(label: &str) -> RunId { } } -fn test_emitter(label: &str) -> EventEmitter { - EventEmitter::new(test_run_id(label)) +fn test_emitter(label: &str) -> Emitter { + Emitter::new(test_run_id(label)) } -fn test_emitter_arc(label: &str) -> Arc { +fn test_emitter_arc(label: &str) -> Arc { Arc::new(test_emitter(label)) } @@ -290,7 +290,7 @@ async fn execute_runs_start_to_exit_and_returns_final_context() { async fn run_with_lifecycle( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, graph: &Graph, run_options: RunOptions, diff --git a/lib/crates/fabro-workflow/src/pipeline/finalize.rs b/lib/crates/fabro-workflow/src/pipeline/finalize.rs index 498cb2ac7..64ec14e96 100644 --- a/lib/crates/fabro-workflow/src/pipeline/finalize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/finalize.rs @@ -1,7 +1,7 @@ use std::sync::Arc; use crate::error::FabroError; -use crate::event::{Event, EventEmitter, RunNoticeLevel}; +use crate::event::{Emitter, Event, RunNoticeLevel}; use crate::git::MetadataStore; use crate::outcome::{Outcome, OutcomeExt, StageStatus}; use crate::records::{Checkpoint, Conclusion, StageSummary}; @@ -15,7 +15,7 @@ use fabro_store::SlateRunStore; use super::types::{Concluded, FinalizeOptions, Retroed}; fn emit_run_notice( - emitter: &EventEmitter, + emitter: &Emitter, level: RunNoticeLevel, code: impl Into, message: impl Into, @@ -365,7 +365,7 @@ mod tests { std::fs::create_dir_all(&run_dir).unwrap(); let inner_store = test_store().create_run(&test_run_id()).await.unwrap(); let run_store = inner_store; - let emitter = Arc::new(EventEmitter::new(test_run_id())); + let emitter = Arc::new(Emitter::new(test_run_id())); let store_logger = StoreProgressLogger::new(run_store.clone()); store_logger.register(&emitter); let retroed = Retroed { diff --git a/lib/crates/fabro-workflow/src/pipeline/initialize.rs b/lib/crates/fabro-workflow/src/pipeline/initialize.rs index ef48e8486..490123a36 100644 --- a/lib/crates/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/initialize.rs @@ -15,7 +15,7 @@ use shlex::try_quote; use crate::devcontainer_bridge::{devcontainer_to_snapshot_config, run_devcontainer_lifecycle}; use crate::error::FabroError; -use crate::event::{Event, EventEmitter, RunNoticeLevel}; +use crate::event::{Emitter, Event, RunNoticeLevel}; use crate::git::{self, GitSyncStatus, MetadataStore}; use crate::handler::llm::{AgentApiBackend, AgentCliBackend, BackendRouter}; use crate::handler::{HandlerRegistry, default_registry}; @@ -47,7 +47,7 @@ async fn run_hooks( } fn emit_run_notice( - emitter: &EventEmitter, + emitter: &Emitter, level: RunNoticeLevel, code: impl Into, message: impl Into, @@ -230,7 +230,7 @@ async fn mint_github_token( async fn build_sandbox_env( spec: &SandboxEnvSpec, github_app: Option<&fabro_github::GitHubAppCredentials>, - emitter: &EventEmitter, + emitter: &Emitter, ) -> Result, FabroError> { let mut env = spec.devcontainer_env.clone(); env.extend(spec.toml_env.clone()); @@ -754,7 +754,7 @@ mod tests { std::fs::create_dir_all(&run_dir).unwrap(); let (graph, source) = simple_graph(); let persisted = test_persisted(graph, source.clone(), &run_dir); - let emitter = Arc::new(crate::event::EventEmitter::new(test_run_id())); + let emitter = Arc::new(crate::event::Emitter::new(test_run_id())); let initialized = initialize( persisted, @@ -822,7 +822,7 @@ mod tests { std::fs::create_dir_all(&run_dir).unwrap(); let (graph, source) = simple_graph(); let persisted = test_persisted(graph, source, &run_dir); - let emitter = Arc::new(crate::event::EventEmitter::new(test_run_id())); + let emitter = Arc::new(crate::event::Emitter::new(test_run_id())); let store = memory_store(); let run_store = store.create_run(&test_run_id()).await.unwrap(); let store_logger = StoreProgressLogger::new(run_store.clone()); diff --git a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs index 0a2108759..0d73535ba 100644 --- a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs +++ b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs @@ -9,7 +9,7 @@ use fabro_llm::generate::{GenerateParams, generate}; use fabro_util::text::strip_goal_decoration; use super::types::{Concluded, Finalized, PullRequestOptions}; -use crate::event::{Event, EventEmitter, RunNoticeLevel}; +use crate::event::{Emitter, Event, RunNoticeLevel}; use crate::outcome::{StageStatus, format_cost as outcome_format_cost}; use crate::records::{Conclusion, RunRecord}; use fabro_retro::retro::Retro; @@ -268,7 +268,7 @@ fn assemble_pr_body( } fn emit_run_notice( - emitter: &EventEmitter, + emitter: &Emitter, level: RunNoticeLevel, code: impl Into, message: impl Into, diff --git a/lib/crates/fabro-workflow/src/pipeline/retro.rs b/lib/crates/fabro-workflow/src/pipeline/retro.rs index 4b1f08162..d61ee3024 100644 --- a/lib/crates/fabro-workflow/src/pipeline/retro.rs +++ b/lib/crates/fabro-workflow/src/pipeline/retro.rs @@ -175,7 +175,7 @@ mod tests { use super::*; use crate::context::Context; - use crate::event::EventEmitter; + use crate::event::Emitter; use crate::event::{Event, StoreProgressLogger, append_event}; use crate::pipeline::types::Executed; use crate::records::{Checkpoint, CheckpointExt, RunRecord}; @@ -305,7 +305,7 @@ mod tests { let checkpoint = build_checkpoint(); let run_store = test_run_store(&run_dir, &checkpoint).await; - let emitter = Arc::new(EventEmitter::new(test_run_id())); + let emitter = Arc::new(Emitter::new(test_run_id())); let store_logger = StoreProgressLogger::new(run_store.clone()); store_logger.register(&emitter); let sandbox: Arc = Arc::new(fabro_agent::LocalSandbox::new( @@ -357,7 +357,7 @@ mod tests { std::fs::create_dir_all(&run_dir).unwrap(); let checkpoint = build_checkpoint(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let seen = Arc::new(Mutex::new(Vec::new())); emitter.on_event({ let seen = Arc::clone(&seen); diff --git a/lib/crates/fabro-workflow/src/pipeline/types.rs b/lib/crates/fabro-workflow/src/pipeline/types.rs index 2753440d3..ec3d05644 100644 --- a/lib/crates/fabro-workflow/src/pipeline/types.rs +++ b/lib/crates/fabro-workflow/src/pipeline/types.rs @@ -17,7 +17,7 @@ use fabro_validate::Diagnostic; use crate::context::Context; use crate::error::FabroError; -use crate::event::EventEmitter; +use crate::event::Emitter; use crate::handler::HandlerRegistry; use crate::outcome::Outcome; use crate::records::{Checkpoint, Conclusion, RunRecord}; @@ -229,7 +229,7 @@ pub struct InitOptions { pub run_id: RunId, pub run_store: SlateRunStore, pub dry_run: bool, - pub emitter: Arc, + pub emitter: Arc, pub sandbox: SandboxSpec, pub llm: LlmSpec, pub interviewer: Arc, @@ -254,7 +254,7 @@ pub struct Initialized { pub run_store: SlateRunStore, pub(crate) checkpoint: Option, pub(crate) seed_context: Option, - pub emitter: Arc, + pub emitter: Arc, pub sandbox: Arc, pub registry: Arc, pub on_node: crate::OnNodeCallback, @@ -274,7 +274,7 @@ pub struct Executed { pub run_options: RunOptions, pub run_store: SlateRunStore, pub hook_runner: Option>, - pub emitter: Arc, + pub emitter: Arc, pub sandbox: Arc, pub duration_ms: u64, pub final_context: Context, @@ -291,7 +291,7 @@ pub struct Retroed { pub run_options: RunOptions, pub run_store: SlateRunStore, pub hook_runner: Option>, - pub emitter: Arc, + pub emitter: Arc, pub sandbox: Arc, pub duration_ms: u64, pub retro: Option, @@ -306,7 +306,7 @@ pub struct Concluded { pub pushed_branch: Option, pub graph: Graph, pub run_options: RunOptions, - pub emitter: Arc, + pub emitter: Arc, } /// Output of the PULL_REQUEST phase. @@ -333,7 +333,7 @@ pub struct RetroOptions { pub goal: String, pub run_dir: PathBuf, pub sandbox: Arc, - pub emitter: Option>, + pub emitter: Option>, pub failed: bool, pub run_duration_ms: u64, pub enabled: bool, diff --git a/lib/crates/fabro-workflow/src/test_support.rs b/lib/crates/fabro-workflow/src/test_support.rs index 2be28e3ba..2a8f83364 100644 --- a/lib/crates/fabro-workflow/src/test_support.rs +++ b/lib/crates/fabro-workflow/src/test_support.rs @@ -9,7 +9,7 @@ use fabro_store::{RunProjection, SlateStore}; use object_store::local::LocalFileSystem; use crate::error::{FabroError, Result}; -use crate::event::{Event, EventEmitter, StoreProgressLogger, append_event}; +use crate::event::{Emitter, Event, StoreProgressLogger, append_event}; use crate::handler::HandlerRegistry; use crate::outcome::Outcome; use crate::pipeline; @@ -28,8 +28,8 @@ struct InitializedState { store_logger: StoreProgressLogger, } -fn bound_emitter(run_id: fabro_types::RunId, observer: &Arc) -> Arc { - let emitter = Arc::new(EventEmitter::new(run_id)); +fn bound_emitter(run_id: fabro_types::RunId, observer: &Arc) -> Arc { + let emitter = Arc::new(Emitter::new(run_id)); let observer_clone = Arc::clone(observer); emitter.on_event(move |event| observer_clone.dispatch_run_event(event)); emitter @@ -37,7 +37,7 @@ fn bound_emitter(run_id: fabro_types::RunId, observer: &Arc) -> Ar async fn initialized( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, graph: &GvGraph, run_options: &RunOptions, @@ -122,7 +122,7 @@ async fn initialized( pub async fn run_graph( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, graph: &GvGraph, run_options: &RunOptions, @@ -146,7 +146,7 @@ pub async fn run_graph( pub async fn run_graph_with_state( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, graph: &GvGraph, run_options: &RunOptions, @@ -177,7 +177,7 @@ pub async fn run_graph_with_state( pub async fn run_graph_with_hooks( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, graph: &GvGraph, run_options: &RunOptions, @@ -203,7 +203,7 @@ pub async fn run_graph_with_hooks( pub async fn run_graph_with_hooks_and_state( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, graph: &GvGraph, run_options: &RunOptions, @@ -236,7 +236,7 @@ pub async fn run_graph_with_hooks_and_state( pub async fn run_graph_from_checkpoint( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, graph: &GvGraph, run_options: &RunOptions, @@ -261,7 +261,7 @@ pub async fn run_graph_from_checkpoint( pub async fn run_graph_from_checkpoint_with_state( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, graph: &GvGraph, run_options: &RunOptions, @@ -293,7 +293,7 @@ pub async fn run_graph_from_checkpoint_with_state( pub struct WorkflowRunner { registry: std::sync::Mutex>, - emitter: Arc, + emitter: Arc, sandbox: Arc, } @@ -301,7 +301,7 @@ impl WorkflowRunner { #[must_use] pub fn new( registry: HandlerRegistry, - emitter: Arc, + emitter: Arc, sandbox: Arc, ) -> Self { Self { diff --git a/lib/crates/fabro-workflow/tests/it/daytona_integration.rs b/lib/crates/fabro-workflow/tests/it/daytona_integration.rs index ef1a65ef8..a319c88d9 100644 --- a/lib/crates/fabro-workflow/tests/it/daytona_integration.rs +++ b/lib/crates/fabro-workflow/tests/it/daytona_integration.rs @@ -26,7 +26,7 @@ use fabro_types::{RunId, Settings}; use fabro_workflow::artifact::sync_artifacts_to_env; use fabro_workflow::context::Context; use fabro_workflow::error::FabroError; -use fabro_workflow::event::EventEmitter; +use fabro_workflow::event::Emitter; use fabro_workflow::handler::exit::ExitHandler; use fabro_workflow::handler::start::StartHandler; use fabro_workflow::handler::{Handler, HandlerRegistry}; @@ -468,7 +468,7 @@ async fn daytona_pipeline_artifact_offload_and_sync() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), env.clone()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), env.clone()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -644,7 +644,7 @@ async fn daytona_git_checkpoint_remote_emits_events() { // Set up event collection let dir = tempfile::tempdir().unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = Arc::new(std::sync::Mutex::new(Vec::new())); { let events_clone = Arc::clone(&events); @@ -818,7 +818,7 @@ async fn daytona_parallel_git_branching_e2e() { graph.edges.push(Edge::new("fan_in", "exit")); let run_tmp = tempfile::tempdir().unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = Arc::new(std::sync::Mutex::new(Vec::new())); { let events_clone = Arc::clone(&events); @@ -1024,7 +1024,7 @@ async fn run_daytona_cli_test(provider: Provider, model: &str, install_command: let backend = AgentCliBackend::new(model.to_string(), provider); let node = Node::new("daytona_cli_test"); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = backend .run( @@ -1182,7 +1182,7 @@ async fn daytona_git_checkpoint_with_shadow_branch() { registry.register("exit", Box::new(ExitHandler)); let meta_branch = MetadataStore::branch_name(&run_id.to_string()); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), env.clone()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), env.clone()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -1287,7 +1287,7 @@ async fn daytona_asset_collection() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), env.clone()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), env.clone()); let mut graph = Graph::new("DaytonaAssetTest"); graph.attrs.insert( @@ -1575,7 +1575,7 @@ async fn daytona_git_push_run_branch_to_origin() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), env.clone()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), env.clone()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), diff --git a/lib/crates/fabro-workflow/tests/it/integration.rs b/lib/crates/fabro-workflow/tests/it/integration.rs index 7d7e33917..f38e57d84 100644 --- a/lib/crates/fabro-workflow/tests/it/integration.rs +++ b/lib/crates/fabro-workflow/tests/it/integration.rs @@ -29,7 +29,7 @@ use fabro_types::{RunEvent, RunId, Settings}; use fabro_validate::{Severity, validate, validate_or_raise}; use fabro_workflow::context::Context; use fabro_workflow::error::{FabroError, FailureSignatureExt}; -use fabro_workflow::event::{Event, EventEmitter}; +use fabro_workflow::event::{Emitter, Event}; use fabro_workflow::handler::agent::{AgentHandler, CodergenBackend, CodergenResult}; use fabro_workflow::handler::command::CommandHandler; use fabro_workflow::handler::conditional::ConditionalHandler; @@ -290,7 +290,7 @@ async fn end_to_end_linear_pipeline() { let dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new( make_linear_registry(), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -422,7 +422,7 @@ async fn end_to_end_branching_pipeline() { registry.register("agent", Box::new(AgentHandler::new(None))); registry.register("conditional", Box::new(ConditionalHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -541,7 +541,7 @@ async fn end_to_end_human_gate_pipeline() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -636,7 +636,7 @@ async fn human_gate_aborted_input_fails_closed_without_fail_route() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -746,7 +746,7 @@ async fn human_gate_aborted_input_routes_via_outcome_fail_condition() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -858,7 +858,7 @@ async fn goal_gate_routes_to_retry_target_on_failure() { registry.register("exit", Box::new(ExitHandler)); registry.register("always_fail", Box::new(AlwaysFailHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -978,7 +978,7 @@ async fn goal_gate_routes_to_retry_target_when_present() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -1289,7 +1289,7 @@ async fn retry_on_failure_then_succeed() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -1361,7 +1361,7 @@ async fn pipeline_with_many_nodes() { let dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new( make_linear_registry(), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -1451,7 +1451,7 @@ impl CodergenBackend for MockCodergenBackend { prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -1545,7 +1545,7 @@ impl Handler for ContextSetterHandler { } } -fn collect_events(emitter: &EventEmitter) -> Arc>> { +fn collect_events(emitter: &Emitter) -> Arc>> { let events = Arc::new(std::sync::Mutex::new(Vec::new())); let events_clone = Arc::clone(&events); emitter.on_event(move |event| { @@ -1683,7 +1683,7 @@ async fn smoke_test_with_mock_codergen_backend() { ); registry.register("conditional", Box::new(ConditionalHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -1784,7 +1784,7 @@ async fn end_to_end_parallel_fan_out_fan_in() { Box::new(FanInHandler::new(Some(Box::new(MockCodergenBackend)))), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -1896,7 +1896,7 @@ async fn resume_from_checkpoint_completes_pipeline() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -1994,7 +1994,7 @@ async fn resume_from_checkpoint_preserves_goal_gate_outcomes() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -2034,7 +2034,7 @@ async fn graph_goal_in_context() { let dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new( make_linear_registry(), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -2072,7 +2072,7 @@ async fn event_streaming_lifecycle() { }"#; let graph = parse(input).expect("parse"); let dir = tempfile::tempdir().unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = collect_events(&emitter); let engine = WorkflowRunner::new(make_linear_registry(), Arc::new(emitter), local_env()); let run_options = RunOptions { @@ -2147,7 +2147,7 @@ async fn context_flow_between_stages() { let dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new( make_linear_registry(), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -2202,7 +2202,7 @@ async fn tool_handler_e2e() { let interviewer = Arc::new(AutoApproveInterviewer); let engine = WorkflowRunner::new( make_full_registry(interviewer), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -2276,7 +2276,7 @@ async fn auto_approve_interviewer_e2e() { let interviewer = Arc::new(AutoApproveInterviewer); let engine = WorkflowRunner::new( make_full_registry(interviewer), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -2315,7 +2315,7 @@ async fn codergen_without_backend_simulated() { let dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new( make_linear_registry(), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -2421,7 +2421,7 @@ async fn branching_loop_back_on_failure() { call_count: std::sync::atomic::AtomicU32::new(0), }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -2506,7 +2506,7 @@ async fn human_gate_loops_back() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -2560,7 +2560,7 @@ async fn scenario_ship_a_feature() { let interviewer = Arc::new(AutoApproveInterviewer); let dir = tempfile::tempdir().unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = collect_events(&emitter); let engine = WorkflowRunner::new( make_full_registry(interviewer), @@ -2650,7 +2650,7 @@ async fn scenario_parallel_expert_review() { ); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -2736,7 +2736,7 @@ async fn scenario_node_retries_on_retry_status() { call_count: std::sync::atomic::AtomicU32::new(0), }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -2800,7 +2800,7 @@ async fn scenario_loop_restart_resets_context() { call_count: Arc::clone(&call_count), }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -2867,7 +2867,7 @@ async fn scenario_bug_triage_router() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); registry.register("conditional", Box::new(ConditionalHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -2928,7 +2928,7 @@ async fn scenario_crash_recovery() { let mut registry = HandlerRegistry::new(Box::new(StartHandler)); registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3036,7 +3036,7 @@ async fn manager_loop_stop_condition_satisfied_e2e() { registry.register("exit", Box::new(ExitHandler)); registry.register("done_setter", Box::new(DoneSetterHandler)); registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3117,7 +3117,7 @@ async fn manager_loop_max_cycles_exceeded_e2e() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3257,7 +3257,7 @@ async fn conditional_branching_success_fail_paths() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); registry.register("always_fail", Box::new(AlwaysFailHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3312,7 +3312,7 @@ async fn edge_selection_condition_match_wins_over_weight() { let mut registry = HandlerRegistry::new(Box::new(StartHandler)); registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3361,7 +3361,7 @@ async fn edge_selection_weight_breaks_ties() { let mut registry = HandlerRegistry::new(Box::new(StartHandler)); registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3402,7 +3402,7 @@ async fn edge_selection_lexical_tiebreak() { let mut registry = HandlerRegistry::new(Box::new(StartHandler)); registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3462,7 +3462,7 @@ async fn context_updates_visible_across_nodes() { registry.register("exit", Box::new(ExitHandler)); registry.register("conditional", Box::new(ConditionalHandler)); registry.register("context_setter", Box::new(ContextSetterHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3506,7 +3506,7 @@ async fn stylesheet_applies_model_override() { let dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new( make_linear_registry(), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -3563,7 +3563,7 @@ async fn custom_handler_registration_and_execution() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); registry.register("my_custom", Box::new(CustomHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3630,7 +3630,7 @@ async fn integration_smoke_plan_implement_review_done() { // Run pipeline let interviewer = Arc::new(AutoApproveInterviewer); let dir = tempfile::tempdir().unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = collect_events(&emitter); let engine = WorkflowRunner::new( make_full_registry(interviewer), @@ -3727,7 +3727,7 @@ async fn manager_loop_runs_child_engine_e2e() { registry.register("exit", Box::new(ExitHandler)); registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3860,7 +3860,7 @@ async fn manager_loop_context_flows_e2e() { registry.register("setter", Box::new(SetterHandler)); registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -3935,7 +3935,7 @@ async fn manager_loop_child_dotfile_e2e() { registry.register("exit", Box::new(ExitHandler)); registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4034,7 +4034,7 @@ async fn import_e2e_through_engine() { let engine = WorkflowRunner::new( make_linear_registry(), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions { @@ -4189,7 +4189,7 @@ async fn fidelity_default_is_compact() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4245,7 +4245,7 @@ async fn fidelity_graph_default_applied() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4297,7 +4297,7 @@ async fn fidelity_node_overrides_graph_default() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4355,7 +4355,7 @@ async fn fidelity_edge_overrides_node_and_graph() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4403,7 +4403,7 @@ async fn fidelity_full_produces_empty_preamble() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4461,7 +4461,7 @@ async fn fidelity_truncate_preamble_minimal() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4532,7 +4532,7 @@ async fn fidelity_summary_low_mode() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4598,7 +4598,7 @@ async fn fidelity_summary_medium_mode() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4664,7 +4664,7 @@ async fn fidelity_summary_high_mode() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4723,7 +4723,7 @@ async fn fidelity_full_sets_thread_id_in_context() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4793,7 +4793,7 @@ async fn fidelity_full_nodes_share_thread_id() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4873,7 +4873,7 @@ async fn fidelity_resume_degrades_full_to_summary_high() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -4969,7 +4969,7 @@ async fn fidelity_resume_degrade_only_affects_first_hop() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5052,7 +5052,7 @@ async fn fidelity_resume_no_degrade_when_not_full() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5093,7 +5093,7 @@ async fn fidelity_stored_in_checkpoint_context() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5181,7 +5181,7 @@ async fn fidelity_precedence_multi_node_pipeline() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5248,7 +5248,7 @@ async fn fidelity_compact_preamble_includes_completed_stages_and_context() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5322,8 +5322,7 @@ async fn fidelity_summary_low_excludes_context_values_in_pipeline() { captures: captures_low.clone(), }), ); - let engine_low = - WorkflowRunner::new(registry_low, Arc::new(EventEmitter::default()), local_env()); + let engine_low = WorkflowRunner::new(registry_low, Arc::new(Emitter::default()), local_env()); let run_options_low = RunOptions { settings: Settings::default(), run_dir: dir_low.path().to_path_buf(), @@ -5389,8 +5388,7 @@ async fn fidelity_summary_low_excludes_context_values_in_pipeline() { captures: captures_med.clone(), }), ); - let engine_med = - WorkflowRunner::new(registry_med, Arc::new(EventEmitter::default()), local_env()); + let engine_med = WorkflowRunner::new(registry_med, Arc::new(Emitter::default()), local_env()); let run_options_med = RunOptions { settings: Settings::default(), run_dir: dir_med.path().to_path_buf(), @@ -5460,7 +5458,7 @@ async fn fidelity_thread_id_fallback_to_previous_node_in_pipeline() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5513,7 +5511,7 @@ async fn fidelity_thread_id_from_node_class_in_pipeline() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5569,7 +5567,7 @@ async fn fidelity_edge_thread_id_override_in_pipeline() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5626,7 +5624,7 @@ async fn fidelity_full_without_explicit_thread_id_uses_previous_node() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5693,7 +5691,7 @@ async fn fidelity_from_parsed_dot_pipeline() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5740,7 +5738,7 @@ async fn fidelity_checkpoint_roundtrip_preserves_fidelity() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5811,7 +5809,7 @@ async fn fidelity_node_thread_id_overrides_edge_thread_id_in_pipeline() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5897,7 +5895,7 @@ async fn fidelity_resume_preserves_context_values_across_checkpoint() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -5965,7 +5963,7 @@ mod real_llm { prompt: &str, _context: &Context, _thread_id: Option<&str>, - _emitter: &Arc, + _emitter: &Arc, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -6064,7 +6062,7 @@ mod real_llm { use super::{load_checkpoint, local_env, test_run_id}; use fabro_graphviz::graph::{AttrValue, Edge, Graph}; use fabro_interview::AutoApproveInterviewer; - use fabro_workflow::event::EventEmitter; + use fabro_workflow::event::Emitter; use fabro_workflow::handler::HandlerRegistry; use fabro_workflow::handler::exit::ExitHandler; use fabro_workflow::handler::human::HumanHandler; @@ -6134,7 +6132,7 @@ mod real_llm { )))), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -6242,7 +6240,7 @@ mod real_llm { Box::new(AgentHandler::new(Some(make_llm_backend(client)))), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -6374,7 +6372,7 @@ mod real_llm { ); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -6474,7 +6472,7 @@ mod real_llm { ))), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -6567,7 +6565,7 @@ async fn human_gate_freeform_only_routes_text() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -6696,7 +6694,7 @@ async fn human_gate_freeform_with_fixed_choice_match() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -6810,7 +6808,7 @@ async fn human_gate_freeform_fallback_on_unmatched_text() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -6937,7 +6935,7 @@ async fn human_gate_freeform_sets_allow_freeform_on_question() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -7044,7 +7042,7 @@ async fn human_gate_without_freeform_sets_allow_freeform_false() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -7280,7 +7278,7 @@ fn hook_runner_from_defs(hooks: Vec) -> Arc, + emitter: Arc, hook_runner: Arc, } @@ -7316,15 +7314,15 @@ impl HookTestRunner { } } -fn emitter_with_events() -> (Arc, Arc>>) { - let emitter = EventEmitter::default(); +fn emitter_with_events() -> (Arc, Arc>>) { + let emitter = Emitter::default(); let events = collect_events(&emitter); (Arc::new(emitter), events) } fn engine_with_hooks(hooks: Vec) -> HookTestRunner { HookTestRunner { - emitter: Arc::new(EventEmitter::default()), + emitter: Arc::new(Emitter::default()), hook_runner: hook_runner_from_defs(hooks), } } @@ -8336,7 +8334,7 @@ async fn run_fidelity_prompt_pipeline(fidelity: &str) -> String { Box::new(AgentHandler::new(Some(Box::new(MockCodergenBackend)))), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -8534,7 +8532,7 @@ async fn large_context_values_are_offloaded_to_artifact_store() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = collect_events(&emitter); let engine = WorkflowRunner::new(registry, Arc::new(emitter), local_env()); let run_options = RunOptions { @@ -8758,11 +8756,7 @@ async fn artifact_pointers_rewritten_for_remote_sandbox() { registry.register("exit", Box::new(ExitHandler)); let remote_env = Arc::new(RemoteMockEnv::new("/sandbox")); - let engine = WorkflowRunner::new( - registry, - Arc::new(EventEmitter::default()), - remote_env.clone(), - ); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), remote_env.clone()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -8891,7 +8885,7 @@ async fn node_dir_uses_visit_count_on_revisit() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -9158,7 +9152,7 @@ async fn cli_backend_run_writes_prompt_and_calls_exec() { let node = Node::new("fix_code"); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = backend .run( @@ -9230,7 +9224,7 @@ async fn cli_backend_run_detects_changed_files() { let node = Node::new("implement"); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = backend .run( @@ -9263,7 +9257,7 @@ async fn cli_backend_run_with_codex_provider() { let node = Node::new("implement"); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = backend .run(&node, "Build the API", &context, None, &emitter, &env, None) @@ -9419,7 +9413,7 @@ async fn cli_backend_run_fails_on_nonzero_exit() { .with_poll_interval(Duration::from_millis(10)); let node = Node::new("step"); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let _ = env; // unused, just for the above struct @@ -9458,7 +9452,7 @@ async fn cli_backend_run_fails_on_unparseable_output() { let node = Node::new("step"); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = backend .run(&node, "do something", &context, None, &emitter, &env, None) @@ -9491,7 +9485,7 @@ async fn cli_backend_run_uses_node_model_override() { ); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); backend .run(&node, "test", &context, None, &emitter, &env, None) @@ -9532,7 +9526,7 @@ async fn cli_backend_run_uses_node_provider_override() { ); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); backend .run(&node, "test", &context, None, &emitter, &env, None) @@ -9557,7 +9551,7 @@ async fn cli_backend_run_returns_text_and_usage() { let node = Node::new("step"); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = backend .run(&node, "test", &context, None, &emitter, &env, None) @@ -9597,7 +9591,7 @@ async fn backend_router_delegates_to_cli_for_cli_node() { ); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = router .run(&node, "Fix the bug", &context, None, &emitter, &env, None) @@ -9631,7 +9625,7 @@ async fn backend_router_delegates_to_api_for_normal_node() { ); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = router .run(&node, "Plan the work", &context, None, &emitter, &env, None) @@ -9668,7 +9662,7 @@ async fn backend_router_delegates_to_cli_for_backend_attr() { ); let context = Context::new(); - let emitter = Arc::new(EventEmitter::default()); + let emitter = Arc::new(Emitter::default()); let result = router .run(&node, "Build it", &context, None, &emitter, &env, None) @@ -9760,7 +9754,7 @@ async fn full_pipeline_with_cli_backend_node() { ); let dir = tempfile::tempdir().unwrap(); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), env); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), env); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -9878,7 +9872,7 @@ async fn stylesheet_backend_property_routes_to_cli() { ); let dir = tempfile::tempdir().unwrap(); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), env); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), env); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -10059,7 +10053,7 @@ async fn git_checkpoint_host_emits_events_and_diff_patch() { // 4. Set up event collection and engine let run_dir = tempfile::tempdir().unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = collect_events(&emitter); let env: Arc = @@ -10234,7 +10228,7 @@ async fn git_checkpoint_host_writes_shadow_branch() { let run_dir = tempfile::tempdir().unwrap(); // Write graph.fabro so init_run can read it std::fs::write(run_dir.path().join("graph.fabro"), "digraph {}").unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let env: Arc = Arc::new(fabro_agent::LocalSandbox::new(worktree_path.clone())); @@ -10424,7 +10418,7 @@ async fn parallel_git_branching_host_e2e() { // 4. Set up engine with FileWriterHandler for branches let run_dir = tempfile::tempdir().unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = collect_events(&emitter); let env: Arc = @@ -10694,7 +10688,7 @@ async fn git_checkpoint_host_skips_empty_diff_patch() { graph.edges.push(Edge::new("work", "exit")); let run_dir = tempfile::tempdir().unwrap(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let _events = collect_events(&emitter); let env: Arc = @@ -11077,7 +11071,7 @@ async fn e2e_circuit_breaker_deterministic_self_loop() { )), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11123,7 +11117,7 @@ async fn e2e_circuit_breaker_custom_limit() { Box::new(DeterministicFailHandler::new("same error every time")), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11162,7 +11156,7 @@ async fn e2e_circuit_breaker_ignores_transient_failures() { registry.register("exit", Box::new(ExitHandler)); registry.register("test_handler", Box::new(TransientInfraFailHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11208,7 +11202,7 @@ async fn e2e_circuit_breaker_different_reasons_separate_counters() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11247,7 +11241,7 @@ async fn e2e_circuit_breaker_loop_restart() { Box::new(DeterministicFailHandler::new("verify step failed")), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11308,7 +11302,7 @@ async fn e2e_failure_signature_persisted_in_context() { Box::new(DeterministicFailHandler::new("test assertion failed")), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11371,7 +11365,7 @@ async fn e2e_failure_signature_hint_overrides_reason_in_context() { registry.register("exit", Box::new(ExitHandler)); registry.register("hint_handler", Box::new(SignatureHintHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11426,7 +11420,7 @@ async fn e2e_signature_maps_persist_in_checkpoint() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11541,7 +11535,7 @@ async fn e2e_circuit_breaker_emits_events_before_abort() { let dir = tempfile::tempdir().unwrap(); let graph = circuit_breaker_self_loop_graph(Some(3)); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = collect_events(&emitter); let mut registry = HandlerRegistry::new(Box::new(StartHandler)); @@ -11616,7 +11610,7 @@ async fn e2e_circuit_breaker_does_not_fire_below_limit() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11711,7 +11705,7 @@ async fn e2e_circuit_breaker_multi_stage_impl_verify_cycle() { )), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11807,7 +11801,7 @@ async fn e2e_loop_restart_blocked_for_deterministic_failure() { Box::new(ClassifiedFailHandler::always("deterministic")), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11846,7 +11840,7 @@ async fn e2e_loop_restart_blocked_for_structural_failure() { Box::new(ClassifiedFailHandler::always("structural")), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11885,7 +11879,7 @@ async fn e2e_loop_restart_blocked_for_budget_exhausted_failure() { Box::new(ClassifiedFailHandler::always("budget_exhausted")), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11924,7 +11918,7 @@ async fn e2e_loop_restart_blocked_for_canceled_failure() { Box::new(ClassifiedFailHandler::always("canceled")), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -11960,7 +11954,7 @@ async fn e2e_loop_restart_blocked_for_compilation_loop_failure() { Box::new(ClassifiedFailHandler::always("compilation_loop")), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -12000,7 +11994,7 @@ async fn e2e_loop_restart_allowed_for_transient_infra() { Box::new(ClassifiedFailHandler::succeed_on("transient_infra", 1)), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -12102,7 +12096,7 @@ async fn e2e_stall_watchdog_triggers_from_dot_parsed_pipeline() { let events = Arc::new(std::sync::Mutex::new(Vec::new())); let events_clone = events.clone(); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); emitter.on_event(move |event| { events_clone.lock().unwrap().push(format!("{event:?}")); }); @@ -12162,7 +12156,7 @@ async fn e2e_stall_watchdog_kept_alive_by_handler_events() { }), ); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -12207,7 +12201,7 @@ async fn e2e_stall_watchdog_disabled_with_zero_timeout() { registry.register("exit", Box::new(ExitHandler)); registry.register("slow", Box::new(SlowTestHandler { sleep_ms: 50 })); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -12271,7 +12265,7 @@ async fn e2e_stall_watchdog_with_explicit_timeout_override() { registry.register("exit", Box::new(ExitHandler)); registry.register("hanging", Box::new(HangingHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), local_env()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), local_env()); let run_options = RunOptions { settings: Settings::default(), run_dir: dir.path().to_path_buf(), @@ -12365,7 +12359,7 @@ async fn asset_collection_local_sandbox_success() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let emitter = EventEmitter::default(); + let emitter = Emitter::default(); let events = collect_events(&emitter); let engine = WorkflowRunner::new(registry, Arc::new(emitter), sandbox.clone()); @@ -12495,7 +12489,7 @@ async fn asset_collection_local_sandbox_on_failure() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), sandbox.clone()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), sandbox.clone()); let mut graph = Graph::new("AssetCollectionFailTest"); graph.attrs.insert( @@ -12586,7 +12580,7 @@ async fn asset_collection_docker_sandbox() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); - let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::default()), sandbox.clone()); + let engine = WorkflowRunner::new(registry, Arc::new(Emitter::default()), sandbox.clone()); let mut graph = Graph::new("DockerAssetTest"); graph.attrs.insert( @@ -12685,7 +12679,7 @@ async fn wait_timer_e2e() { let interviewer = Arc::new(AutoApproveInterviewer); let engine = WorkflowRunner::new( make_full_registry(interviewer), - Arc::new(EventEmitter::default()), + Arc::new(Emitter::default()), local_env(), ); let run_options = RunOptions {