mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-09 03:20:56 +00:00
Simplify: fix buggy JSON sorting, deduplicate event filtering, clean up wait loop
- Replace no-op sort_json_value (IndexMap→IndexMap) in create.rs with normalize_json_value (IndexMap→BTreeMap→Map) from event.rs, fixing RunCreated events having non-deterministic key order - Add AgentEvent::is_streaming_noise() to centralize the 6-variant streaming filter used in api.rs, retro.rs, and subagent.rs - Extract load_file_status closure and merge Ok(None)|Err(_) arms in wait.rs to remove triple-repeated RunStatusRecord::load expression Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
parent
d95dedf711
commit
70bfa53255
7 changed files with 51 additions and 48 deletions
|
|
@ -92,18 +92,14 @@ impl SubAgentManager {
|
|||
tokio::spawn(async move {
|
||||
while let Ok(event) = rx.recv().await {
|
||||
// Skip streaming / noise events
|
||||
if matches!(
|
||||
&event.event,
|
||||
AgentEvent::TextDelta { .. }
|
||||
| AgentEvent::AssistantOutputReplace { .. }
|
||||
| AgentEvent::ReasoningDelta { .. }
|
||||
| AgentEvent::ToolCallOutputDelta { .. }
|
||||
| AgentEvent::AssistantTextStart
|
||||
| AgentEvent::SessionStarted { .. }
|
||||
| AgentEvent::SessionEnded
|
||||
| AgentEvent::ProcessingEnd
|
||||
| AgentEvent::SkillExpanded { .. }
|
||||
) {
|
||||
if event.event.is_streaming_noise()
|
||||
|| matches!(
|
||||
&event.event,
|
||||
AgentEvent::SessionStarted { .. }
|
||||
| AgentEvent::SessionEnded
|
||||
| AgentEvent::ProcessingEnd
|
||||
)
|
||||
{
|
||||
continue;
|
||||
}
|
||||
cb(SubAgentCallbackEvent::Forwarded(event));
|
||||
|
|
|
|||
|
|
@ -201,6 +201,20 @@ pub enum AgentEvent {
|
|||
}
|
||||
|
||||
impl AgentEvent {
|
||||
/// Returns `true` for streaming-delta and UI-noise variants that are
|
||||
/// typically filtered out before forwarding to the workflow event stream.
|
||||
pub fn is_streaming_noise(&self) -> bool {
|
||||
matches!(
|
||||
self,
|
||||
Self::AssistantTextStart
|
||||
| Self::AssistantOutputReplace { .. }
|
||||
| Self::TextDelta { .. }
|
||||
| Self::ReasoningDelta { .. }
|
||||
| Self::ToolCallOutputDelta { .. }
|
||||
| Self::SkillExpanded { .. }
|
||||
)
|
||||
}
|
||||
|
||||
pub fn trace(&self, session_id: &str) {
|
||||
use tracing::{debug, error, info, warn};
|
||||
match self {
|
||||
|
|
|
|||
|
|
@ -36,13 +36,13 @@ pub(crate) async fn run(args: &WaitArgs, styles: &Styles, globals: &GlobalArgs)
|
|||
let started_waiting_at = std::time::Instant::now();
|
||||
|
||||
let final_status = loop {
|
||||
let load_file_status = || RunStatusRecord::load(&status_path).ok().map(|r| r.status);
|
||||
let status = match run_store.as_ref() {
|
||||
Some(run_store) => match run_store.get_status().await {
|
||||
Ok(Some(record)) => Some(record.status),
|
||||
Ok(None) => RunStatusRecord::load(&status_path).ok().map(|record| record.status),
|
||||
Err(_) => RunStatusRecord::load(&status_path).ok().map(|record| record.status),
|
||||
Ok(None) | Err(_) => load_file_status(),
|
||||
},
|
||||
None => RunStatusRecord::load(&status_path).ok().map(|record| record.status),
|
||||
None => load_file_status(),
|
||||
};
|
||||
let status = status.unwrap_or_else(|| {
|
||||
if started_waiting_at.elapsed() < WAIT_STARTUP_GRACE {
|
||||
|
|
|
|||
|
|
@ -462,7 +462,9 @@ impl WorkflowRunEvent {
|
|||
pub fn trace(&self) {
|
||||
use tracing::{debug, error, info, warn};
|
||||
match self {
|
||||
Self::RunCreated { run_id, run_dir, .. } => {
|
||||
Self::RunCreated {
|
||||
run_id, run_dir, ..
|
||||
} => {
|
||||
info!(run_id = %run_id, run_dir, "Run created");
|
||||
}
|
||||
Self::WorkflowRunStarted { name, run_id, .. } => {
|
||||
|
|
@ -665,7 +667,7 @@ impl WorkflowRunEvent {
|
|||
if *success {
|
||||
debug!(branch, "Git fetch succeeded");
|
||||
} else {
|
||||
warn!(branch, "Git fetch failed");
|
||||
warn!(branch, "Git fetch failed");
|
||||
}
|
||||
}
|
||||
Self::GitReset { sha } => {
|
||||
|
|
@ -852,10 +854,7 @@ impl WorkflowRunEvent {
|
|||
} => {
|
||||
debug!(
|
||||
node_id,
|
||||
exit_code,
|
||||
duration_ms,
|
||||
timed_out,
|
||||
"Command completed"
|
||||
exit_code, duration_ms, timed_out, "Command completed"
|
||||
);
|
||||
}
|
||||
Self::AgentCliStarted {
|
||||
|
|
@ -1374,7 +1373,7 @@ fn normalized_envelope_value(envelope: &RunEventEnvelope) -> Result<Value> {
|
|||
Ok(normalize_json_value(value))
|
||||
}
|
||||
|
||||
fn normalize_json_value(value: Value) -> Value {
|
||||
pub(crate) fn normalize_json_value(value: Value) -> Value {
|
||||
match value {
|
||||
Value::Object(map) => Value::Object(
|
||||
map.into_iter()
|
||||
|
|
|
|||
|
|
@ -96,16 +96,9 @@ fn spawn_event_forwarder(
|
|||
track_file_event(&event.event, &mut file_tracking.lock().unwrap());
|
||||
|
||||
// Forward non-streaming agent events to pipeline
|
||||
if !matches!(
|
||||
&event.event,
|
||||
AgentEvent::ProcessingEnd
|
||||
| AgentEvent::AssistantTextStart
|
||||
| AgentEvent::AssistantOutputReplace { .. }
|
||||
| AgentEvent::TextDelta { .. }
|
||||
| AgentEvent::ReasoningDelta { .. }
|
||||
| AgentEvent::ToolCallOutputDelta { .. }
|
||||
| AgentEvent::SkillExpanded { .. }
|
||||
) {
|
||||
if !event.event.is_streaming_noise()
|
||||
&& !matches!(&event.event, AgentEvent::ProcessingEnd)
|
||||
{
|
||||
emitter.emit(&WorkflowRunEvent::Agent {
|
||||
stage: node_id.clone(),
|
||||
event: event.event.clone(),
|
||||
|
|
|
|||
|
|
@ -1,6 +1,6 @@
|
|||
use std::collections::BTreeMap;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::Instant;
|
||||
use std::collections::BTreeMap;
|
||||
|
||||
use async_trait::async_trait;
|
||||
|
||||
|
|
@ -209,16 +209,25 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
|
|||
failure: outcome.failure.clone(),
|
||||
notes: outcome.notes.clone(),
|
||||
files_touched: outcome.files_touched.clone(),
|
||||
context_updates: (!outcome.context_updates.is_empty())
|
||||
.then(|| outcome.context_updates.clone().into_iter().collect::<BTreeMap<_, _>>()),
|
||||
context_updates: (!outcome.context_updates.is_empty()).then(|| {
|
||||
outcome
|
||||
.context_updates
|
||||
.clone()
|
||||
.into_iter()
|
||||
.collect::<BTreeMap<_, _>>()
|
||||
}),
|
||||
jump_to_node: outcome.jump_to_node.clone(),
|
||||
context_values: {
|
||||
let snapshot = state.context.snapshot();
|
||||
(!snapshot.is_empty())
|
||||
.then(|| snapshot.into_iter().collect::<BTreeMap<_, _>>())
|
||||
(!snapshot.is_empty()).then(|| snapshot.into_iter().collect::<BTreeMap<_, _>>())
|
||||
},
|
||||
node_visits: (!state.node_visits.is_empty())
|
||||
.then(|| state.node_visits.clone().into_iter().collect::<BTreeMap<_, _>>()),
|
||||
node_visits: (!state.node_visits.is_empty()).then(|| {
|
||||
state
|
||||
.node_visits
|
||||
.clone()
|
||||
.into_iter()
|
||||
.collect::<BTreeMap<_, _>>()
|
||||
}),
|
||||
loop_failure_signatures: None,
|
||||
restart_failure_signatures: None,
|
||||
attempt: result.attempts as usize,
|
||||
|
|
|
|||
|
|
@ -69,15 +69,7 @@ pub async fn run_retro(options: &RetroOptions, dry_run: bool) -> Option<Retro> {
|
|||
emitter_clone.map(|emitter| -> Arc<dyn Fn(SessionEvent) + Send + Sync> {
|
||||
Arc::new(move |event: SessionEvent| {
|
||||
emitter.touch();
|
||||
if !matches!(
|
||||
&event.event,
|
||||
fabro_agent::AgentEvent::AssistantTextStart
|
||||
| fabro_agent::AgentEvent::AssistantOutputReplace { .. }
|
||||
| fabro_agent::AgentEvent::TextDelta { .. }
|
||||
| fabro_agent::AgentEvent::ReasoningDelta { .. }
|
||||
| fabro_agent::AgentEvent::ToolCallOutputDelta { .. }
|
||||
| fabro_agent::AgentEvent::SkillExpanded { .. }
|
||||
) {
|
||||
if !event.event.is_streaming_noise() {
|
||||
emitter.emit(&WorkflowRunEvent::Agent {
|
||||
stage: "retro".to_string(),
|
||||
event: event.event.clone(),
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue