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}