feat: enhance session activity tracking with directory support and cooldown logic
This commit is contained in:
@@ -24,6 +24,13 @@ struct EventEnvelope {
|
||||
properties: Value,
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct MultiplexedEventEnvelope {
|
||||
#[serde(default)]
|
||||
directory: Option<String>,
|
||||
payload: EventEnvelope,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, PartialEq)]
|
||||
pub enum ActivityPhase {
|
||||
Idle,
|
||||
@@ -31,6 +38,12 @@ pub enum ActivityPhase {
|
||||
Cooldown,
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
enum SseScope {
|
||||
Global,
|
||||
Directory(std::path::PathBuf),
|
||||
}
|
||||
|
||||
pub fn spawn_session_activity_tracker(
|
||||
app: AppHandle,
|
||||
runtime: DesktopRuntime,
|
||||
@@ -85,39 +98,8 @@ async fn run_once(
|
||||
};
|
||||
|
||||
let prefix = opencode.api_prefix();
|
||||
let mut url = format!("http://127.0.0.1:{port}{}/event", prefix);
|
||||
|
||||
if let Some(dir) = opencode.get_working_directory().to_str().map(|s| s.to_string()) {
|
||||
let mut parsed = reqwest::Url::parse(&url)?;
|
||||
parsed
|
||||
.query_pairs_mut()
|
||||
.append_pair("directory", &dir);
|
||||
url = parsed.to_string();
|
||||
}
|
||||
|
||||
debug!("[desktop:activity] Connecting SSE for activity phases: {url}");
|
||||
|
||||
let response = client
|
||||
.get(&url)
|
||||
.header("accept", "text/event-stream")
|
||||
.header("accept-encoding", "identity")
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
debug!(
|
||||
"[desktop:activity] SSE response status={} headers={:?}",
|
||||
response.status(),
|
||||
response.headers()
|
||||
);
|
||||
|
||||
if !response.status().is_success() {
|
||||
warn!(
|
||||
"[desktop:activity] SSE connect failed with status {}",
|
||||
response.status()
|
||||
);
|
||||
tokio::time::sleep(Duration::from_secs(2)).await;
|
||||
return Ok(());
|
||||
}
|
||||
let base = format!("http://127.0.0.1:{port}{prefix}");
|
||||
let (response, scope) = connect_activity_sse(runtime, client, &base).await?;
|
||||
|
||||
use tokio::io::AsyncBufReadExt;
|
||||
|
||||
@@ -130,12 +112,28 @@ async fn run_once(
|
||||
|
||||
loop {
|
||||
buf.clear();
|
||||
let bytes_read = match reader.read_until(b'\n', &mut buf).await {
|
||||
Ok(n) => n,
|
||||
Err(err) => {
|
||||
let bytes_read = match tokio::time::timeout(Duration::from_secs(2), reader.read_until(b'\n', &mut buf)).await
|
||||
{
|
||||
Ok(Ok(n)) => n,
|
||||
Ok(Err(err)) => {
|
||||
warn!("[desktop:activity] Read error in SSE stream: {err:?}");
|
||||
return Err(err.into());
|
||||
}
|
||||
Err(_) => {
|
||||
// No data received recently; if we are connected to a directory-scoped stream and the working directory
|
||||
// has changed, reconnect so activity tracking follows the new directory.
|
||||
if let SseScope::Directory(connected_dir) = &scope {
|
||||
let current_dir = opencode.get_working_directory();
|
||||
if current_dir != *connected_dir {
|
||||
debug!(
|
||||
"[desktop:activity] Working directory changed; reconnecting activity SSE (from {:?} to {:?})",
|
||||
connected_dir, current_dir
|
||||
);
|
||||
return Ok(());
|
||||
}
|
||||
}
|
||||
continue;
|
||||
}
|
||||
};
|
||||
if bytes_read == 0 {
|
||||
break;
|
||||
@@ -156,12 +154,10 @@ async fn run_once(
|
||||
let raw = data_lines.join("\n");
|
||||
data_lines.clear();
|
||||
|
||||
match serde_json::from_str::<EventEnvelope>(&raw) {
|
||||
Ok(event) => handle_event(app, event, phases.clone(), cooldowns.clone()).await,
|
||||
Err(err) => {
|
||||
warn!("[desktop:activity] Failed to parse SSE data: {err}; raw={raw}");
|
||||
}
|
||||
}
|
||||
match parse_event_envelope(&raw) {
|
||||
Ok((event, _directory)) => handle_event(app, event, phases.clone(), cooldowns.clone()).await,
|
||||
Err(err) => warn!("[desktop:activity] Failed to parse SSE data: {err}; raw={raw}"),
|
||||
};
|
||||
continue;
|
||||
}
|
||||
|
||||
@@ -173,6 +169,82 @@ async fn run_once(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn parse_event_envelope(raw: &str) -> Result<(EventEnvelope, Option<String>)> {
|
||||
if let Ok(event) = serde_json::from_str::<EventEnvelope>(raw) {
|
||||
return Ok((event, None));
|
||||
}
|
||||
|
||||
let multiplexed = serde_json::from_str::<MultiplexedEventEnvelope>(raw)?;
|
||||
Ok((multiplexed.payload, multiplexed.directory))
|
||||
}
|
||||
|
||||
async fn connect_activity_sse(
|
||||
runtime: &DesktopRuntime,
|
||||
client: &Client,
|
||||
base: &str,
|
||||
) -> Result<(reqwest::Response, SseScope)> {
|
||||
let opencode = runtime.opencode_manager();
|
||||
|
||||
let global_url = format!("{base}/global/event");
|
||||
match try_connect_sse(client, &global_url, "[desktop:activity]").await {
|
||||
Ok(response) => {
|
||||
debug!("[desktop:activity] Using SSE endpoint: {global_url}");
|
||||
return Ok((response, SseScope::Global));
|
||||
}
|
||||
Err(err) => {
|
||||
debug!(
|
||||
"[desktop:activity] SSE endpoint unavailable: {global_url} ({err:?}); falling back"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
let event_url = format!("{base}/event");
|
||||
match try_connect_sse(client, &event_url, "[desktop:activity]").await {
|
||||
Ok(response) => {
|
||||
debug!("[desktop:activity] Using SSE endpoint: {event_url}");
|
||||
return Ok((response, SseScope::Global));
|
||||
}
|
||||
Err(err) => {
|
||||
debug!(
|
||||
"[desktop:activity] SSE endpoint unavailable: {event_url} ({err:?}); falling back"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
let working_dir = opencode.get_working_directory();
|
||||
let directory = working_dir.to_string_lossy().to_string();
|
||||
let mut parsed = reqwest::Url::parse(&event_url)?;
|
||||
parsed.query_pairs_mut().append_pair("directory", &directory);
|
||||
let directory_url = parsed.to_string();
|
||||
|
||||
let response = try_connect_sse(client, &directory_url, "[desktop:activity]").await?;
|
||||
debug!("[desktop:activity] Using directory-scoped SSE endpoint: {directory_url}");
|
||||
Ok((response, SseScope::Directory(working_dir)))
|
||||
}
|
||||
|
||||
async fn try_connect_sse(client: &Client, url: &str, log_prefix: &str) -> Result<reqwest::Response> {
|
||||
debug!("{log_prefix} Connecting SSE: {url}");
|
||||
|
||||
let response = client
|
||||
.get(url)
|
||||
.header("accept", "text/event-stream")
|
||||
.header("accept-encoding", "identity")
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
debug!(
|
||||
"{log_prefix} SSE response status={} headers={:?}",
|
||||
response.status(),
|
||||
response.headers()
|
||||
);
|
||||
|
||||
if !response.status().is_success() {
|
||||
anyhow::bail!("SSE connect failed with status {}", response.status());
|
||||
}
|
||||
|
||||
Ok(response)
|
||||
}
|
||||
|
||||
async fn handle_event(
|
||||
app: &AppHandle,
|
||||
event: EventEnvelope,
|
||||
@@ -201,6 +273,16 @@ async fn handle_event(
|
||||
set_phase(app, &id, phase, phases.clone(), cooldowns.clone()).await;
|
||||
}
|
||||
}
|
||||
"session.idle" => {
|
||||
let session_id = event
|
||||
.properties
|
||||
.get("sessionID")
|
||||
.and_then(Value::as_str)
|
||||
.map(|s| s.to_string());
|
||||
if let Some(id) = session_id {
|
||||
set_phase(app, &id, ActivityPhase::Idle, phases.clone(), cooldowns.clone()).await;
|
||||
}
|
||||
}
|
||||
"message.updated" => {
|
||||
if let Some(info) = event.properties.get("info") {
|
||||
let role = info.get("role").and_then(Value::as_str).unwrap_or_default();
|
||||
@@ -219,37 +301,109 @@ async fn handle_event(
|
||||
.map(|s| s.to_string());
|
||||
|
||||
if let Some(id) = session_id {
|
||||
// If current phase is busy, move to cooldown for 2s then idle
|
||||
let current = { phases.lock().await.get(&id).cloned() };
|
||||
if matches!(current, Some(ActivityPhase::Busy)) {
|
||||
set_phase(app, &id, ActivityPhase::Cooldown, phases.clone(), cooldowns.clone()).await;
|
||||
|
||||
let app_clone = app.clone();
|
||||
let phases_clone = phases.clone();
|
||||
let cooldowns_clone = cooldowns.clone();
|
||||
let id_clone = id.clone();
|
||||
let handle = tauri::async_runtime::spawn(async move {
|
||||
tokio::time::sleep(Duration::from_secs(2)).await;
|
||||
let current = { phases_clone.lock().await.get(&id_clone).cloned() };
|
||||
if matches!(current, Some(ActivityPhase::Cooldown)) {
|
||||
set_phase(&app_clone, &id_clone, ActivityPhase::Idle, phases_clone, cooldowns_clone).await;
|
||||
}
|
||||
});
|
||||
|
||||
// Store cooldown handle to cancel if phase changes earlier
|
||||
let mut cd = cooldowns.lock().await;
|
||||
if let Some(prev) = cd.remove(&id) {
|
||||
prev.abort();
|
||||
}
|
||||
cd.insert(id, handle);
|
||||
}
|
||||
enter_cooldown_if_busy(app, &id, phases.clone(), cooldowns.clone()).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
"message.part.updated" => {
|
||||
let Some(info) = event.properties.get("info") else {
|
||||
return;
|
||||
};
|
||||
|
||||
let role = info.get("role").and_then(Value::as_str).unwrap_or_default();
|
||||
if role != "assistant" {
|
||||
return;
|
||||
}
|
||||
|
||||
let session_id = info
|
||||
.get("sessionID")
|
||||
.and_then(Value::as_str)
|
||||
.map(|s| s.to_string());
|
||||
|
||||
let Some(id) = session_id else {
|
||||
return;
|
||||
};
|
||||
|
||||
// Mark session busy when we see assistant parts streaming (covers cases where session.status is missing).
|
||||
if is_streaming_assistant_part(&event.properties) {
|
||||
set_phase(app, &id, ActivityPhase::Busy, phases.clone(), cooldowns.clone()).await;
|
||||
}
|
||||
|
||||
// Derive cooldown from "step-finish reason=stop" marker when present.
|
||||
if is_stop_step_finish_part(&event.properties) {
|
||||
enter_cooldown_if_busy(app, &id, phases.clone(), cooldowns.clone()).await;
|
||||
}
|
||||
}
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
|
||||
fn is_streaming_assistant_part(properties: &Value) -> bool {
|
||||
let Some(part) = properties.get("part") else {
|
||||
return false;
|
||||
};
|
||||
let part_type = part.get("type").and_then(Value::as_str).unwrap_or_default();
|
||||
matches!(
|
||||
part_type,
|
||||
"step-start" | "text" | "tool" | "reasoning" | "file" | "patch"
|
||||
)
|
||||
}
|
||||
|
||||
fn is_stop_step_finish_part(properties: &Value) -> bool {
|
||||
let Some(part) = properties.get("part") else {
|
||||
return false;
|
||||
};
|
||||
let part_type = part.get("type").and_then(Value::as_str);
|
||||
let reason = part.get("reason").and_then(Value::as_str);
|
||||
part_type == Some("step-finish") && reason == Some("stop")
|
||||
}
|
||||
|
||||
async fn enter_cooldown_if_busy(
|
||||
app: &AppHandle,
|
||||
session_id: &str,
|
||||
phases: Arc<Mutex<HashMap<String, ActivityPhase>>>,
|
||||
cooldowns: Arc<Mutex<HashMap<String, tauri::async_runtime::JoinHandle<()>>>>,
|
||||
) {
|
||||
let current = { phases.lock().await.get(session_id).cloned() };
|
||||
if !matches!(current, Some(ActivityPhase::Busy)) {
|
||||
return;
|
||||
}
|
||||
|
||||
set_phase(
|
||||
app,
|
||||
session_id,
|
||||
ActivityPhase::Cooldown,
|
||||
phases.clone(),
|
||||
cooldowns.clone(),
|
||||
)
|
||||
.await;
|
||||
|
||||
let app_clone = app.clone();
|
||||
let phases_clone = phases.clone();
|
||||
let cooldowns_clone = cooldowns.clone();
|
||||
let id_clone = session_id.to_string();
|
||||
let handle = tauri::async_runtime::spawn(async move {
|
||||
tokio::time::sleep(Duration::from_secs(2)).await;
|
||||
let current = { phases_clone.lock().await.get(&id_clone).cloned() };
|
||||
if matches!(current, Some(ActivityPhase::Cooldown)) {
|
||||
set_phase(
|
||||
&app_clone,
|
||||
&id_clone,
|
||||
ActivityPhase::Idle,
|
||||
phases_clone,
|
||||
cooldowns_clone,
|
||||
)
|
||||
.await;
|
||||
}
|
||||
});
|
||||
|
||||
let mut cd = cooldowns.lock().await;
|
||||
if let Some(prev) = cd.remove(session_id) {
|
||||
prev.abort();
|
||||
}
|
||||
cd.insert(session_id.to_string(), handle);
|
||||
}
|
||||
|
||||
async fn set_phase(
|
||||
app: &AppHandle,
|
||||
session_id: &str,
|
||||
|
||||
Reference in New Issue
Block a user