Skip to main content

remem/memory/procedure/
mod.rs

1use std::collections::BTreeMap;
2
3use anyhow::{bail, Context, 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 activation_tests;
24#[cfg(test)]
25mod incremental_tests;
26
27const DEFAULT_MIN_VERIFIED_RUNS: usize = 2;
28const DEFAULT_MAX_VERIFICATION_AGE_SECS: i64 = 14 * 24 * 60 * 60;
29
30#[derive(Debug, Clone, PartialEq, Eq)]
31pub struct ProcedureTrace {
32    pub project: String,
33    pub branch: Option<String>,
34    pub workflow_key: String,
35    pub command: String,
36    pub files_touched: Vec<String>,
37    pub succeeded: bool,
38    pub verified_at_epoch: i64,
39    pub source_event_id: Option<i64>,
40}
41
42#[derive(Debug, Clone, PartialEq, Eq)]
43pub struct ProcedurePromotionPolicy {
44    pub min_verified_runs: usize,
45    pub max_verification_age_secs: i64,
46}
47
48impl Default for ProcedurePromotionPolicy {
49    fn default() -> Self {
50        Self {
51            min_verified_runs: DEFAULT_MIN_VERIFIED_RUNS,
52            max_verification_age_secs: DEFAULT_MAX_VERIFICATION_AGE_SECS,
53        }
54    }
55}
56
57#[derive(Debug, Clone, PartialEq)]
58pub struct ProcedureCandidate {
59    pub project: String,
60    pub branch: Option<String>,
61    pub workflow_key: String,
62    pub topic_key: String,
63    pub title: String,
64    pub content: String,
65    pub files: Vec<String>,
66    pub source_event_ids: Vec<i64>,
67    pub verified_runs: usize,
68    pub confidence: f64,
69    pub verified_at_epoch: i64,
70}
71
72pub fn build_procedure_candidate(
73    traces: &[ProcedureTrace],
74    now_epoch: i64,
75    policy: &ProcedurePromotionPolicy,
76) -> Option<ProcedureCandidate> {
77    let mut verified: Vec<&ProcedureTrace> = traces
78        .iter()
79        .filter(|trace| trace.succeeded)
80        .filter(|trace| trace.source_event_id.is_some())
81        .filter(|trace| {
82            now_epoch.saturating_sub(trace.verified_at_epoch) <= policy.max_verification_age_secs
83        })
84        .collect();
85    verified.sort_by_key(|trace| trace.verified_at_epoch);
86    if verified.len() < policy.min_verified_runs {
87        return None;
88    }
89
90    let first = verified[0];
91    if verified.iter().any(|trace| {
92        trace.project != first.project
93            || trace.branch != first.branch
94            || trace.workflow_key != first.workflow_key
95            || trace.command != first.command
96    }) {
97        return None;
98    }
99
100    let mut source_event_ids: Vec<i64> = verified
101        .iter()
102        .filter_map(|trace| trace.source_event_id)
103        .collect();
104    source_event_ids.sort_unstable();
105    source_event_ids.dedup();
106    if source_event_ids.len() < policy.min_verified_runs {
107        return None;
108    }
109
110    let mut files = verified
111        .iter()
112        .flat_map(|trace| trace.files_touched.iter().cloned())
113        .collect::<Vec<_>>();
114    files.sort();
115    files.dedup();
116
117    let verified_at_epoch = verified
118        .iter()
119        .map(|trace| trace.verified_at_epoch)
120        .max()
121        .unwrap_or(now_epoch);
122    let topic_key = procedure_topic_key(first);
123    let confidence = confidence_for_verified_runs(source_event_ids.len());
124    let content = render_procedure_content(
125        first,
126        &files,
127        &source_event_ids,
128        verified.len(),
129        verified_at_epoch,
130    );
131
132    Some(ProcedureCandidate {
133        project: first.project.clone(),
134        branch: first.branch.clone(),
135        workflow_key: first.workflow_key.clone(),
136        title: format!("Procedure: {}", first.workflow_key),
137        topic_key,
138        content,
139        files,
140        source_event_ids,
141        verified_runs: verified.len(),
142        confidence,
143        verified_at_epoch,
144    })
145}
146
147fn confidence_for_verified_runs(verified_runs: usize) -> f64 {
148    (0.7 + (verified_runs as f64 * 0.08)).min(0.95)
149}
150
151pub fn promote_procedure_memory(conn: &Connection, candidate: &ProcedureCandidate) -> Result<i64> {
152    promote_procedure_memory_with_policy(conn, candidate, &ProcedurePromotionPolicy::default())
153}
154
155fn promote_procedure_memory_with_policy(
156    conn: &Connection,
157    candidate: &ProcedureCandidate,
158    policy: &ProcedurePromotionPolicy,
159) -> Result<i64> {
160    let tx = rusqlite::Transaction::new_unchecked(conn, rusqlite::TransactionBehavior::Immediate)?;
161    let files_json = (!candidate.files.is_empty())
162        .then(|| serde_json::to_string(&candidate.files))
163        .transpose()?;
164    let source_events_json = serde_json::to_string(&candidate.source_event_ids)?;
165    let existing_id = procedure_memory_id(
166        &tx,
167        &candidate.project,
168        &candidate.topic_key,
169        candidate.branch.as_deref(),
170    )?;
171    let branch_present = if candidate.branch.is_some() { "1" } else { "0" };
172    let confidence_bits = candidate.confidence.to_bits().to_string();
173    let verified_at_epoch = candidate.verified_at_epoch.to_string();
174    let payload_sha256 = crate::memory::activation::payload_sha256(&[
175        &candidate.project,
176        branch_present,
177        candidate.branch.as_deref().unwrap_or(""),
178        &candidate.topic_key,
179        &candidate.title,
180        &candidate.content,
181        files_json.as_deref().unwrap_or(""),
182        &source_events_json,
183        &confidence_bits,
184        &verified_at_epoch,
185    ]);
186    let activation_id =
187        crate::memory::activation::activation_id_from_key("procedure-promotion", &payload_sha256);
188    let replay_binding = tx
189        .query_row(
190            "SELECT result_memory_id, superseded_ids_json, result_source_trust_class,
191                    result_sha256
192             FROM memory_activation_requests WHERE activation_id = ?1",
193            [&activation_id],
194            |row| {
195                Ok((
196                    row.get::<_, i64>(0)?,
197                    row.get::<_, String>(1)?,
198                    row.get::<_, String>(2)?,
199                    row.get::<_, String>(3)?,
200                ))
201            },
202        )
203        .optional()?
204        .map(
205            |(memory_id, superseded_json, result_source_trust, result_sha256)| {
206                serde_json::from_str::<Vec<i64>>(&superseded_json)
207                    .map(|superseded_ids| {
208                        (
209                            memory_id,
210                            superseded_ids,
211                            result_source_trust,
212                            result_sha256,
213                        )
214                    })
215                    .context("invalid procedure activation superseded ids")
216            },
217        )
218        .transpose()?;
219    if replay_binding.is_none() {
220        evidence::validate_promotion_candidate(&tx, candidate, policy)?;
221    }
222    let binding_id = replay_binding
223        .as_ref()
224        .map(|(memory_id, _, _, _)| *memory_id)
225        .or(existing_id);
226    let retained_provenance = binding_id
227        .map(|memory_id| {
228            crate::memory::activation::ExpectedActiveMemory::from_existing(&tx, memory_id)
229        })
230        .transpose()?;
231    let result_source_trust = replay_binding
232        .as_ref()
233        .map(|(_, _, trust, _)| Ok(trust.clone()))
234        .or_else(|| {
235            binding_id.map(|memory_id| {
236                tx.query_row(
237                    "SELECT source_trust_class FROM memories WHERE id = ?1",
238                    [memory_id],
239                    |row| row.get::<_, String>(0),
240                )
241                .map_err(anyhow::Error::from)
242            })
243        })
244        .transpose()?
245        .map(|trust| {
246            let normalized = trust.strip_prefix("legacy_v086_source_").unwrap_or(&trust);
247            crate::memory::poisoning::SourceTrustClass::parse(normalized).ok_or_else(|| {
248                anyhow::anyhow!("existing procedure memory has invalid source trust: {trust}")
249            })
250        })
251        .transpose()?
252        .unwrap_or(crate::memory::poisoning::SourceTrustClass::LocalToolOutput);
253    let mut expected_memory = crate::memory::activation::ExpectedActiveMemory::new(
254        &candidate.title,
255        &candidate.content,
256        "procedure",
257    )
258    .with_topic_key(Some(&candidate.topic_key))
259    .with_files(files_json.as_deref())
260    .with_candidate_evidence(Some(&source_events_json), None);
261    let retained_candidate_id = if let Some((_, _, _, result_sha256)) = &replay_binding {
262        procedure_candidate_id_for_receipt(&tx, &expected_memory, result_sha256)?
263    } else {
264        retained_provenance
265            .as_ref()
266            .and_then(|memory| memory.source_candidate_id)
267    };
268    expected_memory.source_candidate_id = retained_candidate_id;
269    let superseded_ids = replay_binding
270        .map(|(_, superseded_ids, _, _)| superseded_ids)
271        .unwrap_or_else(|| existing_id.into_iter().collect());
272    let request = crate::memory::activation::ActiveMemoryWriteRequest {
273        activation_id,
274        route_kind: crate::memory::activation::ActivationRouteKind::CandidatePromotion,
275        actor_kind: crate::memory::activation::ActivationActorKind::AutomaticWorker,
276        source_operation: "procedure_promotion".to_string(),
277        source_trust: crate::memory::poisoning::SourceTrustClass::LocalToolOutput,
278        result_source_trust,
279        source_project: candidate.project.clone(),
280        route: crate::memory::activation::ActiveMemoryRoute::default_for(
281            &candidate.project,
282            candidate.branch.as_deref(),
283            "project",
284        ),
285        provenance_kind: crate::memory::activation::ActivationProvenanceKind::Candidate,
286        provenance_ref: format!("verified-procedure:{payload_sha256}"),
287        payload_sha256,
288        expected_memory,
289        poisoning_verdict: crate::memory::activation::ActivationPoisoningVerdict::UpstreamValidated,
290        superseded_ids,
291    };
292    let activation = crate::memory::activation::execute_one(&tx, &request, |permit| {
293        let memory_id = crate::memory::store::insert_memory_replacement_activated(
294            &tx,
295            permit,
296            existing_id,
297            None,
298            &candidate.project,
299            Some(&candidate.topic_key),
300            &candidate.title,
301            &candidate.content,
302            "procedure",
303            files_json.as_deref(),
304            candidate.branch.as_deref(),
305            "project",
306            result_source_trust,
307            Some(candidate.verified_at_epoch),
308            Some(candidate.verified_at_epoch),
309        )?;
310        tx.execute(
311            "UPDATE memories
312             SET evidence_event_ids = ?1,
313                 confidence = ?2,
314                 source_candidate_id = ?3
315             WHERE id = ?4",
316            params![
317                source_events_json,
318                candidate.confidence,
319                retained_candidate_id,
320                memory_id
321            ],
322        )?;
323        Ok(memory_id)
324    })?;
325    tx.commit()?;
326    Ok(activation.memory_id)
327}
328
329fn procedure_candidate_id_for_receipt(
330    conn: &Connection,
331    expected: &crate::memory::activation::ExpectedActiveMemory,
332    result_sha256: &str,
333) -> Result<Option<i64>> {
334    if expected.sha256() == result_sha256 {
335        return Ok(None);
336    }
337    let mut stmt = conn.prepare("SELECT id FROM memory_candidates ORDER BY id ASC")?;
338    let candidate_ids = stmt
339        .query_map([], |row| row.get::<_, i64>(0))?
340        .collect::<std::result::Result<Vec<_>, _>>()?;
341    for candidate_id in candidate_ids {
342        let mut candidate_expected = expected.clone();
343        candidate_expected.source_candidate_id = Some(candidate_id);
344        if candidate_expected.sha256() == result_sha256 {
345            return Ok(Some(candidate_id));
346        }
347    }
348    bail!("procedure activation receipt provenance cannot be reconstructed")
349}
350
351pub(crate) fn promote_verified_procedures_for_task(
352    conn: &Connection,
353    task: &crate::db::ExtractionTask,
354    policy: &ProcedurePromotionPolicy,
355) -> Result<usize> {
356    let now_epoch = chrono::Utc::now().timestamp();
357    let traces = trace_store::load_verified_procedure_traces(conn, task, policy, now_epoch)?;
358    let mut groups: BTreeMap<(String, Option<String>, String, String), Vec<ProcedureTrace>> =
359        BTreeMap::new();
360    for trace in traces {
361        groups
362            .entry((
363                trace.project.clone(),
364                trace.branch.clone(),
365                trace.workflow_key.clone(),
366                trace.command.clone(),
367            ))
368            .or_default()
369            .push(trace);
370    }
371
372    let mut promoted = 0usize;
373    for traces in groups.into_values() {
374        let Some(candidate) = build_procedure_candidate(&traces, now_epoch, policy) else {
375            continue;
376        };
377        let existed = procedure_memory_exists(
378            conn,
379            &candidate.project,
380            &candidate.topic_key,
381            candidate.branch.as_deref(),
382        )?;
383        promote_procedure_memory_with_policy(conn, &candidate, policy)?;
384        if !existed {
385            promoted += 1;
386        }
387    }
388    Ok(promoted)
389}
390
391fn procedure_topic_key(trace: &ProcedureTrace) -> String {
392    crate::memory::slugify_for_topic(
393        &format!(
394            "procedure {} branch {} command {}",
395            trace.workflow_key,
396            trace.branch.as_deref().unwrap_or("no-branch"),
397            trace.command
398        ),
399        96,
400    )
401}
402
403fn procedure_memory_id(
404    conn: &Connection,
405    project: &str,
406    topic_key: &str,
407    branch: Option<&str>,
408) -> Result<Option<i64>> {
409    conn.query_row(
410        "SELECT id FROM memories
411             WHERE project = ?1
412               AND topic_key = ?2
413               AND scope = 'project'
414               AND memory_type = 'procedure'
415               AND status = 'active'
416               AND branch IS ?3
417               AND COALESCE(owner_scope, 'repo') = 'repo'
418               AND COALESCE(owner_key, project) = ?1
419               AND COALESCE(target_project, project) = ?1
420             ORDER BY updated_at_epoch DESC, id DESC
421             LIMIT 1",
422        params![project, topic_key, branch],
423        |row| row.get(0),
424    )
425    .optional()
426    .map_err(Into::into)
427}
428
429fn procedure_memory_exists(
430    conn: &Connection,
431    project: &str,
432    topic_key: &str,
433    branch: Option<&str>,
434) -> Result<bool> {
435    Ok(procedure_memory_id(conn, project, topic_key, branch)?.is_some())
436}
437
438fn render_procedure_content(
439    trace: &ProcedureTrace,
440    files: &[String],
441    source_event_ids: &[i64],
442    verified_runs: usize,
443    verified_at_epoch: i64,
444) -> String {
445    let files_line = if files.is_empty() {
446        "Files: none recorded".to_string()
447    } else {
448        format!("Files: {}", files.join(", "))
449    };
450    format!(
451        "Procedure: {}\nCommand: {}\n{}\nVerified runs: {}\nVerified at: {}\nSource events: {}\nReuse when: the same project and branch need this verified workflow.",
452        trace.workflow_key,
453        trace.command,
454        files_line,
455        verified_runs,
456        verified_at_epoch,
457        source_event_ids
458            .iter()
459            .map(i64::to_string)
460            .collect::<Vec<_>>()
461            .join(",")
462    )
463}
464
465#[cfg(test)]
466mod tests {
467    use super::*;
468
469    fn trace(event_id: i64, verified_at_epoch: i64) -> ProcedureTrace {
470        ProcedureTrace {
471            project: "/tmp/remem".to_string(),
472            branch: Some("main".to_string()),
473            workflow_key: "pr-review-loop".to_string(),
474            command: "cargo test".to_string(),
475            files_touched: vec!["src/lib.rs".to_string()],
476            succeeded: true,
477            verified_at_epoch,
478            source_event_id: Some(event_id),
479        }
480    }
481
482    #[test]
483    fn repeated_verified_workflow_promotes_procedure_memory() -> Result<()> {
484        let policy = ProcedurePromotionPolicy::default();
485        let candidate =
486            build_procedure_candidate(&[trace(10, 1_000), trace(11, 1_100)], 1_200, &policy)
487                .expect("two verified traces should promote");
488
489        assert_eq!(candidate.project, "/tmp/remem");
490        assert_eq!(candidate.branch.as_deref(), Some("main"));
491        assert_eq!(candidate.source_event_ids, vec![10, 11]);
492        assert_eq!(candidate.verified_runs, 2);
493        assert!(candidate.topic_key.contains("branch-main"));
494        assert!(candidate.topic_key.contains("command-cargo-test"));
495
496        Ok(())
497    }
498
499    #[test]
500    fn one_off_verified_workflow_does_not_promote() {
501        let policy = ProcedurePromotionPolicy::default();
502        let candidate = build_procedure_candidate(&[trace(10, 1_000)], 1_200, &policy);
503        assert!(candidate.is_none());
504    }
505
506    #[test]
507    fn missing_fresh_source_refs_do_not_promote() {
508        let policy = ProcedurePromotionPolicy::default();
509        let mut missing_source = trace(10, 1_000);
510        missing_source.source_event_id = None;
511        assert!(
512            build_procedure_candidate(&[missing_source, trace(11, 1_050)], 1_100, &policy)
513                .is_none()
514        );
515
516        let old = trace(12, 1_000);
517        let stale_now = 1_000 + DEFAULT_MAX_VERIFICATION_AGE_SECS + 1;
518        assert!(
519            build_procedure_candidate(&[old, trace(13, stale_now)], stale_now, &policy).is_none()
520        );
521    }
522
523    #[test]
524    fn mixed_project_or_branch_does_not_promote() {
525        let policy = ProcedurePromotionPolicy::default();
526        let mut other_project = trace(11, 1_100);
527        other_project.project = "/tmp/other".to_string();
528        assert!(
529            build_procedure_candidate(&[trace(10, 1_000), other_project], 1_200, &policy).is_none()
530        );
531
532        let mut other_branch = trace(12, 1_100);
533        other_branch.branch = Some("feature".to_string());
534        assert!(
535            build_procedure_candidate(&[trace(10, 1_000), other_branch], 1_200, &policy).is_none()
536        );
537    }
538
539    #[test]
540    fn production_task_promotes_repeated_successful_bash_procedure() -> Result<()> {
541        let mut conn = Connection::open_in_memory()?;
542        conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
543        crate::migrate::run_migrations(&conn)?;
544        let command = "cargo test";
545        for seq in [1, 2] {
546            crate::db::record_captured_event(
547                &conn,
548                &crate::db::CaptureEventInput {
549                    host: "codex-cli",
550                    session_id: "sess-procedure-runtime",
551                    project: "/tmp/remem",
552                    cwd: None,
553                    event_type: "tool_result",
554                    role: None,
555                    tool_name: Some("Bash"),
556                    content: &serde_json::json!({
557                        "seq": seq,
558                        "event_type": "bash",
559                        "exit_code": 0,
560                        "tool_input": { "command": command },
561                        "files": "[\"src/lib.rs\"]",
562                        "git_branch": "main"
563                    })
564                    .to_string(),
565                    task_kind: Some(crate::db::ExtractionTaskKind::ObservationExtract),
566                },
567            )?;
568        }
569        conn.execute("UPDATE workspaces SET git_branch = 'feature'", [])?;
570        let task = crate::db::claim_next_extraction_task(&mut conn, "worker-a", 60)?
571            .expect("task should be claimed");
572
573        let promoted = promote_verified_procedures_for_task(
574            &conn,
575            &task,
576            &ProcedurePromotionPolicy::default(),
577        )?;
578
579        assert_eq!(promoted, 1);
580        let (memory_type, topic_key, branch, evidence): (String, String, Option<String>, String) = conn.query_row(
581            "SELECT memory_type, topic_key, branch, evidence_event_ids FROM memories WHERE memory_type = 'procedure'",
582            [],
583            |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
584        )?;
585        assert_eq!(memory_type, "procedure");
586        assert_eq!(branch.as_deref(), Some("main"));
587        assert!(topic_key.contains("command-cargo-test"));
588        assert_eq!(serde_json::from_str::<Vec<i64>>(&evidence)?.len(), 2);
589        Ok(())
590    }
591
592    #[test]
593    fn production_task_ignores_procedure_events_outside_evidence_window() -> Result<()> {
594        let conn = Connection::open_in_memory()?;
595        conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
596        crate::migrate::run_migrations(&conn)?;
597        let command = "cargo test";
598        let mut old_high_watermark = 0;
599        for seq in [1, 2] {
600            let outcome = crate::db::record_captured_event(
601                &conn,
602                &crate::db::CaptureEventInput {
603                    host: "codex-cli",
604                    session_id: "sess-procedure-old",
605                    project: "/tmp/remem",
606                    cwd: None,
607                    event_type: "tool_result",
608                    role: None,
609                    tool_name: Some("Bash"),
610                    content: &serde_json::json!({
611                        "seq": seq,
612                        "event_type": "bash",
613                        "exit_code": 0,
614                        "tool_input": { "command": command },
615                        "files": "[\"src/lib.rs\"]",
616                        "git_branch": "main"
617                    })
618                    .to_string(),
619                    task_kind: None,
620                },
621            )?;
622            old_high_watermark = outcome.event_row_id;
623        }
624        let current = crate::db::record_captured_event(
625            &conn,
626            &crate::db::CaptureEventInput {
627                host: "codex-cli",
628                session_id: "sess-procedure-current",
629                project: "/tmp/remem",
630                cwd: None,
631                event_type: "tool_result",
632                role: None,
633                tool_name: Some("Bash"),
634                content: &serde_json::json!({
635                    "seq": 3,
636                    "event_type": "bash",
637                    "exit_code": 0,
638                    "tool_input": { "command": command },
639                    "files": "[\"src/lib.rs\"]",
640                    "git_branch": "main"
641                })
642                .to_string(),
643                task_kind: None,
644            },
645        )?;
646        let (host_id, workspace_id, project_id, session_row_id): (i64, i64, i64, i64) = conn
647            .query_row(
648                "SELECT host_id, workspace_id, project_id, session_row_id
649                 FROM captured_events
650                 WHERE id = ?1",
651                [current.event_row_id],
652                |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
653            )?;
654        let task = crate::db::ExtractionTask {
655            id: 1,
656            task_kind: crate::db::ExtractionTaskKind::ObservationExtract,
657            host_id,
658            workspace_id,
659            project_id,
660            session_row_id: Some(session_row_id),
661            host: "codex-cli".to_string(),
662            project: "/tmp/remem".to_string(),
663            session_id: Some("sess-procedure-current".to_string()),
664            ai_profile: None,
665            priority: crate::db::ExtractionTaskKind::ObservationExtract.priority(),
666            cursor_event_id: Some(old_high_watermark),
667            high_watermark_event_id: Some(current.event_row_id),
668            attempts: 0,
669            replay_range_id: None,
670        };
671
672        let promoted = promote_verified_procedures_for_task(
673            &conn,
674            &task,
675            &ProcedurePromotionPolicy::default(),
676        )?;
677
678        assert_eq!(promoted, 0);
679        let procedure_count: i64 = conn.query_row(
680            "SELECT COUNT(*) FROM memories WHERE memory_type = 'procedure'",
681            [],
682            |row| row.get(0),
683        )?;
684        assert_eq!(procedure_count, 0);
685        Ok(())
686    }
687
688    #[test]
689    fn production_task_accumulates_verified_runs_across_windows() -> Result<()> {
690        let mut conn = Connection::open_in_memory()?;
691        conn.execute_batch("PRAGMA journal_mode=WAL; PRAGMA foreign_keys=ON;")?;
692        crate::migrate::run_migrations(&conn)?;
693        let command = "cargo test";
694
695        crate::db::record_captured_event(
696            &conn,
697            &crate::db::CaptureEventInput {
698                host: "codex-cli",
699                session_id: "sess-procedure-windowed",
700                project: "/tmp/remem",
701                cwd: None,
702                event_type: "tool_result",
703                role: None,
704                tool_name: Some("Bash"),
705                content: &serde_json::json!({
706                    "event_type": "bash",
707                    "exit_code": 0,
708                    "tool_input": { "command": command },
709                    "files": "[\"src/lib.rs\"]",
710                    "git_branch": "main"
711                })
712                .to_string(),
713                task_kind: Some(crate::db::ExtractionTaskKind::ObservationExtract),
714            },
715        )?;
716        let first_task = crate::db::claim_next_extraction_task(&mut conn, "worker-a", 60)?
717            .ok_or_else(|| anyhow::anyhow!("first task should be claimed"))?;
718        assert_eq!(
719            promote_verified_procedures_for_task(
720                &conn,
721                &first_task,
722                &ProcedurePromotionPolicy::default(),
723            )?,
724            0
725        );
726        crate::db::mark_extraction_task_done(
727            &conn,
728            first_task.id,
729            "worker-a",
730            first_task.high_watermark_event_id,
731        )?;
732
733        crate::db::record_captured_event(
734            &conn,
735            &crate::db::CaptureEventInput {
736                host: "codex-cli",
737                session_id: "sess-procedure-windowed",
738                project: "/tmp/remem",
739                cwd: None,
740                event_type: "tool_result",
741                role: None,
742                tool_name: Some("Bash"),
743                content: &serde_json::json!({
744                    "event_type": "bash",
745                    "exit_code": 0,
746                    "tool_input": { "command": command },
747                    "files": "[\"src/lib.rs\"]",
748                    "git_branch": "main"
749                })
750                .to_string(),
751                task_kind: Some(crate::db::ExtractionTaskKind::ObservationExtract),
752            },
753        )?;
754        let second_task = crate::db::claim_next_extraction_task(&mut conn, "worker-b", 60)?
755            .ok_or_else(|| anyhow::anyhow!("second task should be claimed"))?;
756
757        assert_eq!(
758            promote_verified_procedures_for_task(
759                &conn,
760                &second_task,
761                &ProcedurePromotionPolicy::default(),
762            )?,
763            1
764        );
765        let evidence: String = conn.query_row(
766            "SELECT evidence_event_ids FROM memories WHERE memory_type = 'procedure'",
767            [],
768            |row| row.get(0),
769        )?;
770        assert_eq!(serde_json::from_str::<Vec<i64>>(&evidence)?.len(), 2);
771        Ok(())
772    }
773}