Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
116 changes: 113 additions & 3 deletions crates/spacetop-core/src/index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -157,9 +157,39 @@ impl WorkflowIndex {
.collect();
}

pub fn clear_entity_activities(&mut self) {
self.entity_activities.clear();
self.session_scan_error = None;
/// Carry the last published session activity through a workflow reload,
/// but only for active entities whose stable id and source path still
/// identify the same workflow item.
pub fn retain_session_activity_from(&mut self, previous: &Self) {
let previous_active_paths: HashMap<&str, &Path> = previous
.active
.iter()
.map(|entity| (entity.id.as_str(), entity.path.as_path()))
.collect();
let mut has_matching_identity = false;
let retained = self
.active
.iter()
.filter_map(|entity| {
let previous_path = previous_active_paths.get(entity.id.as_str())?;
if *previous_path != entity.path.as_path() {
return None;
}
has_matching_identity = true;
previous
.entity_activities
.get(&entity.id)
.cloned()
.map(|activity| (entity.id.clone(), activity))
})
.collect();

self.entity_activities = retained;
self.session_scan_error = if has_matching_identity {
previous.session_scan_error.clone()
} else {
None
};
}

pub fn set_session_scan_error(&mut self, message: String) {
Expand Down Expand Up @@ -666,6 +696,86 @@ mod tests {
);
}

#[test]
fn session_activity_transfer_preserves_matching_active_entity_identity() {
let mut previous = index();
previous.replace_session_scan_report(crate::domain::SessionScanReport {
workflow_dir: PathBuf::from("/tmp/workflow"),
repo_root: PathBuf::from("/tmp"),
scanned_roots: Vec::new(),
errors: vec!["one session record was malformed".to_string()],
attributions: vec![
crate::domain::EntityActivityAttribution {
entity_id: "010".to_string(),
activity: crate::domain::EntityActivity::Running {
handler: crate::domain::ActivityHandler::Worker,
runtime: crate::domain::AgentRuntime::Codex,
session_id: "session-010".to_string(),
updated_unix: 1_718_000_000,
},
},
crate::domain::EntityActivityAttribution {
entity_id: "002".to_string(),
activity: crate::domain::EntityActivity::HumanGate {
runtime: crate::domain::AgentRuntime::ClaudeCode,
session_id: "session-002".to_string(),
updated_unix: 1_718_000_001,
},
},
],
});
let mut reloaded = index();

reloaded.retain_session_activity_from(&previous);

assert_eq!(
reloaded
.entity_activity_for_entity_id("010")
.and_then(crate::domain::EntityActivity::session_id),
Some("session-010")
);
assert!(matches!(
reloaded.entity_activity_for_entity_id("002"),
Some(crate::domain::EntityActivity::HumanGate { .. })
));
assert_eq!(
reloaded.session_scan_error(),
Some("one session record was malformed")
);
}

#[test]
fn session_activity_transfer_does_not_cross_changed_entity_identity() {
let mut previous = index();
previous.replace_session_scan_report(crate::domain::SessionScanReport {
workflow_dir: PathBuf::from("/tmp/workflow"),
repo_root: PathBuf::from("/tmp"),
scanned_roots: Vec::new(),
errors: vec!["old scanner warning".to_string()],
attributions: vec![crate::domain::EntityActivityAttribution {
entity_id: "010".to_string(),
activity: crate::domain::EntityActivity::Running {
handler: crate::domain::ActivityHandler::Worker,
runtime: crate::domain::AgentRuntime::Codex,
session_id: "session-010".to_string(),
updated_unix: 1_718_000_000,
},
}],
});
let mut reused_id = entity("010", "Reused id", "plan");
reused_id.path = PathBuf::from("renamed-010.md");
let mut reloaded = index();
reloaded.active = vec![reused_id, entity("003", "New entity", "plan")];
reloaded.rebuild_lookup_maps();

reloaded.retain_session_activity_from(&previous);

assert!(!reloaded.entity_has_current_activity("010"));
assert!(!reloaded.entity_has_current_activity("002"));
assert!(!reloaded.entity_has_current_activity("003"));
assert_eq!(reloaded.session_scan_error(), None);
}

#[test]
fn related_returns_issue_pr_and_feedback_relations() {
let mut index = index();
Expand Down
4 changes: 2 additions & 2 deletions crates/spacetop/src/app/overview.rs
Original file line number Diff line number Diff line change
Expand Up @@ -370,13 +370,13 @@ impl OverviewState {
self.reload_from_index(index);
}

pub fn reload_from_index(&mut self, index: WorkflowIndex) {
pub fn reload_from_index(&mut self, mut index: WorkflowIndex) {
let prior_slug = self
.selected_item()
.and_then(|entity| entity_slug(&entity.path));

index.retain_session_activity_from(&self.index);
self.index = index;
self.index.clear_entity_activities();
// Invalidate archive view — a watcher-driven reload may have touched
// `_archive/` too. Dropping the cached list forces a rescan the next
// time the user toggles to archived scope.
Expand Down
37 changes: 37 additions & 0 deletions crates/spacetop/src/app/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2508,6 +2508,43 @@ fn session_activity_scan_failure_is_non_fatal_and_preserves_last_snapshot() {
);
}

#[test]
fn workflow_reload_preserves_running_activity_for_same_entity_identity() {
let workflow_dir = PathBuf::from("/tmp/spacetop-session-activity/workflow");
let snapshot = snapshot_with_items(1);
let mut app = App::from_snapshot(workflow_dir.clone(), snapshot.clone());
let repo_root = app.repo_root().expect("repo root").to_path_buf();

for report in [
session_report(&workflow_dir, &repo_root, "000"),
session_report(&workflow_dir, &repo_root, "000"),
] {
app.apply_session_activity_result(SessionActivityWorkerResult {
workflow_dir: workflow_dir.clone(),
repo_root: repo_root.clone(),
result: Ok(report),
state: Default::default(),
retry_immediately: false,
});
assert!(
app.as_overview()
.expect("overview")
.index()
.entity_has_current_activity("000"),
"each successful periodic report should publish running"
);

app.reload_from_snapshot(snapshot.clone());
assert!(
app.as_overview()
.expect("overview")
.index()
.entity_has_current_activity("000"),
"an unchanged workflow reload must retain the last good running attribution"
);
}
}

fn session_report(workflow_dir: &Path, repo_root: &Path, entity_id: &str) -> SessionScanReport {
SessionScanReport {
workflow_dir: workflow_dir.to_path_buf(),
Expand Down
20 changes: 12 additions & 8 deletions crates/spacetop/src/ui/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,17 @@ mod worktree;

fn app_with_items(items: Vec<Entity>) -> App {
let root = PathBuf::from("/tmp/spacetop-test");
let snapshot = WorkflowSnapshot {
let mut app = App::from_snapshot(root, snapshot_with_items(items));
app.handle_key(crossterm::event::KeyEvent::new(
crossterm::event::KeyCode::Enter,
crossterm::event::KeyModifiers::NONE,
));
app
}

fn snapshot_with_items(items: Vec<Entity>) -> WorkflowSnapshot {
let root = PathBuf::from("/tmp/spacetop-test");
WorkflowSnapshot {
definition: WorkflowDefinition {
root: root.clone(),
state: None,
Expand All @@ -49,13 +59,7 @@ fn app_with_items(items: Vec<Entity>) -> App {
},
items,
parse_errors: Vec::new(),
};
let mut app = App::from_snapshot(root, snapshot);
app.handle_key(crossterm::event::KeyEvent::new(
crossterm::event::KeyCode::Enter,
crossterm::event::KeyModifiers::NONE,
));
app
}
}

fn app_with_storage(storage: spacetop_core::domain::WorkflowStorage) -> App {
Expand Down
43 changes: 42 additions & 1 deletion crates/spacetop/src/ui/tests/task_list.rs
Original file line number Diff line number Diff line change
Expand Up @@ -309,6 +309,47 @@ fn task_row_renders_scanner_replay_then_clears_on_terminal_report() {
}));
assert!(buffer_text(running_buffer).contains("running · worker"));

let mut scan_state = running.state;
for cycle in 1..=5 {
app.reload_from_snapshot(snapshot_with_items(vec![active.clone()]));
terminal.draw(|frame| render(frame, &app)).expect("render");
let reloaded_buffer = terminal.backend().buffer();
assert!(
find_styled_text(reloaded_buffer, "\u{25CF}", |style| {
style.fg == Some(Color::Green)
}),
"cycle {cycle}: reload must preserve the running marker"
);
assert!(buffer_text(reloaded_buffer).contains("running · worker"));

let unchanged = scan_local_sessions_with_state(
&SessionScanRequest {
previous_state: scan_state,
..request.clone()
},
&StdProcessProbe,
SystemTime::now(),
)
.expect("unchanged scan");
scan_state = unchanged.state.clone();
app.apply_session_activity_result(crate::app::SessionActivityWorkerResult {
workflow_dir: workflow_dir.clone(),
repo_root: repo_root.clone(),
result: Ok(unchanged.report),
state: unchanged.state,
retry_immediately: false,
});
terminal.draw(|frame| render(frame, &app)).expect("render");
let unchanged_buffer = terminal.backend().buffer();
assert!(
find_styled_text(unchanged_buffer, "\u{25CF}", |style| {
style.fg == Some(Color::Green)
}),
"cycle {cycle}: unchanged scan must keep the running marker"
);
assert!(buffer_text(unchanged_buffer).contains("running · worker"));
}

let mut append = fs::OpenOptions::new()
.append(true)
.open(&child_path)
Expand All @@ -320,7 +361,7 @@ fn task_row_renders_scanner_replay_then_clears_on_terminal_report() {
.expect("terminal");
let stopped = scan_local_sessions_with_state(
&SessionScanRequest {
previous_state: running.state,
previous_state: scan_state,
..request
},
&StdProcessProbe,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -72,6 +72,13 @@ prompt and transcript bodies are not exposed in the UI. JSONL artifacts are
streamed record by record without a size cutoff and projected directly into
typed, privacy-safe facts.

A successful workflow snapshot reload preserves the last published activity
and scanner diagnostic only for active entities whose id and source path both
match the prior snapshot. Removed entities, new entities, and reused ids at a
different path inherit no prior attribution. A workflow reload is not a
lifecycle event; only the correlated structured session evidence below can
start, change, or stop activity.

`SessionScanState` crosses the background/app boundary. It contains
`SessionFileCursor` values and a `SessionEvidenceStore` keyed by stable runtime
session identity (falling back to a typed source identity until a session id is
Expand Down
Loading