1use std::collections::BTreeMap;
2
3use anyhow::Result;
4use rusqlite::{params, Connection, OptionalExtension};
5
6mod evidence;
7mod export;
8mod list;
9mod registry;
10mod trace_store;
11
12pub(crate) use export::{
13 load_export_eligible_procedure, procedure_export_slug, render_procedure_export,
14 ProcedureExportFormat, ProcedureExportSource, PROCEDURE_EXPORT_DRAFT_MARKER,
15};
16pub use list::{list_promoted_procedures, ProcedureListItem};
17pub(crate) use registry::{
18 ensure_existing_export_registry_match, load_procedure_export_doctor_report,
19 procedure_export_registry_exists, record_procedure_export, ProcedureExportRecordRequest,
20};
21
22#[cfg(test)]
23mod incremental_tests;
24
25const DEFAULT_MIN_VERIFIED_RUNS: usize = 2;
26const DEFAULT_MAX_VERIFICATION_AGE_SECS: i64 = 14 * 24 * 60 * 60;
27
28#[derive(Debug, Clone, PartialEq, Eq)]
29pub struct ProcedureTrace {
30 pub project: String,
31 pub branch: Option<String>,
32 pub workflow_key: String,
33 pub command: String,
34 pub files_touched: Vec<String>,
35 pub succeeded: bool,
36 pub verified_at_epoch: i64,
37 pub source_event_id: Option<i64>,
38}
39
40#[derive(Debug, Clone, PartialEq, Eq)]
41pub struct ProcedurePromotionPolicy {
42 pub min_verified_runs: usize,
43 pub max_verification_age_secs: i64,
44}
45
46impl Default for ProcedurePromotionPolicy {
47 fn default() -> Self {
48 Self {
49 min_verified_runs: DEFAULT_MIN_VERIFIED_RUNS,
50 max_verification_age_secs: DEFAULT_MAX_VERIFICATION_AGE_SECS,
51 }
52 }
53}
54
55#[derive(Debug, Clone, PartialEq)]
56pub struct ProcedureCandidate {
57 pub project: String,
58 pub branch: Option<String>,
59 pub workflow_key: String,
60 pub topic_key: String,
61 pub title: String,
62 pub content: String,
63 pub files: Vec<String>,
64 pub source_event_ids: Vec<i64>,
65 pub verified_runs: usize,
66 pub confidence: f64,
67 pub verified_at_epoch: i64,
68}
69
70pub fn build_procedure_candidate(
71 traces: &[ProcedureTrace],
72 now_epoch: i64,
73 policy: &ProcedurePromotionPolicy,
74) -> Option<ProcedureCandidate> {
75 let mut verified: Vec<&ProcedureTrace> = traces
76 .iter()
77 .filter(|trace| trace.succeeded)
78 .filter(|trace| trace.source_event_id.is_some())
79 .filter(|trace| {
80 now_epoch.saturating_sub(trace.verified_at_epoch) <= policy.max_verification_age_secs
81 })
82 .collect();
83 verified.sort_by_key(|trace| trace.verified_at_epoch);
84 if verified.len() < policy.min_verified_runs {
85 return None;
86 }
87
88 let first = verified[0];
89 if verified.iter().any(|trace| {
90 trace.project != first.project
91 || trace.branch != first.branch
92 || trace.workflow_key != first.workflow_key
93 || trace.command != first.command
94 }) {
95 return None;
96 }
97
98 let mut source_event_ids: Vec<i64> = verified
99 .iter()
100 .filter_map(|trace| trace.source_event_id)
101 .collect();
102 source_event_ids.sort_unstable();
103 source_event_ids.dedup();
104 if source_event_ids.len() < policy.min_verified_runs {
105 return None;
106 }
107
108 let mut files = verified
109 .iter()
110 .flat_map(|trace| trace.files_touched.iter().cloned())
111 .collect::<Vec<_>>();
112 files.sort();
113 files.dedup();
114
115 let verified_at_epoch = verified
116 .iter()
117 .map(|trace| trace.verified_at_epoch)
118 .max()
119 .unwrap_or(now_epoch);
120 let topic_key = procedure_topic_key(first);
121 let confidence = confidence_for_verified_runs(source_event_ids.len());
122 let content = render_procedure_content(
123 first,
124 &files,
125 &source_event_ids,
126 verified.len(),
127 verified_at_epoch,
128 );
129
130 Some(ProcedureCandidate {
131 project: first.project.clone(),
132 branch: first.branch.clone(),
133 workflow_key: first.workflow_key.clone(),
134 title: format!("Procedure: {}", first.workflow_key),
135 topic_key,
136 content,
137 files,
138 source_event_ids,
139 verified_runs: verified.len(),
140 confidence,
141 verified_at_epoch,
142 })
143}
144
145fn confidence_for_verified_runs(verified_runs: usize) -> f64 {
146 (0.7 + (verified_runs as f64 * 0.08)).min(0.95)
147}
148
149pub fn promote_procedure_memory(conn: &Connection, candidate: &ProcedureCandidate) -> Result<i64> {
150 let tx = conn.unchecked_transaction()?;
151 let files_json = (!candidate.files.is_empty())
152 .then(|| serde_json::to_string(&candidate.files))
153 .transpose()?;
154 let source_events_json = serde_json::to_string(&candidate.source_event_ids)?;
155 let memory_id = crate::memory::insert_memory_full(
156 &tx,
157 None,
158 &candidate.project,
159 Some(&candidate.topic_key),
160 &candidate.title,
161 &candidate.content,
162 "procedure",
163 files_json.as_deref(),
164 candidate.branch.as_deref(),
165 "project",
166 Some(candidate.verified_at_epoch),
167 )?;
168 tx.execute(
169 "UPDATE memories
170 SET evidence_event_ids = ?1,
171 confidence = ?2
172 WHERE id = ?3",
173 params![source_events_json, candidate.confidence, memory_id],
174 )?;
175 tx.commit()?;
176 Ok(memory_id)
177}
178
179pub(crate) fn promote_verified_procedures_for_task(
180 conn: &Connection,
181 task: &crate::db::ExtractionTask,
182 policy: &ProcedurePromotionPolicy,
183) -> Result<usize> {
184 let now_epoch = chrono::Utc::now().timestamp();
185 let traces = trace_store::load_verified_procedure_traces(conn, task, policy, now_epoch)?;
186 let mut groups: BTreeMap<(String, Option<String>, String, String), Vec<ProcedureTrace>> =
187 BTreeMap::new();
188 for trace in traces {
189 groups
190 .entry((
191 trace.project.clone(),
192 trace.branch.clone(),
193 trace.workflow_key.clone(),
194 trace.command.clone(),
195 ))
196 .or_default()
197 .push(trace);
198 }
199
200 let mut promoted = 0usize;
201 for traces in groups.into_values() {
202 let Some(candidate) = build_procedure_candidate(&traces, now_epoch, policy) else {
203 continue;
204 };
205 let existed = procedure_memory_exists(conn, &candidate.project, &candidate.topic_key)?;
206 promote_procedure_memory(conn, &candidate)?;
207 if !existed {
208 promoted += 1;
209 }
210 }
211 Ok(promoted)
212}
213
214fn procedure_topic_key(trace: &ProcedureTrace) -> String {
215 crate::memory::slugify_for_topic(
216 &format!(
217 "procedure {} branch {} command {}",
218 trace.workflow_key,
219 trace.branch.as_deref().unwrap_or("no-branch"),
220 trace.command
221 ),
222 96,
223 )
224}
225
226fn procedure_memory_exists(conn: &Connection, project: &str, topic_key: &str) -> Result<bool> {
227 let existing: Option<i64> = conn
228 .query_row(
229 "SELECT id FROM memories
230 WHERE project = ?1
231 AND topic_key = ?2
232 AND scope = 'project'
233 LIMIT 1",
234 params![project, topic_key],
235 |row| row.get(0),
236 )
237 .optional()?;
238 Ok(existing.is_some())
239}
240
241fn render_procedure_content(
242 trace: &ProcedureTrace,
243 files: &[String],
244 source_event_ids: &[i64],
245 verified_runs: usize,
246 verified_at_epoch: i64,
247) -> String {
248 let files_line = if files.is_empty() {
249 "Files: none recorded".to_string()
250 } else {
251 format!("Files: {}", files.join(", "))
252 };
253 format!(
254 "Procedure: {}\nCommand: {}\n{}\nVerified runs: {}\nVerified at: {}\nSource events: {}\nReuse when: the same project and branch need this verified workflow.",
255 trace.workflow_key,
256 trace.command,
257 files_line,
258 verified_runs,
259 verified_at_epoch,
260 source_event_ids
261 .iter()
262 .map(i64::to_string)
263 .collect::<Vec<_>>()
264 .join(",")
265 )
266}
267
268#[cfg(test)]
269mod tests {
270 use super::*;
271
272 fn trace(event_id: i64, verified_at_epoch: i64) -> ProcedureTrace {
273 ProcedureTrace {
274 project: "/tmp/remem".to_string(),
275 branch: Some("main".to_string()),
276 workflow_key: "pr-review-loop".to_string(),
277 command: "cargo test".to_string(),
278 files_touched: vec!["src/lib.rs".to_string()],
279 succeeded: true,
280 verified_at_epoch,
281 source_event_id: Some(event_id),
282 }
283 }
284
285 #[test]
286 fn repeated_verified_workflow_promotes_procedure_memory() -> Result<()> {
287 let conn = Connection::open_in_memory()?;
288 conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
289 crate::migrate::run_migrations(&conn)?;
290 let policy = ProcedurePromotionPolicy::default();
291 let candidate =
292 build_procedure_candidate(&[trace(10, 1_000), trace(11, 1_100)], 1_200, &policy)
293 .expect("two verified traces should promote");
294
295 assert_eq!(candidate.project, "/tmp/remem");
296 assert_eq!(candidate.branch.as_deref(), Some("main"));
297 assert_eq!(candidate.source_event_ids, vec![10, 11]);
298 assert_eq!(candidate.verified_runs, 2);
299 assert!(candidate.topic_key.contains("branch-main"));
300 assert!(candidate.topic_key.contains("command-cargo-test"));
301
302 let memory_id = promote_procedure_memory(&conn, &candidate)?;
303 let (memory_type, branch, evidence): (String, Option<String>, String) = conn.query_row(
304 "SELECT memory_type, branch, evidence_event_ids FROM memories WHERE id = ?1",
305 [memory_id],
306 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
307 )?;
308 assert_eq!(memory_type, "procedure");
309 assert_eq!(branch.as_deref(), Some("main"));
310 assert_eq!(serde_json::from_str::<Vec<i64>>(&evidence)?, vec![10, 11]);
311 Ok(())
312 }
313
314 #[test]
315 fn one_off_verified_workflow_does_not_promote() {
316 let policy = ProcedurePromotionPolicy::default();
317 let candidate = build_procedure_candidate(&[trace(10, 1_000)], 1_200, &policy);
318 assert!(candidate.is_none());
319 }
320
321 #[test]
322 fn missing_fresh_source_refs_do_not_promote() {
323 let policy = ProcedurePromotionPolicy::default();
324 let mut missing_source = trace(10, 1_000);
325 missing_source.source_event_id = None;
326 assert!(
327 build_procedure_candidate(&[missing_source, trace(11, 1_050)], 1_100, &policy)
328 .is_none()
329 );
330
331 let old = trace(12, 1_000);
332 let stale_now = 1_000 + DEFAULT_MAX_VERIFICATION_AGE_SECS + 1;
333 assert!(
334 build_procedure_candidate(&[old, trace(13, stale_now)], stale_now, &policy).is_none()
335 );
336 }
337
338 #[test]
339 fn mixed_project_or_branch_does_not_promote() {
340 let policy = ProcedurePromotionPolicy::default();
341 let mut other_project = trace(11, 1_100);
342 other_project.project = "/tmp/other".to_string();
343 assert!(
344 build_procedure_candidate(&[trace(10, 1_000), other_project], 1_200, &policy).is_none()
345 );
346
347 let mut other_branch = trace(12, 1_100);
348 other_branch.branch = Some("feature".to_string());
349 assert!(
350 build_procedure_candidate(&[trace(10, 1_000), other_branch], 1_200, &policy).is_none()
351 );
352 }
353
354 #[test]
355 fn production_task_promotes_repeated_successful_bash_procedure() -> Result<()> {
356 let mut conn = Connection::open_in_memory()?;
357 conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
358 crate::migrate::run_migrations(&conn)?;
359 let command = "cargo test";
360 for seq in [1, 2] {
361 crate::db::record_captured_event(
362 &conn,
363 &crate::db::CaptureEventInput {
364 host: "codex-cli",
365 session_id: "sess-procedure-runtime",
366 project: "/tmp/remem",
367 cwd: None,
368 event_type: "tool_result",
369 role: None,
370 tool_name: Some("Bash"),
371 content: &serde_json::json!({
372 "seq": seq,
373 "event_type": "bash",
374 "exit_code": 0,
375 "tool_input": { "command": command },
376 "files": "[\"src/lib.rs\"]",
377 "git_branch": "main"
378 })
379 .to_string(),
380 task_kind: Some(crate::db::ExtractionTaskKind::ObservationExtract),
381 },
382 )?;
383 }
384 conn.execute("UPDATE workspaces SET git_branch = 'feature'", [])?;
385 let task = crate::db::claim_next_extraction_task(&mut conn, "worker-a", 60)?
386 .expect("task should be claimed");
387
388 let promoted = promote_verified_procedures_for_task(
389 &conn,
390 &task,
391 &ProcedurePromotionPolicy::default(),
392 )?;
393
394 assert_eq!(promoted, 1);
395 let (memory_type, topic_key, branch, evidence): (String, String, Option<String>, String) = conn.query_row(
396 "SELECT memory_type, topic_key, branch, evidence_event_ids FROM memories WHERE memory_type = 'procedure'",
397 [],
398 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
399 )?;
400 assert_eq!(memory_type, "procedure");
401 assert_eq!(branch.as_deref(), Some("main"));
402 assert!(topic_key.contains("command-cargo-test"));
403 assert_eq!(serde_json::from_str::<Vec<i64>>(&evidence)?.len(), 2);
404 Ok(())
405 }
406
407 #[test]
408 fn production_task_ignores_procedure_events_outside_evidence_window() -> Result<()> {
409 let conn = Connection::open_in_memory()?;
410 conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
411 crate::migrate::run_migrations(&conn)?;
412 let command = "cargo test";
413 let mut old_high_watermark = 0;
414 for seq in [1, 2] {
415 let outcome = crate::db::record_captured_event(
416 &conn,
417 &crate::db::CaptureEventInput {
418 host: "codex-cli",
419 session_id: "sess-procedure-old",
420 project: "/tmp/remem",
421 cwd: None,
422 event_type: "tool_result",
423 role: None,
424 tool_name: Some("Bash"),
425 content: &serde_json::json!({
426 "seq": seq,
427 "event_type": "bash",
428 "exit_code": 0,
429 "tool_input": { "command": command },
430 "files": "[\"src/lib.rs\"]",
431 "git_branch": "main"
432 })
433 .to_string(),
434 task_kind: None,
435 },
436 )?;
437 old_high_watermark = outcome.event_row_id;
438 }
439 let current = crate::db::record_captured_event(
440 &conn,
441 &crate::db::CaptureEventInput {
442 host: "codex-cli",
443 session_id: "sess-procedure-current",
444 project: "/tmp/remem",
445 cwd: None,
446 event_type: "tool_result",
447 role: None,
448 tool_name: Some("Bash"),
449 content: &serde_json::json!({
450 "seq": 3,
451 "event_type": "bash",
452 "exit_code": 0,
453 "tool_input": { "command": command },
454 "files": "[\"src/lib.rs\"]",
455 "git_branch": "main"
456 })
457 .to_string(),
458 task_kind: None,
459 },
460 )?;
461 let (host_id, workspace_id, project_id, session_row_id): (i64, i64, i64, i64) = conn
462 .query_row(
463 "SELECT host_id, workspace_id, project_id, session_row_id
464 FROM captured_events
465 WHERE id = ?1",
466 [current.event_row_id],
467 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
468 )?;
469 let task = crate::db::ExtractionTask {
470 id: 1,
471 task_kind: crate::db::ExtractionTaskKind::ObservationExtract,
472 host_id,
473 workspace_id,
474 project_id,
475 session_row_id: Some(session_row_id),
476 host: "codex-cli".to_string(),
477 project: "/tmp/remem".to_string(),
478 session_id: Some("sess-procedure-current".to_string()),
479 ai_profile: None,
480 priority: crate::db::ExtractionTaskKind::ObservationExtract.priority(),
481 cursor_event_id: Some(old_high_watermark),
482 high_watermark_event_id: Some(current.event_row_id),
483 attempts: 0,
484 replay_range_id: None,
485 };
486
487 let promoted = promote_verified_procedures_for_task(
488 &conn,
489 &task,
490 &ProcedurePromotionPolicy::default(),
491 )?;
492
493 assert_eq!(promoted, 0);
494 let procedure_count: i64 = conn.query_row(
495 "SELECT COUNT(*) FROM memories WHERE memory_type = 'procedure'",
496 [],
497 |row| row.get(0),
498 )?;
499 assert_eq!(procedure_count, 0);
500 Ok(())
501 }
502
503 #[test]
504 fn production_task_accumulates_verified_runs_across_windows() -> Result<()> {
505 let mut conn = Connection::open_in_memory()?;
506 conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
507 crate::migrate::run_migrations(&conn)?;
508 let command = "cargo test";
509
510 crate::db::record_captured_event(
511 &conn,
512 &crate::db::CaptureEventInput {
513 host: "codex-cli",
514 session_id: "sess-procedure-windowed",
515 project: "/tmp/remem",
516 cwd: None,
517 event_type: "tool_result",
518 role: None,
519 tool_name: Some("Bash"),
520 content: &serde_json::json!({
521 "event_type": "bash",
522 "exit_code": 0,
523 "tool_input": { "command": command },
524 "files": "[\"src/lib.rs\"]",
525 "git_branch": "main"
526 })
527 .to_string(),
528 task_kind: Some(crate::db::ExtractionTaskKind::ObservationExtract),
529 },
530 )?;
531 let first_task = crate::db::claim_next_extraction_task(&mut conn, "worker-a", 60)?
532 .ok_or_else(|| anyhow::anyhow!("first task should be claimed"))?;
533 assert_eq!(
534 promote_verified_procedures_for_task(
535 &conn,
536 &first_task,
537 &ProcedurePromotionPolicy::default(),
538 )?,
539 0
540 );
541 crate::db::mark_extraction_task_done(
542 &conn,
543 first_task.id,
544 "worker-a",
545 first_task.high_watermark_event_id,
546 )?;
547
548 crate::db::record_captured_event(
549 &conn,
550 &crate::db::CaptureEventInput {
551 host: "codex-cli",
552 session_id: "sess-procedure-windowed",
553 project: "/tmp/remem",
554 cwd: None,
555 event_type: "tool_result",
556 role: None,
557 tool_name: Some("Bash"),
558 content: &serde_json::json!({
559 "event_type": "bash",
560 "exit_code": 0,
561 "tool_input": { "command": command },
562 "files": "[\"src/lib.rs\"]",
563 "git_branch": "main"
564 })
565 .to_string(),
566 task_kind: Some(crate::db::ExtractionTaskKind::ObservationExtract),
567 },
568 )?;
569 let second_task = crate::db::claim_next_extraction_task(&mut conn, "worker-b", 60)?
570 .ok_or_else(|| anyhow::anyhow!("second task should be claimed"))?;
571
572 assert_eq!(
573 promote_verified_procedures_for_task(
574 &conn,
575 &second_task,
576 &ProcedurePromotionPolicy::default(),
577 )?,
578 1
579 );
580 let evidence: String = conn.query_row(
581 "SELECT evidence_event_ids FROM memories WHERE memory_type = 'procedure'",
582 [],
583 |row| row.get(0),
584 )?;
585 assert_eq!(serde_json::from_str::<Vec<i64>>(&evidence)?.len(), 2);
586 Ok(())
587 }
588}