use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Instant;
use tokio::sync::broadcast;
use crate::analysis::tier_allocator::ContextTier;
#[derive(Debug, Clone)]
pub enum EvolutionEvent {
AgentFocus {
agent_id: String,
agent_role: String,
file_path: String,
tier: ContextTier,
},
AgentDefocus { agent_id: String },
AgentStream {
agent_id: String,
content_preview: String,
tokens: usize,
},
ToolInvoked {
agent_id: String,
tool_name: String,
target_file: Option<String>,
},
ToolCompleted {
agent_id: String,
tool_name: String,
success: bool,
duration_ms: u64,
},
AgentStateChange {
agent_id: String,
state: AgentActivityState,
},
Throughput(ThroughputSnapshot),
TierUpdate {
focus_node: String,
total_files: usize,
tiers: Vec<TierEntry>,
},
BuildResult {
agent_id: String,
worktree: String,
success: bool,
error_count: usize,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AgentActivityState {
Idle,
Thinking,
ToolCall,
Streaming,
Verifying,
Done,
}
impl std::fmt::Display for AgentActivityState {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Idle => write!(f, "idle"),
Self::Thinking => write!(f, "thinking"),
Self::ToolCall => write!(f, "tool_call"),
Self::Streaming => write!(f, "streaming"),
Self::Verifying => write!(f, "verifying"),
Self::Done => write!(f, "done"),
}
}
}
#[derive(Debug, Clone)]
pub struct TierEntry {
pub file_path: String,
pub tier: ContextTier,
pub hops: usize,
}
#[derive(Debug, Clone)]
pub struct ThroughputSnapshot {
pub tokens_in_per_sec: f64,
pub tokens_out_per_sec: f64,
pub concurrent_requests: usize,
pub active_agents: Vec<String>,
pub total_tokens_session: u64,
pub timestamp: Instant,
}
#[derive(Debug)]
pub struct ThroughputTracker {
tokens_in: AtomicU64,
tokens_out: AtomicU64,
concurrent_requests: AtomicUsize,
total_tokens: AtomicU64,
last_snapshot: std::sync::Mutex<Instant>,
last_in: AtomicU64,
last_out: AtomicU64,
}
impl ThroughputTracker {
pub fn new() -> Self {
Self {
tokens_in: AtomicU64::new(0),
tokens_out: AtomicU64::new(0),
concurrent_requests: AtomicUsize::new(0),
total_tokens: AtomicU64::new(0),
last_snapshot: std::sync::Mutex::new(Instant::now()),
last_in: AtomicU64::new(0),
last_out: AtomicU64::new(0),
}
}
pub fn record_tokens_in(&self, count: usize) {
self.tokens_in.fetch_add(count as u64, Ordering::Relaxed);
self.total_tokens.fetch_add(count as u64, Ordering::Relaxed);
}
pub fn record_tokens_out(&self, count: usize) {
self.tokens_out.fetch_add(count as u64, Ordering::Relaxed);
self.total_tokens.fetch_add(count as u64, Ordering::Relaxed);
}
pub fn request_started(&self) {
self.concurrent_requests.fetch_add(1, Ordering::Relaxed);
}
pub fn request_completed(&self) {
self.concurrent_requests.fetch_sub(1, Ordering::Relaxed);
}
pub fn concurrent_requests(&self) -> usize {
self.concurrent_requests.load(Ordering::Relaxed)
}
pub fn snapshot(&self, active_agents: Vec<String>) -> ThroughputSnapshot {
let now = Instant::now();
let current_in = self.tokens_in.load(Ordering::Relaxed);
let current_out = self.tokens_out.load(Ordering::Relaxed);
let (elapsed_secs, prev_in, prev_out) = {
let mut last = self.last_snapshot.lock().unwrap_or_else(|e| e.into_inner());
let elapsed = now.duration_since(*last).as_secs_f64().max(0.001);
let prev_in = self.last_in.swap(current_in, Ordering::Relaxed);
let prev_out = self.last_out.swap(current_out, Ordering::Relaxed);
*last = now;
(elapsed, prev_in, prev_out)
};
let delta_in = current_in.saturating_sub(prev_in) as f64;
let delta_out = current_out.saturating_sub(prev_out) as f64;
ThroughputSnapshot {
tokens_in_per_sec: delta_in / elapsed_secs,
tokens_out_per_sec: delta_out / elapsed_secs,
concurrent_requests: self.concurrent_requests.load(Ordering::Relaxed),
active_agents,
total_tokens_session: self.total_tokens.load(Ordering::Relaxed),
timestamp: now,
}
}
}
impl Default for ThroughputTracker {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone)]
pub struct EvolutionBus {
tx: broadcast::Sender<EvolutionEvent>,
throughput: Arc<ThroughputTracker>,
}
impl EvolutionBus {
pub fn new(capacity: usize) -> Self {
let (tx, _) = broadcast::channel(capacity);
Self {
tx,
throughput: Arc::new(ThroughputTracker::new()),
}
}
pub fn emit(&self, event: EvolutionEvent) {
let _ = self.tx.send(event);
}
pub fn subscribe(&self) -> broadcast::Receiver<EvolutionEvent> {
self.tx.subscribe()
}
pub fn throughput(&self) -> &ThroughputTracker {
&self.throughput
}
pub fn emit_throughput(&self, active_agents: Vec<String>) {
let snapshot = self.throughput.snapshot(active_agents);
self.emit(EvolutionEvent::Throughput(snapshot));
}
pub fn subscriber_count(&self) -> usize {
self.tx.receiver_count()
}
}
impl Default for EvolutionBus {
fn default() -> Self {
Self::new(256)
}
}
pub struct EvolutionBridgeEmitter {
bus: EvolutionBus,
agent_id: String,
inner: Arc<dyn super::tui_events::EventEmitter>,
}
impl EvolutionBridgeEmitter {
pub fn new(
bus: EvolutionBus,
agent_id: String,
inner: Arc<dyn super::tui_events::EventEmitter>,
) -> Self {
Self {
bus,
agent_id,
inner,
}
}
}
impl super::tui_events::EventEmitter for EvolutionBridgeEmitter {
fn emit(&self, event: super::tui_events::AgentEvent) {
self.inner.emit(event.clone());
match event {
super::tui_events::AgentEvent::ToolStarted { name } => {
self.bus.emit(EvolutionEvent::ToolInvoked {
agent_id: self.agent_id.clone(),
tool_name: name,
target_file: None,
});
self.bus.emit(EvolutionEvent::AgentStateChange {
agent_id: self.agent_id.clone(),
state: AgentActivityState::ToolCall,
});
}
super::tui_events::AgentEvent::ToolCompleted {
name,
success,
duration_ms,
} => {
self.bus.emit(EvolutionEvent::ToolCompleted {
agent_id: self.agent_id.clone(),
tool_name: name,
success,
duration_ms,
});
}
super::tui_events::AgentEvent::TokenUsage {
prompt_tokens,
completion_tokens,
} => {
self.bus
.throughput()
.record_tokens_out(prompt_tokens as usize);
self.bus
.throughput()
.record_tokens_in(completion_tokens as usize);
}
super::tui_events::AgentEvent::AssistantDelta { ref text } => {
self.bus.emit(EvolutionEvent::AgentStream {
agent_id: self.agent_id.clone(),
content_preview: text.chars().take(80).collect(),
tokens: text.len() / 4, });
}
super::tui_events::AgentEvent::Started => {
self.bus.emit(EvolutionEvent::AgentStateChange {
agent_id: self.agent_id.clone(),
state: AgentActivityState::Thinking,
});
}
super::tui_events::AgentEvent::Completed { .. } => {
self.bus.emit(EvolutionEvent::AgentStateChange {
agent_id: self.agent_id.clone(),
state: AgentActivityState::Done,
});
}
_ => {} }
}
}
#[cfg(test)]
#[path = "../../tests/unit/agent/evolution_events/evolution_events_test.rs"]
mod tests;