Skip to main content

remem/memory/procedure/
mod.rs

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}