1use std::io::Write;
21use std::path::{Path, PathBuf};
22use std::time::{SystemTime, UNIX_EPOCH};
23use tracing::warn;
24
25pub fn events_path(project_root: &Path) -> PathBuf {
27 project_root.join(".devflow").join("events.jsonl")
28}
29
30const SCHEMA_VERSION: u32 = 1;
32
33pub fn emit(project_root: &Path, phase: u32, event: &str, fields: serde_json::Value) {
36 let mut line = serde_json::json!({
37 "v": SCHEMA_VERSION,
38 "ts": unix_now(),
39 "phase": phase,
40 "event": event,
41 });
42 match fields {
43 serde_json::Value::Object(map) => {
44 let base = line.as_object_mut().expect("line is an object");
45 for (key, value) in map {
46 base.entry(key).or_insert(value);
49 }
50 }
51 serde_json::Value::Null => {}
52 other => {
53 line["data"] = other;
54 }
55 }
56 let path = events_path(project_root);
57 if let Some(parent) = path.parent()
58 && let Err(err) = crate::workflow::ensure_devflow_dir(parent)
59 {
60 warn!("could not create events dir: {err}");
61 return;
62 }
63 let result = std::fs::OpenOptions::new()
64 .create(true)
65 .append(true)
66 .open(&path)
67 .and_then(|mut f| f.write_all(format!("{line}\n").as_bytes()));
68 if let Err(err) = result {
69 warn!("could not append to {}: {err}", path.display());
70 }
71}
72
73#[derive(Debug, Clone)]
79pub struct PhaseEventSummary {
80 pub event: serde_json::Value,
81 pub stage_launched_ts: Option<u64>,
82}
83
84pub fn last_events_by_phase(
91 project_root: &Path,
92) -> std::collections::HashMap<u32, PhaseEventSummary> {
93 use std::collections::hash_map::Entry;
94
95 let mut latest: std::collections::HashMap<u32, PhaseEventSummary> =
96 std::collections::HashMap::new();
97 let Ok(contents) = std::fs::read_to_string(events_path(project_root)) else {
98 return latest;
99 };
100 for event in contents
101 .lines()
102 .filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
103 {
104 let Some(phase) = event.get("phase").and_then(|p| p.as_u64()) else {
105 continue;
106 };
107 let phase = phase as u32;
108 let launch_ts = (event.get("event").and_then(|e| e.as_str()) == Some("stage_launched"))
109 .then(|| event.get("ts").and_then(|t| t.as_u64()))
110 .flatten();
111 match latest.entry(phase) {
115 Entry::Occupied(mut occupied) => {
116 let summary = occupied.get_mut();
117 summary.event = event;
118 if let Some(ts) = launch_ts {
119 summary.stage_launched_ts = Some(ts);
120 }
121 }
122 Entry::Vacant(vacant) => {
123 vacant.insert(PhaseEventSummary {
124 event,
125 stage_launched_ts: launch_ts,
126 });
127 }
128 }
129 }
130 latest
131}
132
133pub fn last_event_for_phase(project_root: &Path, phase: u32) -> Option<serde_json::Value> {
135 last_events_by_phase(project_root)
136 .remove(&phase)
137 .map(|summary| summary.event)
138}
139
140pub fn last_event_of_kind_for_phase(
151 project_root: &Path,
152 phase: u32,
153 event: &str,
154) -> Option<serde_json::Value> {
155 let contents = std::fs::read_to_string(events_path(project_root)).ok()?;
156 let mut last = None;
157 for value in contents
158 .lines()
159 .filter_map(|line| serde_json::from_str::<serde_json::Value>(line).ok())
160 {
161 if value.get("phase").and_then(|p| p.as_u64()) == Some(phase as u64)
162 && value.get("event").and_then(|e| e.as_str()) == Some(event)
163 {
164 last = Some(value);
165 }
166 }
167 last
168}
169
170pub fn has_event_for_phase(project_root: &Path, phase: u32, event: &str) -> bool {
174 last_event_of_kind_for_phase(project_root, phase, event).is_some()
175}
176
177pub fn describe(event: &serde_json::Value) -> String {
179 let kind = event
180 .get("event")
181 .and_then(|e| e.as_str())
182 .unwrap_or("unknown");
183 let detail = ["to", "stage", "status", "hook", "reason"]
184 .iter()
185 .find_map(|key| event.get(*key).and_then(|v| v.as_str()));
186 match detail {
187 Some(detail) => format!("{kind} ({detail})"),
188 None => kind.to_string(),
189 }
190}
191
192fn unix_now() -> u64 {
193 SystemTime::now()
194 .duration_since(UNIX_EPOCH)
195 .map(|d| d.as_secs())
196 .unwrap_or(0)
197}
198
199#[cfg(test)]
200mod tests {
201 use super::*;
202
203 fn read_lines(root: &Path) -> Vec<serde_json::Value> {
204 std::fs::read_to_string(events_path(root))
205 .unwrap_or_default()
206 .lines()
207 .map(|l| serde_json::from_str(l).expect("every line parses as JSON"))
208 .collect()
209 }
210
211 #[test]
212 fn emit_appends_parseable_lines_with_envelope_fields() {
213 let dir = tempfile::tempdir().unwrap();
214 emit(
215 dir.path(),
216 14,
217 "transition",
218 serde_json::json!({"from": "code", "to": "validate"}),
219 );
220 emit(
221 dir.path(),
222 15,
223 "gate_fired",
224 serde_json::json!({"stage": "ship"}),
225 );
226
227 let lines = read_lines(dir.path());
228 assert_eq!(lines.len(), 2);
229 for line in &lines {
230 assert_eq!(line["v"], 1);
231 assert!(line["ts"].as_u64().is_some());
232 assert!(line["phase"].as_u64().is_some());
233 assert!(line["event"].as_str().is_some());
234 }
235 assert_eq!(lines[0]["phase"], 14);
236 assert_eq!(lines[0]["from"], "code");
237 assert_eq!(lines[1]["phase"], 15);
238 assert_eq!(lines[1]["stage"], "ship");
239 }
240
241 #[test]
242 fn emit_never_lets_payload_forge_envelope_keys() {
243 let dir = tempfile::tempdir().unwrap();
244 emit(
245 dir.path(),
246 7,
247 "transition",
248 serde_json::json!({"phase": 99, "event": "forged", "note": "kept"}),
249 );
250
251 let lines = read_lines(dir.path());
252 assert_eq!(lines[0]["phase"], 7, "envelope phase must win");
253 assert_eq!(lines[0]["event"], "transition", "envelope event must win");
254 assert_eq!(lines[0]["note"], "kept");
255 }
256
257 #[test]
258 fn last_event_for_phase_filters_by_phase() {
259 let dir = tempfile::tempdir().unwrap();
260 emit(dir.path(), 1, "workflow_started", serde_json::Value::Null);
261 emit(dir.path(), 2, "workflow_started", serde_json::Value::Null);
262 emit(
263 dir.path(),
264 1,
265 "transition",
266 serde_json::json!({"to": "plan"}),
267 );
268
269 let last = last_event_for_phase(dir.path(), 1).expect("phase 1 events exist");
270 assert_eq!(last["event"], "transition");
271 let other = last_event_for_phase(dir.path(), 2).expect("phase 2 events exist");
272 assert_eq!(other["event"], "workflow_started");
273 assert!(last_event_for_phase(dir.path(), 3).is_none());
274 }
275
276 #[test]
278 fn last_events_by_phase_collects_latest_per_phase_in_one_pass() {
279 let dir = tempfile::tempdir().unwrap();
280 emit(dir.path(), 1, "workflow_started", serde_json::Value::Null);
281 emit(dir.path(), 2, "workflow_started", serde_json::Value::Null);
282 emit(
283 dir.path(),
284 1,
285 "transition",
286 serde_json::json!({"to": "plan"}),
287 );
288 emit(
289 dir.path(),
290 2,
291 "gate_fired",
292 serde_json::json!({"stage": "ship"}),
293 );
294
295 let latest = last_events_by_phase(dir.path());
296 assert_eq!(latest.len(), 2);
297 assert_eq!(latest[&1].event["event"], "transition");
298 assert_eq!(latest[&2].event["event"], "gate_fired");
299 assert!(last_events_by_phase(&dir.path().join("empty")).is_empty());
300 }
301
302 #[test]
307 fn last_events_by_phase_tracks_newest_stage_launched_ts_across_the_pass() {
308 let dir = tempfile::tempdir().unwrap();
309
310 emit(dir.path(), 1, "workflow_started", serde_json::Value::Null);
312 assert_eq!(last_events_by_phase(dir.path())[&1].stage_launched_ts, None);
313
314 emit(
316 dir.path(),
317 1,
318 "stage_launched",
319 serde_json::json!({"stage": "define"}),
320 );
321 let first_ts = last_events_by_phase(dir.path())[&1]
322 .stage_launched_ts
323 .expect("stage_launched recorded a timestamp");
324
325 emit(
328 dir.path(),
329 1,
330 "transition",
331 serde_json::json!({"to": "plan"}),
332 );
333 let after_transition = last_events_by_phase(dir.path());
334 assert_eq!(after_transition[&1].stage_launched_ts, Some(first_ts));
335 assert_eq!(after_transition[&1].event["event"], "transition");
336
337 emit(
339 dir.path(),
340 1,
341 "stage_launched",
342 serde_json::json!({"stage": "plan"}),
343 );
344 let path = events_path(dir.path());
345 let mut contents = std::fs::read_to_string(&path).unwrap();
346 contents.push_str("{truncated\n");
348 std::fs::write(&path, contents).unwrap();
349
350 let latest = last_events_by_phase(dir.path());
351 let newest_ts = latest[&1]
352 .stage_launched_ts
353 .expect("newest stage_launched timestamp present");
354 assert!(newest_ts >= first_ts, "newest launch must not be older");
355 assert_eq!(latest[&1].event["event"], "stage_launched");
356 }
357
358 #[test]
359 fn last_event_skips_corrupt_lines() {
360 let dir = tempfile::tempdir().unwrap();
361 emit(dir.path(), 4, "workflow_started", serde_json::Value::Null);
362 let path = events_path(dir.path());
363 let mut contents = std::fs::read_to_string(&path).unwrap();
364 contents.push_str("{truncated\n");
365 std::fs::write(&path, contents).unwrap();
366
367 let last = last_event_for_phase(dir.path(), 4).expect("valid line still found");
368 assert_eq!(last["event"], "workflow_started");
369 }
370
371 const MARKER_EVENT: &str = "marker_event";
376
377 #[test]
378 fn last_event_of_kind_for_phase_filters_by_phase_and_event_name() {
379 let dir = tempfile::tempdir().unwrap();
380 emit(dir.path(), 1, "workflow_finished", serde_json::Value::Null);
381 emit(
382 dir.path(),
383 1,
384 MARKER_EVENT,
385 serde_json::json!({"stage": "ship"}),
386 );
387 emit(
388 dir.path(),
389 2,
390 MARKER_EVENT,
391 serde_json::json!({"stage": "ship"}),
392 );
393
394 let phase1 = last_event_of_kind_for_phase(dir.path(), 1, MARKER_EVENT)
395 .expect("phase 1 marker event exists");
396 assert_eq!(phase1["stage"], "ship");
397 assert!(has_event_for_phase(dir.path(), 1, MARKER_EVENT));
398 assert!(!has_event_for_phase(dir.path(), 3, MARKER_EVENT));
399 assert!(last_event_of_kind_for_phase(dir.path(), 3, MARKER_EVENT).is_none());
400 }
401
402 #[test]
403 fn last_event_of_kind_for_phase_skips_corrupt_lines() {
404 let dir = tempfile::tempdir().unwrap();
405 emit(dir.path(), 4, MARKER_EVENT, serde_json::Value::Null);
406 let path = events_path(dir.path());
407 let mut contents = std::fs::read_to_string(&path).unwrap();
408 contents.push_str("{truncated\n");
409 std::fs::write(&path, contents).unwrap();
410
411 assert!(has_event_for_phase(dir.path(), 4, MARKER_EVENT));
412 }
413
414 #[test]
415 fn describe_prefers_detail_fields() {
416 assert_eq!(
417 describe(&serde_json::json!({"event": "transition", "to": "ship"})),
418 "transition (ship)"
419 );
420 assert_eq!(
421 describe(&serde_json::json!({"event": "workflow_finished"})),
422 "workflow_finished"
423 );
424 assert_eq!(describe(&serde_json::json!({})), "unknown");
425 }
426
427 #[test]
428 fn emit_is_fail_soft_on_unwritable_path() {
429 let dir = tempfile::tempdir().unwrap();
432 std::fs::write(dir.path().join(".devflow"), "not a dir").unwrap();
433 emit(dir.path(), 1, "transition", serde_json::Value::Null);
434 }
435}