use anyhow::{anyhow, Result};
use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use crate::durable::{Author, SteerComment, WorkRef};
use crate::lf::commands::runs::{collect_run_activity_since, RunSnapshot};
use crate::lf::commands::util::parse_since;
use crate::lf::commands::waves::PrMergeRequestSnapshot;
use crate::lf::commands::work_catalog::{WorkCatalog, WorkOwner};
use crate::lf::commands::WorkFilter;
use crate::store::sqlite::SqliteStore;
use crate::work::task::{GithubPr, TaskPr, TaskPrId};
const MAX_LIMIT: usize = 200;
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkActivitySnapshot {
pub generated_at: i64,
pub since: i64,
pub limit: usize,
pub truncated: bool,
pub items: Vec<WorkActivityEntry>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkActivityEntry {
pub id: String,
pub recorded_at: i64,
pub summary: String,
pub work: WorkRef,
pub subject: String,
pub fact: WorkActivityFact,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum WorkActivityFact {
WorkCreated,
RunStarted {
run_id: String,
},
RunFinished {
run_id: String,
status: String,
},
PrStarted {
id: TaskPrId,
},
PrPublishRequested {
id: TaskPrId,
github: Option<GithubPr>,
},
PrMergeRequested {
id: TaskPrId,
request: PrMergeRequestSnapshot,
github: Option<GithubPr>,
},
PrMerged {
id: TaskPrId,
github: Option<GithubPr>,
merge_commit: String,
},
PrAbandoned {
id: TaskPrId,
github: Option<GithubPr>,
},
SteerIssued {
id: i64,
author: Author,
},
}
pub fn run(
since_value: &str,
limit: usize,
wave: Option<&str>,
project: Option<&str>,
task: Option<&str>,
json: bool,
) -> Result<()> {
let generated_at = OffsetDateTime::now_utc();
let since = parse_since(since_value, generated_at)?.unix_timestamp();
if limit == 0 || limit > MAX_LIMIT {
return Err(anyhow!(
"--limit must be between 1 and {MAX_LIMIT}; got {}",
limit
));
}
let filter = WorkFilter {
wave,
project,
task,
};
let path = crate::store::database_path_from_env()?;
let snapshot = if path.exists() {
let store = SqliteStore::new(&path)?;
build_snapshot(&store, generated_at.unix_timestamp(), since, limit, filter)?
} else {
WorkActivitySnapshot {
generated_at: generated_at.unix_timestamp(),
since,
limit,
truncated: false,
items: Vec::new(),
}
};
if json {
println!("{}", serde_json::to_string(&snapshot)?);
} else {
print_snapshot(&snapshot);
}
Ok(())
}
fn build_snapshot(
store: &SqliteStore,
generated_at: i64,
since: i64,
limit: usize,
filter: WorkFilter<'_>,
) -> Result<WorkActivitySnapshot> {
let tasks = store
.list_tasks(None)
.map_err(|error| anyhow!("failed to read Tasks: {error}"))?;
let catalog = WorkCatalog::new(store.work_identities()?)?;
let mut entries = catalog.creation_entries(since);
for task in &tasks {
let work = catalog
.owners
.get(&WorkRef::Task(task.id.clone()))
.ok_or_else(|| anyhow!("Task {} is missing from the Work catalog", task.id))?;
for pr in store
.task_prs(&task.id)
.map_err(|error| anyhow!("failed to read PRs for {}: {error}", task.plan.identifier))?
{
entries.extend(pr_entries(&pr, work, since));
}
}
let runs = collect_run_activity_since(filter, since, &catalog)?;
for run in runs {
if let Some(work) = catalog.resolve_run(&run) {
entries.extend(run_entries(&run, work, since));
}
}
let steers = store
.steers_since(since)
.map_err(|error| anyhow!("failed to read Steers: {error}"))?;
for comment in steers {
if let Some(work) = catalog.owners.get(&comment.work) {
entries.push(steer_entry(&comment, work));
}
}
Ok(finalize_snapshot(
entries,
generated_at,
since,
limit,
&catalog,
filter,
))
}
impl WorkCatalog {
fn creation_entries(&self, since: i64) -> Vec<WorkActivityEntry> {
self.owners
.values()
.filter_map(|owner| {
let recorded_at = owner.created_at.filter(|value| *value >= since)?;
let kind = match &owner.work {
WorkRef::Wave(_) => "Wave",
WorkRef::Project(_) => "Project",
WorkRef::Task(_) => "Task",
};
Some(activity_entry(
format!("{}:{}:created", owner.work.kind(), owner.work.id()),
recorded_at,
format!("{kind} {} created", owner.subject),
owner,
WorkActivityFact::WorkCreated,
))
})
.collect()
}
}
fn activity_entry(
id: String,
recorded_at: i64,
summary: String,
owner: &WorkOwner,
fact: WorkActivityFact,
) -> WorkActivityEntry {
WorkActivityEntry {
id,
recorded_at,
summary,
work: owner.work.clone(),
subject: owner.subject.clone(),
fact,
}
}
fn run_entries(run: &RunSnapshot, work: &WorkOwner, since: i64) -> Vec<WorkActivityEntry> {
let label = run.label();
let mut entries = Vec::new();
if run.started >= since {
entries.push(activity_entry(
format!("run:{}:started", run.id),
run.started,
format!("{label} started"),
work,
WorkActivityFact::RunStarted {
run_id: run.id.clone(),
},
));
}
if let Some(ended) = run.ended.filter(|ended| *ended >= since) {
entries.push(activity_entry(
format!("run:{}:finished", run.id),
ended,
format!("{label} finished {}", run.status()),
work,
WorkActivityFact::RunFinished {
run_id: run.id.clone(),
status: run.status().to_string(),
},
));
}
entries
}
fn pr_entries(pr: &TaskPr, work: &WorkOwner, since: i64) -> Vec<WorkActivityEntry> {
let github = pr.github();
let reference = github
.map(|github| format!("PR #{}", github.number))
.unwrap_or_else(|| format!("PR {}", pr.slug));
let mut entries = Vec::new();
if pr.created_at.unix_timestamp() >= since {
entries.push(activity_entry(
format!("pr:{}:started", pr.id),
pr.created_at.unix_timestamp(),
format!("PR work started on {}", pr.branch),
work,
WorkActivityFact::PrStarted { id: pr.id.clone() },
));
}
if let Some(publication) = pr
.publication
.as_ref()
.filter(|publication| publication.requested_at.unix_timestamp() >= since)
{
entries.push(activity_entry(
format!("pr:{}:publish_requested", pr.id),
publication.requested_at.unix_timestamp(),
format!("Publication requested for {reference}"),
work,
WorkActivityFact::PrPublishRequested {
id: pr.id.clone(),
github: publication.github.clone(),
},
));
}
if let Some(request) = pr
.publication
.as_ref()
.and_then(|publication| publication.merge.as_ref())
.filter(|request| request.requested_at.unix_timestamp() >= since)
{
entries.push(activity_entry(
format!("pr:{}:merge_requested", pr.id),
request.requested_at.unix_timestamp(),
format!("Merge requested for {reference}"),
work,
WorkActivityFact::PrMergeRequested {
id: pr.id.clone(),
request: PrMergeRequestSnapshot::from(request),
github: github.cloned(),
},
));
}
if let Some(merge_commit) = pr
.merge_commit
.as_ref()
.filter(|_| pr.updated_at.unix_timestamp() >= since)
{
entries.push(activity_entry(
format!("pr:{}:merged", pr.id),
pr.updated_at.unix_timestamp(),
format!("{reference} merged"),
work,
WorkActivityFact::PrMerged {
id: pr.id.clone(),
github: github.cloned(),
merge_commit: merge_commit.clone(),
},
));
}
if let Some(abandoned_at) = pr
.abandoned_at
.filter(|abandoned_at| abandoned_at.unix_timestamp() >= since)
{
entries.push(activity_entry(
format!("pr:{}:abandoned", pr.id),
abandoned_at.unix_timestamp(),
format!("{reference} abandoned"),
work,
WorkActivityFact::PrAbandoned {
id: pr.id.clone(),
github: github.cloned(),
},
));
}
entries
}
fn steer_entry(comment: &SteerComment, work: &WorkOwner) -> WorkActivityEntry {
activity_entry(
format!(
"steer:{}:{}:{}",
comment.work.kind(),
comment.work.id(),
comment.steer.id
),
comment.issued_at.unix_timestamp(),
format!("Steered: {}", comment.steer.text),
work,
WorkActivityFact::SteerIssued {
id: comment.steer.id,
author: comment.steer.author.clone(),
},
)
}
fn finalize_snapshot(
mut entries: Vec<WorkActivityEntry>,
generated_at: i64,
since: i64,
limit: usize,
catalog: &WorkCatalog,
filter: WorkFilter<'_>,
) -> WorkActivitySnapshot {
entries.retain(|entry| {
catalog
.owners
.get(&entry.work)
.is_some_and(|owner| owner.matches(filter))
});
entries.sort_by(|left, right| {
right
.recorded_at
.cmp(&left.recorded_at)
.then_with(|| right.id.cmp(&left.id))
});
let truncated = entries.len() > limit;
entries.truncate(limit);
WorkActivitySnapshot {
generated_at,
since,
limit,
truncated,
items: entries,
}
}
fn print_snapshot(snapshot: &WorkActivitySnapshot) {
if snapshot.items.is_empty() {
println!("No durable Work activity recorded in this window.");
return;
}
println!("TIME WORK ACTIVITY");
for item in &snapshot.items {
let time = OffsetDateTime::from_unix_timestamp(item.recorded_at)
.map(|value| {
format!(
"{:04}-{:02}-{:02} {:02}:{:02}",
value.year(),
u8::from(value.month()),
value.day(),
value.hour(),
value.minute()
)
})
.unwrap_or_else(|_| item.recorded_at.to_string());
println!(
"{time:<17} {work:<16} {summary}",
work = item.subject,
summary = item.summary
);
}
if snapshot.truncated {
println!("… more activity exists; increase --limit to inspect it.");
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::work::task::{
AfterMerge, GithubPr, PrMergeMode, PrMergeRequest, PrPresentation, PrPublication, TaskId,
TaskPrId,
};
#[test]
fn work_activity_fixture_round_trips() {
let fixture = include_str!("../../../../../tests/fixtures/dto/work_activity_snapshot.json");
let snapshot = serde_json::from_str::<WorkActivitySnapshot>(fixture).unwrap();
assert_eq!(snapshot.limit, 50);
assert_eq!(snapshot.items.len(), 5);
assert!(matches!(
snapshot.items[0].fact,
WorkActivityFact::RunFinished { ref status, .. } if status == "ok"
));
assert!(matches!(
snapshot.items[1].fact,
WorkActivityFact::PrMerged { ref github, .. }
if github.as_ref().map(|pr| pr.number) == Some(1144)
));
assert!(matches!(
snapshot.items[2].fact,
WorkActivityFact::PrMergeRequested { ref request, .. }
if request.requested_at == "2026-07-21T18:35:51Z"
));
assert_eq!(
serde_json::from_str::<WorkActivitySnapshot>(
&serde_json::to_string(&snapshot).unwrap()
)
.unwrap(),
snapshot
);
}
fn task_owner(identifier: &str, project: &str, wave: &str) -> WorkOwner {
WorkOwner {
work: WorkRef::Task(TaskId::new()),
subject: identifier.to_string(),
selectors: vec![
format!("wave:{wave}"),
format!("project:{project}"),
format!("task:{identifier}"),
],
created_at: None,
}
}
#[test]
fn steer_activity_ids_include_their_work_identity() {
let first = task_owner("W2-1", "control", "live");
let second = task_owner("W2-2", "control", "live");
let comment = |work: WorkRef| SteerComment {
work,
steer: crate::durable::Steer {
id: 1,
author: Author::User,
text: "direction".to_string(),
},
issued_at: OffsetDateTime::UNIX_EPOCH,
};
let first_entry = steer_entry(&comment(first.work.clone()), &first);
let second_entry = steer_entry(&comment(second.work.clone()), &second);
assert_ne!(first_entry.id, second_entry.id);
assert!(first_entry.id.contains(first.work.id()));
assert!(second_entry.id.contains(second.work.id()));
}
fn catalog(owners: &[WorkOwner]) -> WorkCatalog {
WorkCatalog {
owners: owners
.iter()
.map(|owner| (owner.work.clone(), owner.clone()))
.collect(),
}
}
fn entry(id: &str, recorded_at: i64, owner: &WorkOwner) -> WorkActivityEntry {
activity_entry(
id.to_string(),
recorded_at,
id.to_string(),
owner,
WorkActivityFact::WorkCreated,
)
}
#[test]
fn filters_apply_before_ordering_and_cap() {
let other = task_owner("W2-2", "other", "live");
let target = task_owner("W2-1", "control", "live");
let catalog = catalog(&[other.clone(), target.clone()]);
let entries = vec![
entry("other-new", 30, &other),
entry("target-old", 10, &target),
entry("target-new", 20, &target),
];
let snapshot = finalize_snapshot(
entries,
40,
0,
1,
&catalog,
WorkFilter {
wave: Some("live"),
project: Some("control"),
task: Some("W2-1"),
},
);
assert!(snapshot.truncated);
assert_eq!(snapshot.items[0].id, "target-new");
}
#[test]
fn ordering_has_a_stable_tie_breaker() {
let owner = task_owner("W2-1", "control", "live");
let catalog = catalog(std::slice::from_ref(&owner));
let snapshot = finalize_snapshot(
vec![entry("a", 10, &owner), entry("b", 10, &owner)],
20,
0,
50,
&catalog,
WorkFilter::default(),
);
assert_eq!(
snapshot
.items
.iter()
.map(|item| item.id.as_str())
.collect::<Vec<_>>(),
["b", "a"]
);
}
#[test]
fn conflicting_filters_return_no_activity() {
let owner = task_owner("W2-1", "control", "live");
let catalog = catalog(std::slice::from_ref(&owner));
let snapshot = finalize_snapshot(
vec![entry("target", 10, &owner)],
20,
0,
50,
&catalog,
WorkFilter {
wave: Some("another-wave"),
project: Some("control"),
task: Some("W2-1"),
},
);
assert!(snapshot.items.is_empty());
assert!(!snapshot.truncated);
}
#[test]
fn run_start_and_finish_name_the_same_generic_run() {
let run = RunSnapshot {
id: "run_00000000000000000000000000000001".to_string(),
parent_run_id: None,
repo: Some("/repo".to_string()),
worktree: Some("/repo.task".to_string()),
subjects: vec![crate::run_record::SubjectAttribution::declared(
"task:W2-1".to_string(),
)],
skill: Some("implement".to_string()),
outcome: Some("completed".to_string()),
started: 10,
ended: Some(20),
usage: crate::run_record::RunUsage::empty(),
evidence_gaps: 0,
harness: "codex".to_string(),
model: None,
surface: "cli".to_string(),
};
let entries = run_entries(&run, &task_owner("W2-1", "control", "live"), 0);
assert!(matches!(
entries[0].fact,
WorkActivityFact::RunStarted { .. }
));
assert!(matches!(
&entries[1].fact,
WorkActivityFact::RunFinished {
run_id,
status,
} if run_id == "run_00000000000000000000000000000001"
&& status == "completed"
));
}
#[test]
fn wire_reuses_work_reference_and_one_fact_tag() {
let owner = task_owner("W2-1", "control", "live");
let value = serde_json::to_value(entry("created", 10, &owner)).unwrap();
assert_eq!(value["work"]["kind"], "task");
assert_eq!(value["work"]["id"], owner.work.id());
assert_eq!(value["subject"], "W2-1");
assert_eq!(value["fact"]["kind"], "work_created");
assert!(value.get("evidence").is_none());
assert!(value.get("kind").is_none());
}
#[test]
fn serial_pr_stages_preserve_their_distinct_evidence() {
let pr = TaskPr {
id: TaskPrId::new(),
task_id: TaskId::new(),
sequence: 2,
slug: "activity-query".to_string(),
branch: "jack/activity-query".to_string(),
base_commit: "base".to_string(),
parent_pr_id: None,
publication: Some(PrPublication {
requested_at: OffsetDateTime::from_unix_timestamp(20).unwrap(),
presentation: Some(PrPresentation {
title: "Activity query".to_string(),
body: "Reviewer context".to_string(),
head_sha: "head".to_string(),
}),
github: Some(GithubPr {
number: 140,
url: "https://github.com/loopflowstudio/loopflow/pull/140".to_string(),
head_sha: Some("head".to_string()),
}),
merge: Some(PrMergeRequest {
mode: PrMergeMode::Auto,
requested_at: OffsetDateTime::from_unix_timestamp(30).unwrap(),
head_sha: "head".to_string(),
after_merge: AfterMerge::ContinueTask,
next_slug: None,
}),
}),
merge_commit: Some("merge".to_string()),
abandoned_at: None,
ci_observation: None,
github_observation: None,
linear_attachment_id: None,
linear_comment_id: None,
linear_link_error: None,
created_at: OffsetDateTime::from_unix_timestamp(10).unwrap(),
updated_at: OffsetDateTime::from_unix_timestamp(40).unwrap(),
};
let entries = pr_entries(&pr, &task_owner("W2-1", "control", "live"), 0);
assert!(matches!(
entries[0].fact,
WorkActivityFact::PrStarted { .. }
));
assert!(matches!(
entries[1].fact,
WorkActivityFact::PrPublishRequested { .. }
));
assert!(matches!(
entries[2].fact,
WorkActivityFact::PrMergeRequested { .. }
));
let merge_request = serde_json::to_value(&entries[2]).unwrap();
assert_eq!(
merge_request["fact"]["request"]["requested_at"],
"1970-01-01T00:00:30Z"
);
assert!(matches!(
&entries[3].fact,
WorkActivityFact::PrMerged {
github: Some(github),
merge_commit,
..
} if github.url.ends_with("/140") && merge_commit == "merge"
));
assert!(entries
.iter()
.all(|entry| entry.id.contains(pr.id.as_str())));
}
}