1use std::collections::BTreeMap;
2#[cfg(not(target_arch = "wasm32"))]
3use std::path::{Path, PathBuf};
4use std::sync::Arc;
5
6use async_trait::async_trait;
7use chrono::{DateTime, Utc};
8#[cfg(not(target_arch = "wasm32"))]
9use rusqlite::{
10 Connection, Error, ErrorCode, OptionalExtension, Transaction, TransactionBehavior, params,
11};
12
13use crate::WorkGraphError;
14use crate::types::{
15 AttentionListRequest, AttentionPruneRequest, WorkAttentionBinding, WorkAttentionBindingId,
16 WorkAttentionStatus, WorkEdge, WorkExecutionBinding, WorkExecutionBindingFilter,
17 WorkExecutionBindingId, WorkGraphEvent, WorkGraphEventKind, WorkItem, WorkItemFilter,
18 WorkItemId, WorkNamespace,
19};
20use crate::{WorkAttentionMachine, WorkGraphMachine};
21
22#[cfg(target_arch = "wasm32")]
23use crate::tokio::sync::RwLock;
24#[cfg(not(target_arch = "wasm32"))]
25use tokio::sync::RwLock;
26
27#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub enum WorkGraphStoreKind {
29 Disabled,
30 Memory,
31 Sqlite,
32 Custom,
33}
34
35impl WorkGraphStoreKind {
36 pub fn as_str(self) -> &'static str {
37 match self {
38 Self::Disabled => "disabled",
39 Self::Memory => "memory",
40 Self::Sqlite => "sqlite",
41 Self::Custom => "custom",
42 }
43 }
44}
45
46impl std::fmt::Display for WorkGraphStoreKind {
47 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
48 f.write_str(self.as_str())
49 }
50}
51
52#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
53#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
54pub struct WorkGraphEventFilter {
55 pub realm_id: Option<String>,
56 pub namespace: Option<WorkNamespace>,
57 #[serde(default)]
58 pub all_namespaces: bool,
59 pub after_seq: Option<i64>,
60 pub limit: Option<usize>,
61}
62
63#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
64#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
65pub trait WorkGraphStore: Send + Sync {
66 fn kind(&self) -> WorkGraphStoreKind;
67
68 async fn get_store_time_utc(&self) -> Result<DateTime<Utc>, WorkGraphError>;
69
70 async fn insert_item(
71 &self,
72 item: WorkItem,
73 event: WorkGraphEvent,
74 ) -> Result<WorkItem, WorkGraphError>;
75
76 async fn update_item_cas(
77 &self,
78 item: WorkItem,
79 expected_previous_revision: u64,
80 event: WorkGraphEvent,
81 ) -> Result<WorkItem, WorkGraphError>;
82
83 async fn update_item_and_attention_cas(
84 &self,
85 item: WorkItem,
86 expected_previous_revision: u64,
87 item_event: WorkGraphEvent,
88 attention_updates: Vec<(WorkAttentionBinding, u64, WorkGraphEvent)>,
89 ) -> Result<WorkItem, WorkGraphError>;
90
91 async fn get_item(
92 &self,
93 realm_id: &str,
94 namespace: &WorkNamespace,
95 id: &WorkItemId,
96 ) -> Result<Option<WorkItem>, WorkGraphError>;
97
98 async fn list_items(&self, filter: WorkItemFilter) -> Result<Vec<WorkItem>, WorkGraphError>;
99
100 async fn insert_goal(
101 &self,
102 _item: WorkItem,
103 _item_event: WorkGraphEvent,
104 _attention: WorkAttentionBinding,
105 _attention_event: WorkGraphEvent,
106 ) -> Result<(WorkItem, WorkAttentionBinding), WorkGraphError> {
107 Err(unsupported(self.kind()))
108 }
109
110 async fn update_attention_cas(
111 &self,
112 _attention: WorkAttentionBinding,
113 _expected_previous_revision: u64,
114 _event: WorkGraphEvent,
115 ) -> Result<WorkAttentionBinding, WorkGraphError> {
116 Err(unsupported(self.kind()))
117 }
118
119 async fn reassign_attention_cas(
120 &self,
121 _previous: WorkAttentionBinding,
122 _expected_previous_revision: u64,
123 _previous_event: WorkGraphEvent,
124 _replacement: WorkAttentionBinding,
125 _replacement_event: WorkGraphEvent,
126 ) -> Result<(WorkAttentionBinding, WorkAttentionBinding), WorkGraphError> {
127 Err(unsupported(self.kind()))
128 }
129
130 async fn get_attention(
131 &self,
132 _realm_id: &str,
133 _namespace: &WorkNamespace,
134 _binding_id: &WorkAttentionBindingId,
135 ) -> Result<Option<WorkAttentionBinding>, WorkGraphError> {
136 Err(unsupported(self.kind()))
137 }
138
139 async fn list_attention(
140 &self,
141 _filter: AttentionListRequest,
142 ) -> Result<Vec<WorkAttentionBinding>, WorkGraphError> {
143 Err(unsupported(self.kind()))
144 }
145
146 async fn insert_execution_binding(
150 &self,
151 _commit: crate::WorkExecutionBindCommit,
152 _expected_item_revision: u64,
153 _event: WorkGraphEvent,
154 ) -> Result<WorkExecutionBinding, WorkGraphError> {
155 Err(unsupported(self.kind()))
156 }
157
158 async fn get_execution_binding(
159 &self,
160 _realm_id: &str,
161 _namespace: &WorkNamespace,
162 _binding_id: &WorkExecutionBindingId,
163 ) -> Result<Option<WorkExecutionBinding>, WorkGraphError> {
164 Err(unsupported(self.kind()))
165 }
166
167 async fn get_execution_binding_by_target_run(
171 &self,
172 _realm_id: &str,
173 _run_id: &str,
174 ) -> Result<Option<WorkExecutionBinding>, WorkGraphError> {
175 Err(unsupported(self.kind()))
176 }
177
178 async fn update_execution_binding_cas(
179 &self,
180 _commit: crate::WorkExecutionObservationCommit,
181 _expected_previous_revision: u64,
182 _event: WorkGraphEvent,
183 ) -> Result<WorkExecutionBinding, WorkGraphError> {
184 Err(unsupported(self.kind()))
185 }
186
187 async fn list_execution_bindings(
188 &self,
189 _filter: WorkExecutionBindingFilter,
190 ) -> Result<Vec<WorkExecutionBinding>, WorkGraphError> {
191 Err(unsupported(self.kind()))
192 }
193
194 async fn list_execution_bindings_for_recovery(
198 &self,
199 realm_id: &str,
200 ) -> Result<Vec<WorkExecutionBinding>, WorkGraphError> {
201 let mut bindings = self
202 .list_execution_bindings(WorkExecutionBindingFilter {
203 realm_id: Some(realm_id.to_string()),
204 namespace: None,
205 item_id: None,
206 current_only: true,
207 limit: None,
208 })
209 .await?;
210 let mut active = Vec::with_capacity(bindings.len());
211 for binding in bindings.drain(..) {
212 if !crate::WorkExecutionMachine::retry_eligible(&binding)? {
213 active.push(binding);
214 }
215 }
216 Ok(active)
217 }
218
219 async fn list_attention_bounded(
223 &self,
224 filter: AttentionListRequest,
225 limit: usize,
226 ) -> Result<Vec<WorkAttentionBinding>, WorkGraphError> {
227 let mut bindings = self.list_attention(filter).await?;
228 bindings.truncate(limit);
229 Ok(bindings)
230 }
231
232 async fn prune_terminal_attention(
236 &self,
237 _filter: AttentionPruneRequest,
238 ) -> Result<u64, WorkGraphError> {
239 Err(unsupported(self.kind()))
240 }
241
242 async fn insert_edge(
243 &self,
244 edge: WorkEdge,
245 event: WorkGraphEvent,
246 ) -> Result<WorkEdge, WorkGraphError>;
247
248 async fn insert_edge_validated(
249 &self,
250 _edge: WorkEdge,
251 _event: WorkGraphEvent,
252 ) -> Result<WorkEdge, WorkGraphError> {
253 Err(unsupported(self.kind()))
254 }
255
256 async fn list_edges(
257 &self,
258 realm_id: &str,
259 namespace: &WorkNamespace,
260 ) -> Result<Vec<WorkEdge>, WorkGraphError>;
261
262 async fn list_edges_bounded(
264 &self,
265 realm_id: &str,
266 namespace: &WorkNamespace,
267 limit: usize,
268 ) -> Result<Vec<WorkEdge>, WorkGraphError> {
269 let mut edges = self.list_edges(realm_id, namespace).await?;
270 edges.truncate(limit);
271 Ok(edges)
272 }
273
274 async fn list_events(
275 &self,
276 filter: WorkGraphEventFilter,
277 ) -> Result<Vec<WorkGraphEvent>, WorkGraphError>;
278
279 async fn list_public_events(
282 &self,
283 mut filter: WorkGraphEventFilter,
284 ) -> Result<Vec<WorkGraphEvent>, WorkGraphError> {
285 let visible_limit = filter.limit.unwrap_or(usize::MAX);
286 if visible_limit == 0 {
287 return Ok(Vec::new());
288 }
289 filter.limit = Some(visible_limit);
293 Ok(self
294 .list_events(filter)
295 .await?
296 .into_iter()
297 .filter(|event| !is_internal_execution_event(event.kind))
298 .collect())
299 }
300
301 async fn latest_event_seq(
303 &self,
304 filter: WorkGraphEventFilter,
305 ) -> Result<Option<i64>, WorkGraphError> {
306 Ok(self
307 .list_events(filter)
308 .await?
309 .into_iter()
310 .filter_map(|event| event.seq)
311 .max())
312 }
313}
314
315#[derive(Default)]
316pub struct DisabledWorkGraphStore;
317
318#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
319#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
320impl WorkGraphStore for DisabledWorkGraphStore {
321 fn kind(&self) -> WorkGraphStoreKind {
322 WorkGraphStoreKind::Disabled
323 }
324
325 async fn get_store_time_utc(&self) -> Result<DateTime<Utc>, WorkGraphError> {
326 Err(unsupported(self.kind()))
327 }
328
329 async fn insert_item(
330 &self,
331 _item: WorkItem,
332 _event: WorkGraphEvent,
333 ) -> Result<WorkItem, WorkGraphError> {
334 Err(unsupported(self.kind()))
335 }
336
337 async fn update_item_cas(
338 &self,
339 _item: WorkItem,
340 _expected_previous_revision: u64,
341 _event: WorkGraphEvent,
342 ) -> Result<WorkItem, WorkGraphError> {
343 Err(unsupported(self.kind()))
344 }
345
346 async fn update_item_and_attention_cas(
347 &self,
348 _item: WorkItem,
349 _expected_previous_revision: u64,
350 _item_event: WorkGraphEvent,
351 _attention_updates: Vec<(WorkAttentionBinding, u64, WorkGraphEvent)>,
352 ) -> Result<WorkItem, WorkGraphError> {
353 Err(unsupported(self.kind()))
354 }
355
356 async fn get_item(
357 &self,
358 _realm_id: &str,
359 _namespace: &WorkNamespace,
360 _id: &WorkItemId,
361 ) -> Result<Option<WorkItem>, WorkGraphError> {
362 Err(unsupported(self.kind()))
363 }
364
365 async fn list_items(&self, _filter: WorkItemFilter) -> Result<Vec<WorkItem>, WorkGraphError> {
366 Err(unsupported(self.kind()))
367 }
368
369 async fn insert_goal(
370 &self,
371 _item: WorkItem,
372 _item_event: WorkGraphEvent,
373 _attention: WorkAttentionBinding,
374 _attention_event: WorkGraphEvent,
375 ) -> Result<(WorkItem, WorkAttentionBinding), WorkGraphError> {
376 Err(unsupported(self.kind()))
377 }
378
379 async fn update_attention_cas(
380 &self,
381 _attention: WorkAttentionBinding,
382 _expected_previous_revision: u64,
383 _event: WorkGraphEvent,
384 ) -> Result<WorkAttentionBinding, WorkGraphError> {
385 Err(unsupported(self.kind()))
386 }
387
388 async fn get_attention(
389 &self,
390 _realm_id: &str,
391 _namespace: &WorkNamespace,
392 _binding_id: &WorkAttentionBindingId,
393 ) -> Result<Option<WorkAttentionBinding>, WorkGraphError> {
394 Err(unsupported(self.kind()))
395 }
396
397 async fn list_attention(
398 &self,
399 _filter: AttentionListRequest,
400 ) -> Result<Vec<WorkAttentionBinding>, WorkGraphError> {
401 Err(unsupported(self.kind()))
402 }
403
404 async fn insert_edge(
405 &self,
406 _edge: WorkEdge,
407 _event: WorkGraphEvent,
408 ) -> Result<WorkEdge, WorkGraphError> {
409 Err(unsupported(self.kind()))
410 }
411
412 async fn insert_edge_validated(
413 &self,
414 _edge: WorkEdge,
415 _event: WorkGraphEvent,
416 ) -> Result<WorkEdge, WorkGraphError> {
417 Err(unsupported(self.kind()))
418 }
419
420 async fn list_edges(
421 &self,
422 _realm_id: &str,
423 _namespace: &WorkNamespace,
424 ) -> Result<Vec<WorkEdge>, WorkGraphError> {
425 Err(unsupported(self.kind()))
426 }
427
428 async fn list_events(
429 &self,
430 _filter: WorkGraphEventFilter,
431 ) -> Result<Vec<WorkGraphEvent>, WorkGraphError> {
432 Err(unsupported(self.kind()))
433 }
434}
435
436fn unsupported(kind: WorkGraphStoreKind) -> WorkGraphError {
437 WorkGraphError::UnsupportedBackend(kind.to_string())
438}
439
440#[derive(Default)]
441pub struct MemoryWorkGraphStore {
442 inner: Arc<RwLock<MemoryWorkGraphState>>,
443}
444
445#[derive(Default)]
446struct MemoryWorkGraphState {
447 items: BTreeMap<(String, WorkNamespace, WorkItemId), WorkItem>,
448 attention: BTreeMap<(String, WorkNamespace, WorkAttentionBindingId), WorkAttentionBinding>,
449 execution_bindings:
450 BTreeMap<(String, WorkNamespace, WorkExecutionBindingId), WorkExecutionBinding>,
451 execution_recovery: std::collections::BTreeSet<(String, WorkNamespace, WorkExecutionBindingId)>,
452 edges: Vec<WorkEdge>,
453 events: Vec<WorkGraphEvent>,
454 next_event_seq: i64,
455}
456
457impl MemoryWorkGraphStore {
458 pub fn new() -> Self {
459 Self::default()
460 }
461}
462
463#[cfg_attr(target_arch = "wasm32", async_trait(?Send))]
464#[cfg_attr(not(target_arch = "wasm32"), async_trait)]
465impl WorkGraphStore for MemoryWorkGraphStore {
466 fn kind(&self) -> WorkGraphStoreKind {
467 WorkGraphStoreKind::Memory
468 }
469
470 async fn get_store_time_utc(&self) -> Result<DateTime<Utc>, WorkGraphError> {
471 Ok(Utc::now())
472 }
473
474 async fn insert_item(
475 &self,
476 item: WorkItem,
477 event: WorkGraphEvent,
478 ) -> Result<WorkItem, WorkGraphError> {
479 WorkGraphMachine::validate_item_projection(&item)?;
480 let mut guard = self.inner.write().await;
481 let key = item_key(&item.realm_id, &item.namespace, &item.id);
482 if guard.items.contains_key(&key) {
483 return Err(WorkGraphError::Conflict(format!(
484 "work item {} already exists",
485 item.id
486 )));
487 }
488 guard.items.insert(key, item.clone());
489 guard.append_event(event);
490 Ok(item)
491 }
492
493 async fn update_item_cas(
494 &self,
495 item: WorkItem,
496 expected_previous_revision: u64,
497 event: WorkGraphEvent,
498 ) -> Result<WorkItem, WorkGraphError> {
499 WorkGraphMachine::validate_item_projection(&item)?;
500 let mut guard = self.inner.write().await;
501 let key = item_key(&item.realm_id, &item.namespace, &item.id);
502 let Some(current) = guard.items.get(&key) else {
503 return Err(WorkGraphError::not_found(
504 item.realm_id.clone(),
505 item.namespace.clone(),
506 item.id.clone(),
507 ));
508 };
509 if current.revision != expected_previous_revision {
510 return Err(WorkGraphError::StaleRevision {
511 id: item.id.clone(),
512 expected: expected_previous_revision,
513 actual: current.revision,
514 });
515 }
516 guard.items.insert(key, item.clone());
517 guard.append_event(event);
518 Ok(item)
519 }
520
521 async fn get_item(
522 &self,
523 realm_id: &str,
524 namespace: &WorkNamespace,
525 id: &WorkItemId,
526 ) -> Result<Option<WorkItem>, WorkGraphError> {
527 let guard = self.inner.read().await;
528 Ok(guard.items.get(&item_key(realm_id, namespace, id)).cloned())
529 }
530
531 async fn list_items(&self, filter: WorkItemFilter) -> Result<Vec<WorkItem>, WorkGraphError> {
532 let guard = self.inner.read().await;
533 let compare = |left: &WorkItem, right: &WorkItem| {
534 left.updated_at
535 .cmp(&right.updated_at)
536 .then_with(|| left.id.cmp(&right.id))
537 };
538 if let Some(limit) = filter.limit {
539 let mut items = Vec::with_capacity(limit.min(1024));
540 for item in guard
541 .items
542 .values()
543 .filter(|item| item_matches_filter(item, &filter))
544 {
545 let index = items
546 .binary_search_by(|existing| compare(existing, item))
547 .unwrap_or_else(|index| index);
548 if index < limit {
549 items.insert(index, item.clone());
550 if items.len() > limit {
551 items.pop();
552 }
553 }
554 }
555 return Ok(items);
556 }
557 let mut items = guard
558 .items
559 .values()
560 .filter(|item| item_matches_filter(item, &filter))
561 .cloned()
562 .collect::<Vec<_>>();
563 items.sort_by(compare);
564 Ok(items)
565 }
566
567 async fn insert_execution_binding(
568 &self,
569 commit: crate::WorkExecutionBindCommit,
570 expected_item_revision: u64,
571 event: WorkGraphEvent,
572 ) -> Result<WorkExecutionBinding, WorkGraphError> {
573 let (binding, effect) = commit.into_parts();
574 binding.validate()?;
575 crate::WorkExecutionMachine::validate_projection(&binding)?;
576 let (expected_state, expected_effect) =
577 crate::WorkExecutionMachine::bind(&binding.binding_id, binding.target.run_id())?;
578 if binding.machine_state != expected_state || effect != expected_effect {
579 return Err(WorkGraphError::InvalidInput(format!(
580 "work execution binding {} lacks canonical bind authority",
581 binding.binding_id
582 )));
583 }
584 let mut guard = self.inner.write().await;
585 let key = execution_binding_key(
586 &binding.work_ref.realm_id,
587 &binding.work_ref.namespace,
588 &binding.binding_id,
589 );
590 if let Some(existing) = guard.execution_bindings.get(&key) {
591 return if existing == &binding {
592 Ok(existing.clone())
593 } else {
594 Err(WorkGraphError::Conflict(format!(
595 "work execution binding {} already exists with different content",
596 binding.binding_id
597 )))
598 };
599 }
600 validate_execution_binding_insert(
601 &binding,
602 expected_item_revision,
603 guard.items.values(),
604 guard.execution_bindings.values(),
605 )?;
606 if !crate::WorkExecutionMachine::retry_eligible(&binding)? {
607 guard.execution_recovery.insert(key.clone());
608 }
609 guard.execution_bindings.insert(key, binding.clone());
610 guard.append_event(event);
611 Ok(binding)
612 }
613
614 async fn get_execution_binding(
615 &self,
616 realm_id: &str,
617 namespace: &WorkNamespace,
618 binding_id: &WorkExecutionBindingId,
619 ) -> Result<Option<WorkExecutionBinding>, WorkGraphError> {
620 let guard = self.inner.read().await;
621 Ok(guard
622 .execution_bindings
623 .get(&execution_binding_key(realm_id, namespace, binding_id))
624 .cloned())
625 }
626
627 async fn get_execution_binding_by_target_run(
628 &self,
629 realm_id: &str,
630 run_id: &str,
631 ) -> Result<Option<WorkExecutionBinding>, WorkGraphError> {
632 let guard = self.inner.read().await;
633 Ok(guard
634 .execution_bindings
635 .values()
636 .find(|binding| {
637 binding.work_ref.realm_id == realm_id && binding.target.run_id() == run_id
638 })
639 .cloned())
640 }
641
642 async fn update_execution_binding_cas(
643 &self,
644 commit: crate::WorkExecutionObservationCommit,
645 expected_previous_revision: u64,
646 event: WorkGraphEvent,
647 ) -> Result<WorkExecutionBinding, WorkGraphError> {
648 let (previous, observation, binding, effect) = commit.into_parts();
649 crate::WorkExecutionMachine::validate_projection(&binding)?;
650 let mut guard = self.inner.write().await;
651 let key = execution_binding_key(
652 &binding.work_ref.realm_id,
653 &binding.work_ref.namespace,
654 &binding.binding_id,
655 );
656 let current = guard.execution_bindings.get(&key).ok_or_else(|| {
657 WorkGraphError::Conflict(format!(
658 "work execution binding {} does not exist",
659 binding.binding_id
660 ))
661 })?;
662 if current.machine_state.revision != expected_previous_revision {
663 return Err(WorkGraphError::Conflict(format!(
664 "stale work execution revision for {}: expected {}, actual {}",
665 binding.binding_id, expected_previous_revision, current.machine_state.revision
666 )));
667 }
668 if current != &previous {
669 return Err(WorkGraphError::Conflict(format!(
670 "work execution transition authority for {} was minted from a different predecessor",
671 binding.binding_id
672 )));
673 }
674 let (expected_binding, expected_effect) = crate::WorkExecutionMachine::observe(
675 current.clone(),
676 expected_previous_revision,
677 observation,
678 )?;
679 if binding != expected_binding || effect != expected_effect {
680 return Err(WorkGraphError::Conflict(format!(
681 "work execution transition authority for {} does not match the generated machine result",
682 binding.binding_id
683 )));
684 }
685 if !current.has_same_immutable_spec(&binding) {
686 return Err(WorkGraphError::Conflict(format!(
687 "immutable work execution specification changed for {}",
688 binding.binding_id
689 )));
690 }
691 if crate::WorkExecutionMachine::retry_eligible(&binding)? {
692 guard.execution_recovery.remove(&key);
693 } else {
694 guard.execution_recovery.insert(key.clone());
695 }
696 guard.execution_bindings.insert(key, binding.clone());
697 guard.append_event(event);
698 Ok(binding)
699 }
700
701 async fn list_execution_bindings(
702 &self,
703 filter: WorkExecutionBindingFilter,
704 ) -> Result<Vec<WorkExecutionBinding>, WorkGraphError> {
705 let guard = self.inner.read().await;
706 let superseded = guard
707 .execution_bindings
708 .values()
709 .filter_map(|binding| {
710 binding.supersedes.clone().map(|supersedes| {
711 (
712 binding.work_ref.realm_id.clone(),
713 binding.work_ref.namespace.clone(),
714 supersedes,
715 )
716 })
717 })
718 .collect::<std::collections::BTreeSet<_>>();
719 let mut bindings = guard
720 .execution_bindings
721 .values()
722 .filter(|binding| execution_binding_matches_filter(binding, &filter, &superseded))
723 .cloned()
724 .collect::<Vec<_>>();
725 bindings.sort_by(|left, right| {
726 left.created_at
727 .cmp(&right.created_at)
728 .then_with(|| left.binding_id.cmp(&right.binding_id))
729 });
730 if let Some(limit) = filter.limit {
731 bindings.truncate(limit);
732 }
733 Ok(bindings)
734 }
735
736 async fn list_execution_bindings_for_recovery(
737 &self,
738 realm_id: &str,
739 ) -> Result<Vec<WorkExecutionBinding>, WorkGraphError> {
740 let guard = self.inner.read().await;
741 Ok(guard
742 .execution_recovery
743 .iter()
744 .filter(|(realm, _, _)| realm == realm_id)
745 .filter_map(|key| guard.execution_bindings.get(key).cloned())
746 .collect())
747 }
748
749 async fn insert_goal(
750 &self,
751 item: WorkItem,
752 item_event: WorkGraphEvent,
753 attention: WorkAttentionBinding,
754 attention_event: WorkGraphEvent,
755 ) -> Result<(WorkItem, WorkAttentionBinding), WorkGraphError> {
756 WorkGraphMachine::validate_item_projection(&item)?;
757 let mut guard = self.inner.write().await;
758 let item_key = item_key(&item.realm_id, &item.namespace, &item.id);
759 if guard.items.contains_key(&item_key) {
760 return Err(WorkGraphError::Conflict(format!(
761 "work item {} already exists",
762 item.id
763 )));
764 }
765 let attention_key = attention_key(
766 &attention.work_ref.realm_id,
767 &attention.work_ref.namespace,
768 &attention.binding_id,
769 );
770 if guard.attention.contains_key(&attention_key) {
771 return Err(WorkGraphError::Conflict(format!(
772 "work attention binding {} already exists",
773 attention.binding_id
774 )));
775 }
776 if let Some(occupant) = active_target_occupant_in(guard.attention.values(), &attention) {
777 return Err(active_target_conflict(&attention, &occupant));
778 }
779 guard.items.insert(item_key, item.clone());
780 guard.attention.insert(attention_key, attention.clone());
781 guard.append_event(item_event);
782 guard.append_event(attention_event);
783 Ok((item, attention))
784 }
785
786 async fn update_attention_cas(
787 &self,
788 attention: WorkAttentionBinding,
789 expected_previous_revision: u64,
790 event: WorkGraphEvent,
791 ) -> Result<WorkAttentionBinding, WorkGraphError> {
792 let mut guard = self.inner.write().await;
793 let key = attention_key(
794 &attention.work_ref.realm_id,
795 &attention.work_ref.namespace,
796 &attention.binding_id,
797 );
798 let Some(current) = guard.attention.get(&key) else {
799 return Err(WorkGraphError::not_found(
800 attention.work_ref.realm_id.clone(),
801 attention.work_ref.namespace.clone(),
802 attention.work_ref.item_id.clone(),
803 ));
804 };
805 if current.machine_state.revision != expected_previous_revision {
806 return Err(WorkGraphError::StaleRevision {
807 id: attention.work_ref.item_id.clone(),
808 expected: expected_previous_revision,
809 actual: current.machine_state.revision,
810 });
811 }
812 if let Some(occupant) = active_target_occupant_in(guard.attention.values(), &attention) {
813 return Err(active_target_conflict(&attention, &occupant));
814 }
815 guard.attention.insert(key, attention.clone());
816 guard.append_event(event);
817 Ok(attention)
818 }
819
820 async fn reassign_attention_cas(
821 &self,
822 previous: WorkAttentionBinding,
823 expected_previous_revision: u64,
824 previous_event: WorkGraphEvent,
825 replacement: WorkAttentionBinding,
826 replacement_event: WorkGraphEvent,
827 ) -> Result<(WorkAttentionBinding, WorkAttentionBinding), WorkGraphError> {
828 let mut guard = self.inner.write().await;
829 let previous_key = attention_key(
830 &previous.work_ref.realm_id,
831 &previous.work_ref.namespace,
832 &previous.binding_id,
833 );
834 let Some(current) = guard.attention.get(&previous_key) else {
835 return Err(WorkGraphError::attention_not_found(
836 previous.work_ref.realm_id.clone(),
837 previous.work_ref.namespace.clone(),
838 previous.binding_id.clone(),
839 ));
840 };
841 if current.machine_state.revision != expected_previous_revision {
842 return Err(WorkGraphError::StaleRevision {
843 id: previous.work_ref.item_id.clone(),
844 expected: expected_previous_revision,
845 actual: current.machine_state.revision,
846 });
847 }
848 let replacement_key = attention_key(
849 &replacement.work_ref.realm_id,
850 &replacement.work_ref.namespace,
851 &replacement.binding_id,
852 );
853 if guard.attention.contains_key(&replacement_key) {
854 return Err(WorkGraphError::Conflict(format!(
855 "work attention binding {} already exists",
856 replacement.binding_id
857 )));
858 }
859 if let Some(occupant) = active_target_occupant_in(
862 guard
863 .attention
864 .values()
865 .filter(|binding| binding.binding_id != previous.binding_id),
866 &replacement,
867 ) {
868 return Err(active_target_conflict(&replacement, &occupant));
869 }
870 guard.attention.insert(previous_key, previous.clone());
871 guard.attention.insert(replacement_key, replacement.clone());
872 guard.append_event(previous_event);
873 guard.append_event(replacement_event);
874 Ok((previous, replacement))
875 }
876
877 async fn update_item_and_attention_cas(
878 &self,
879 item: WorkItem,
880 expected_previous_revision: u64,
881 item_event: WorkGraphEvent,
882 attention_updates: Vec<(WorkAttentionBinding, u64, WorkGraphEvent)>,
883 ) -> Result<WorkItem, WorkGraphError> {
884 WorkGraphMachine::validate_item_projection(&item)?;
885 let mut guard = self.inner.write().await;
886 let key = item_key(&item.realm_id, &item.namespace, &item.id);
887 let Some(current) = guard.items.get(&key) else {
888 return Err(WorkGraphError::not_found(
889 item.realm_id.clone(),
890 item.namespace.clone(),
891 item.id.clone(),
892 ));
893 };
894 if current.revision != expected_previous_revision {
895 return Err(WorkGraphError::StaleRevision {
896 id: item.id.clone(),
897 expected: expected_previous_revision,
898 actual: current.revision,
899 });
900 }
901 for (attention, expected_revision, _) in &attention_updates {
902 let key = attention_key(
903 &attention.work_ref.realm_id,
904 &attention.work_ref.namespace,
905 &attention.binding_id,
906 );
907 let Some(current) = guard.attention.get(&key) else {
908 return Err(WorkGraphError::not_found(
909 attention.work_ref.realm_id.clone(),
910 attention.work_ref.namespace.clone(),
911 attention.work_ref.item_id.clone(),
912 ));
913 };
914 if current.machine_state.revision != *expected_revision {
915 return Err(WorkGraphError::StaleRevision {
916 id: attention.work_ref.item_id.clone(),
917 expected: *expected_revision,
918 actual: current.machine_state.revision,
919 });
920 }
921 }
922 let batch_ids: Vec<WorkAttentionBindingId> = attention_updates
926 .iter()
927 .map(|(attention, _, _)| attention.binding_id.clone())
928 .collect();
929 for (index, (attention, _, _)) in attention_updates.iter().enumerate() {
930 let occupant = active_target_occupant_in(
931 guard
932 .attention
933 .values()
934 .filter(|binding| !batch_ids.contains(&binding.binding_id))
935 .chain(
936 attention_updates[..index]
937 .iter()
938 .map(|(applied, _, _)| applied),
939 ),
940 attention,
941 );
942 if let Some(occupant) = occupant {
943 return Err(active_target_conflict(attention, &occupant));
944 }
945 }
946 guard.items.insert(key, item.clone());
947 guard.append_event(item_event);
948 for (attention, _, event) in attention_updates {
949 let key = attention_key(
950 &attention.work_ref.realm_id,
951 &attention.work_ref.namespace,
952 &attention.binding_id,
953 );
954 guard.attention.insert(key, attention);
955 guard.append_event(event);
956 }
957 Ok(item)
958 }
959
960 async fn get_attention(
961 &self,
962 realm_id: &str,
963 namespace: &WorkNamespace,
964 binding_id: &WorkAttentionBindingId,
965 ) -> Result<Option<WorkAttentionBinding>, WorkGraphError> {
966 let guard = self.inner.read().await;
967 Ok(guard
968 .attention
969 .get(&attention_key(realm_id, namespace, binding_id))
970 .cloned())
971 }
972
973 async fn list_attention(
974 &self,
975 filter: AttentionListRequest,
976 ) -> Result<Vec<WorkAttentionBinding>, WorkGraphError> {
977 let guard = self.inner.read().await;
978 let mut bindings = guard
979 .attention
980 .values()
981 .filter(|binding| attention_matches_filter(binding, &filter))
982 .cloned()
983 .collect::<Vec<_>>();
984 bindings.sort_by(|left, right| {
985 left.updated_at
986 .cmp(&right.updated_at)
987 .then_with(|| left.binding_id.cmp(&right.binding_id))
988 });
989 Ok(bindings)
990 }
991
992 async fn list_attention_bounded(
993 &self,
994 filter: AttentionListRequest,
995 limit: usize,
996 ) -> Result<Vec<WorkAttentionBinding>, WorkGraphError> {
997 let guard = self.inner.read().await;
998 let compare = |left: &WorkAttentionBinding, right: &WorkAttentionBinding| {
999 left.updated_at
1000 .cmp(&right.updated_at)
1001 .then_with(|| left.binding_id.cmp(&right.binding_id))
1002 };
1003 let mut bindings = Vec::with_capacity(limit.min(1024));
1004 for binding in guard
1005 .attention
1006 .values()
1007 .filter(|binding| attention_matches_filter(binding, &filter))
1008 {
1009 let index = bindings
1010 .binary_search_by(|existing| compare(existing, binding))
1011 .unwrap_or_else(|index| index);
1012 if index < limit {
1013 bindings.insert(index, binding.clone());
1014 if bindings.len() > limit {
1015 bindings.pop();
1016 }
1017 }
1018 }
1019 Ok(bindings)
1020 }
1021
1022 async fn prune_terminal_attention(
1023 &self,
1024 filter: AttentionPruneRequest,
1025 ) -> Result<u64, WorkGraphError> {
1026 let mut guard = self.inner.write().await;
1027 let before = guard.attention.len();
1028 guard.attention.retain(|_, binding| {
1029 let in_scope = filter
1030 .realm_id
1031 .as_ref()
1032 .is_none_or(|realm_id| &binding.work_ref.realm_id == realm_id)
1033 && filter
1034 .namespace
1035 .as_ref()
1036 .is_none_or(|namespace| &binding.work_ref.namespace == namespace)
1037 && filter
1038 .updated_before
1039 .is_none_or(|updated_before| binding.updated_at < updated_before);
1040 !(in_scope && binding.status.is_terminal())
1041 });
1042 Ok((before - guard.attention.len()) as u64)
1043 }
1044
1045 async fn insert_edge(
1046 &self,
1047 edge: WorkEdge,
1048 event: WorkGraphEvent,
1049 ) -> Result<WorkEdge, WorkGraphError> {
1050 let mut guard = self.inner.write().await;
1051 if guard.edges.iter().any(|existing| existing == &edge) {
1052 return Err(duplicate_edge_error(&edge));
1053 }
1054 guard.edges.push(edge.clone());
1055 guard.append_event(event);
1056 Ok(edge)
1057 }
1058
1059 async fn insert_edge_validated(
1060 &self,
1061 edge: WorkEdge,
1062 event: WorkGraphEvent,
1063 ) -> Result<WorkEdge, WorkGraphError> {
1064 let mut guard = self.inner.write().await;
1065 if guard.edges.iter().any(|existing| existing == &edge) {
1066 return Err(duplicate_edge_error(&edge));
1067 }
1068 let existing_edges = guard
1069 .edges
1070 .iter()
1071 .filter(|existing| {
1072 existing.realm_id == edge.realm_id && existing.namespace == edge.namespace
1073 })
1074 .cloned()
1075 .collect::<Vec<_>>();
1076 let existing_items = guard
1077 .items
1078 .values()
1079 .filter(|item| item.realm_id == edge.realm_id && item.namespace == edge.namespace)
1080 .cloned()
1081 .collect::<Vec<_>>();
1082 WorkGraphMachine::validate_link(&edge, &existing_items, &existing_edges)?;
1083 guard.edges.push(edge.clone());
1084 guard.append_event(event);
1085 Ok(edge)
1086 }
1087
1088 async fn list_edges(
1089 &self,
1090 realm_id: &str,
1091 namespace: &WorkNamespace,
1092 ) -> Result<Vec<WorkEdge>, WorkGraphError> {
1093 let guard = self.inner.read().await;
1094 Ok(guard
1095 .edges
1096 .iter()
1097 .filter(|edge| edge.realm_id == realm_id && edge.namespace == *namespace)
1098 .cloned()
1099 .collect())
1100 }
1101
1102 async fn list_edges_bounded(
1103 &self,
1104 realm_id: &str,
1105 namespace: &WorkNamespace,
1106 limit: usize,
1107 ) -> Result<Vec<WorkEdge>, WorkGraphError> {
1108 let guard = self.inner.read().await;
1109 Ok(guard
1110 .edges
1111 .iter()
1112 .filter(|edge| edge.realm_id == realm_id && edge.namespace == *namespace)
1113 .take(limit)
1114 .cloned()
1115 .collect())
1116 }
1117
1118 async fn list_events(
1119 &self,
1120 filter: WorkGraphEventFilter,
1121 ) -> Result<Vec<WorkGraphEvent>, WorkGraphError> {
1122 let guard = self.inner.read().await;
1123 let events = guard
1124 .events
1125 .iter()
1126 .filter(|event| event_matches_filter(event, &filter))
1127 .take(filter.limit.unwrap_or(usize::MAX))
1128 .cloned()
1129 .collect::<Vec<_>>();
1130 Ok(events)
1131 }
1132
1133 async fn list_public_events(
1134 &self,
1135 filter: WorkGraphEventFilter,
1136 ) -> Result<Vec<WorkGraphEvent>, WorkGraphError> {
1137 let limit = filter.limit.unwrap_or(usize::MAX);
1138 if limit == 0 {
1139 return Ok(Vec::new());
1140 }
1141 let guard = self.inner.read().await;
1142 Ok(guard
1143 .events
1144 .iter()
1145 .filter(|event| event_matches_filter(event, &filter))
1146 .filter(|event| !is_internal_execution_event(event.kind))
1147 .take(limit)
1148 .cloned()
1149 .collect())
1150 }
1151
1152 async fn latest_event_seq(
1153 &self,
1154 filter: WorkGraphEventFilter,
1155 ) -> Result<Option<i64>, WorkGraphError> {
1156 let guard = self.inner.read().await;
1157 Ok(guard
1158 .events
1159 .iter()
1160 .filter(|event| event_matches_filter(event, &filter))
1161 .filter_map(|event| event.seq)
1162 .max())
1163 }
1164}
1165
1166impl MemoryWorkGraphState {
1167 fn append_event(&mut self, mut event: WorkGraphEvent) {
1168 self.next_event_seq += 1;
1169 event.seq = Some(self.next_event_seq);
1170 self.events.push(event);
1171 }
1172}
1173
1174fn item_key(
1175 realm_id: &str,
1176 namespace: &WorkNamespace,
1177 id: &WorkItemId,
1178) -> (String, WorkNamespace, WorkItemId) {
1179 (realm_id.to_string(), namespace.clone(), id.clone())
1180}
1181
1182fn attention_key(
1183 realm_id: &str,
1184 namespace: &WorkNamespace,
1185 id: &WorkAttentionBindingId,
1186) -> (String, WorkNamespace, WorkAttentionBindingId) {
1187 (realm_id.to_string(), namespace.clone(), id.clone())
1188}
1189
1190fn execution_binding_key(
1191 realm_id: &str,
1192 namespace: &WorkNamespace,
1193 id: &WorkExecutionBindingId,
1194) -> (String, WorkNamespace, WorkExecutionBindingId) {
1195 (realm_id.to_string(), namespace.clone(), id.clone())
1196}
1197
1198fn validate_execution_binding_insert<'a>(
1199 binding: &WorkExecutionBinding,
1200 expected_item_revision: u64,
1201 items: impl Iterator<Item = &'a WorkItem>,
1202 bindings: impl Iterator<Item = &'a WorkExecutionBinding>,
1203) -> Result<(), WorkGraphError> {
1204 let item = items
1205 .filter(|item| {
1206 item.realm_id == binding.work_ref.realm_id
1207 && item.namespace == binding.work_ref.namespace
1208 && item.id == binding.work_ref.item_id
1209 })
1210 .last()
1211 .ok_or_else(|| {
1212 WorkGraphError::not_found(
1213 binding.work_ref.realm_id.clone(),
1214 binding.work_ref.namespace.clone(),
1215 binding.work_ref.item_id.clone(),
1216 )
1217 })?;
1218 if item.revision != expected_item_revision {
1219 return Err(WorkGraphError::StaleRevision {
1220 id: item.id.clone(),
1221 expected: expected_item_revision,
1222 actual: item.revision,
1223 });
1224 }
1225 if WorkGraphMachine::classify_terminality(item)? {
1226 return Err(WorkGraphError::InvalidTransition(format!(
1227 "terminal work item {} cannot bind a new execution",
1228 item.id
1229 )));
1230 }
1231
1232 let bindings = bindings.collect::<Vec<_>>();
1233 if bindings
1234 .iter()
1235 .any(|existing| existing.target.run_id() == binding.target.run_id())
1236 {
1237 return Err(WorkGraphError::Conflict(format!(
1238 "work execution binding {} reuses target run id {}",
1239 binding.binding_id,
1240 binding.target.run_id()
1241 )));
1242 }
1243 let scoped = bindings
1244 .into_iter()
1245 .filter(|existing| {
1246 existing.work_ref.realm_id == binding.work_ref.realm_id
1247 && existing.work_ref.namespace == binding.work_ref.namespace
1248 && existing.work_ref.item_id == binding.work_ref.item_id
1249 })
1250 .collect::<Vec<_>>();
1251 if scoped.iter().any(|existing| {
1252 existing.idempotency_key == binding.idempotency_key
1253 || existing.target.run_id() == binding.target.run_id()
1254 }) {
1255 return Err(WorkGraphError::Conflict(format!(
1256 "work execution binding {} reuses an idempotency key or run id",
1257 binding.binding_id
1258 )));
1259 }
1260
1261 match &binding.supersedes {
1262 None if !scoped.is_empty() => Err(WorkGraphError::Conflict(format!(
1263 "work item {} already has an execution chain",
1264 binding.work_ref.item_id
1265 ))),
1266 None => Ok(()),
1267 Some(predecessor) => {
1268 if predecessor == &binding.binding_id {
1269 return Err(WorkGraphError::InvalidInput(
1270 "work execution binding cannot supersede itself".to_string(),
1271 ));
1272 }
1273 let Some(predecessor_binding) = scoped
1274 .iter()
1275 .copied()
1276 .find(|existing| &existing.binding_id == predecessor)
1277 else {
1278 return Err(WorkGraphError::InvalidInput(format!(
1279 "superseded work execution binding {predecessor} is not in the same work item chain"
1280 )));
1281 };
1282 if !crate::WorkExecutionMachine::retry_eligible(predecessor_binding)? {
1283 return Err(WorkGraphError::InvalidTransition(format!(
1284 "work execution binding {predecessor} is not terminal and cannot be superseded"
1285 )));
1286 }
1287 if scoped
1288 .iter()
1289 .any(|existing| existing.supersedes.as_ref() == Some(predecessor))
1290 {
1291 return Err(WorkGraphError::Conflict(format!(
1292 "work execution binding {predecessor} is already superseded"
1293 )));
1294 }
1295 Ok(())
1296 }
1297 }
1298}
1299
1300fn execution_binding_matches_filter(
1301 binding: &WorkExecutionBinding,
1302 filter: &WorkExecutionBindingFilter,
1303 superseded: &std::collections::BTreeSet<(String, WorkNamespace, WorkExecutionBindingId)>,
1304) -> bool {
1305 filter
1306 .realm_id
1307 .as_ref()
1308 .is_none_or(|realm_id| &binding.work_ref.realm_id == realm_id)
1309 && filter
1310 .namespace
1311 .as_ref()
1312 .is_none_or(|namespace| &binding.work_ref.namespace == namespace)
1313 && filter
1314 .item_id
1315 .as_ref()
1316 .is_none_or(|item_id| &binding.work_ref.item_id == item_id)
1317 && (!filter.current_only
1318 || !superseded.contains(&(
1319 binding.work_ref.realm_id.clone(),
1320 binding.work_ref.namespace.clone(),
1321 binding.binding_id.clone(),
1322 )))
1323}
1324
1325fn item_matches_filter(item: &WorkItem, filter: &WorkItemFilter) -> bool {
1326 if let Some(realm_id) = &filter.realm_id
1327 && &item.realm_id != realm_id
1328 {
1329 return false;
1330 }
1331 if !filter.all_namespaces
1332 && let Some(namespace) = &filter.namespace
1333 && &item.namespace != namespace
1334 {
1335 return false;
1336 }
1337 if !filter.statuses.is_empty() && !filter.statuses.contains(&item.status) {
1338 return false;
1339 }
1340 if !filter.include_terminal && WorkGraphMachine::classify_terminality(item).unwrap_or(true) {
1346 return false;
1347 }
1348 filter
1349 .labels
1350 .iter()
1351 .all(|label| item.labels.contains(label))
1352}
1353
1354fn attention_matches_filter(binding: &WorkAttentionBinding, filter: &AttentionListRequest) -> bool {
1355 if let Some(realm_id) = &filter.realm_id
1356 && &binding.work_ref.realm_id != realm_id
1357 {
1358 return false;
1359 }
1360 if let Some(namespace) = &filter.namespace
1361 && &binding.work_ref.namespace != namespace
1362 {
1363 return false;
1364 }
1365 if let Some(target) = &filter.target
1366 && &binding.target != target
1367 {
1368 return false;
1369 }
1370 if let Some(status) = &filter.status
1371 && !attention_status_matches_filter(&binding.status, status)
1372 {
1373 return false;
1374 }
1375 true
1376}
1377
1378fn attention_status_matches_filter(
1379 actual: &crate::types::WorkAttentionStatus,
1380 filter: &crate::types::WorkAttentionStatus,
1381) -> bool {
1382 use crate::types::WorkAttentionStatus;
1383
1384 match (actual, filter) {
1385 (WorkAttentionStatus::Active, WorkAttentionStatus::Active)
1386 | (WorkAttentionStatus::Superseded, WorkAttentionStatus::Superseded)
1387 | (WorkAttentionStatus::Stopped, WorkAttentionStatus::Stopped) => true,
1388 (WorkAttentionStatus::Paused { .. }, WorkAttentionStatus::Paused { until: None }) => true,
1389 (
1390 WorkAttentionStatus::Paused {
1391 until: Some(actual_until),
1392 },
1393 WorkAttentionStatus::Paused {
1394 until: Some(filter_until),
1395 },
1396 ) => actual_until == filter_until,
1397 _ => false,
1398 }
1399}
1400
1401fn event_matches_filter(event: &WorkGraphEvent, filter: &WorkGraphEventFilter) -> bool {
1402 if let Some(after_seq) = filter.after_seq
1403 && event.seq.unwrap_or_default() <= after_seq
1404 {
1405 return false;
1406 }
1407 if let Some(realm_id) = &filter.realm_id
1408 && &event.realm_id != realm_id
1409 {
1410 return false;
1411 }
1412 if !filter.all_namespaces
1413 && let Some(namespace) = &filter.namespace
1414 && &event.namespace != namespace
1415 {
1416 return false;
1417 }
1418 true
1419}
1420
1421fn is_internal_execution_event(kind: WorkGraphEventKind) -> bool {
1422 matches!(
1423 kind,
1424 WorkGraphEventKind::ExecutionBound | WorkGraphEventKind::ExecutionTransitioned
1425 )
1426}
1427
1428#[cfg(not(target_arch = "wasm32"))]
1429pub struct SqliteWorkGraphStore {
1430 path: PathBuf,
1431}
1432
1433#[cfg(not(target_arch = "wasm32"))]
1434impl SqliteWorkGraphStore {
1435 pub fn open(path: impl Into<PathBuf>) -> Result<Self, WorkGraphError> {
1436 let store = Self { path: path.into() };
1437 store.with_connection(|_| Ok(()))?;
1439 Ok(store)
1440 }
1441
1442 pub fn path(&self) -> &Path {
1443 &self.path
1444 }
1445
1446 pub fn rebuild_projection_from_events(&self) -> Result<(), WorkGraphError> {
1447 self.with_connection(|conn| {
1448 let tx = conn
1452 .transaction_with_behavior(TransactionBehavior::Immediate)
1453 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1454 tx.execute("DELETE FROM workgraph_items", [])
1455 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1456 tx.execute("DELETE FROM workgraph_edges", [])
1457 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1458 tx.execute("DELETE FROM workgraph_attention", [])
1459 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1460 tx.execute("DELETE FROM workgraph_execution_bindings", [])
1461 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1462
1463 let events = {
1464 let mut stmt = tx
1465 .prepare("SELECT event_json FROM workgraph_events ORDER BY seq ASC")
1466 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1467 let rows = stmt
1468 .query_map([], |row| row_json::<WorkGraphEvent>(row, 0))
1469 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1470 let mut events = Vec::new();
1471 for row in rows {
1472 events.push(row.map_err(|err| WorkGraphError::Store(err.to_string()))?);
1473 }
1474 events
1475 };
1476
1477 for event in events {
1478 replay_event_tx(&tx, &event)?;
1479 }
1480 normalize_attention_for_terminal_items_tx(&tx)?;
1481 tx.commit()
1482 .map_err(|err| WorkGraphError::Store(err.to_string()))
1483 })
1484 }
1485
1486 fn with_connection<T>(
1487 &self,
1488 f: impl FnOnce(&mut Connection) -> Result<T, WorkGraphError>,
1489 ) -> Result<T, WorkGraphError> {
1490 let _guard = meerkat_sqlite::OperationGuard::for_database(&self.path)
1493 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1494 let mut conn = meerkat_sqlite::open_with(
1495 &self.path,
1496 meerkat_sqlite::ConnectionProfile::PRIMARY,
1497 meerkat_sqlite::OpenOptions {
1498 schema_preflight: &[&WORKGRAPH_DOMAIN],
1499 ..Default::default()
1500 },
1501 )
1502 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1503 meerkat_sqlite::apply_domain_migrations(&mut conn, &WORKGRAPH_DOMAIN)
1504 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1505 f(&mut conn)
1506 }
1507}
1508
1509#[cfg(not(target_arch = "wasm32"))]
1510#[async_trait]
1511impl WorkGraphStore for SqliteWorkGraphStore {
1512 fn kind(&self) -> WorkGraphStoreKind {
1513 WorkGraphStoreKind::Sqlite
1514 }
1515
1516 async fn get_store_time_utc(&self) -> Result<DateTime<Utc>, WorkGraphError> {
1517 Ok(Utc::now())
1518 }
1519
1520 async fn insert_item(
1521 &self,
1522 item: WorkItem,
1523 event: WorkGraphEvent,
1524 ) -> Result<WorkItem, WorkGraphError> {
1525 WorkGraphMachine::validate_item_projection(&item)?;
1526 self.with_connection(|conn| {
1527 let tx = conn
1528 .transaction_with_behavior(TransactionBehavior::Immediate)
1529 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1530 insert_item_tx(&tx, &item)?;
1531 insert_event_tx(&tx, &event)?;
1532 tx.commit()
1533 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1534 Ok(item)
1535 })
1536 }
1537
1538 async fn update_item_cas(
1539 &self,
1540 item: WorkItem,
1541 expected_previous_revision: u64,
1542 event: WorkGraphEvent,
1543 ) -> Result<WorkItem, WorkGraphError> {
1544 WorkGraphMachine::validate_item_projection(&item)?;
1545 self.with_connection(|conn| {
1546 let tx = conn
1547 .transaction_with_behavior(TransactionBehavior::Immediate)
1548 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1549 let changed = update_item_tx(&tx, &item, expected_previous_revision)?;
1550 if changed == 0 {
1551 let actual = current_revision_tx(&tx, &item.realm_id, &item.namespace, &item.id)?;
1552 return match actual {
1553 Some(actual) => Err(WorkGraphError::StaleRevision {
1554 id: item.id,
1555 expected: expected_previous_revision,
1556 actual,
1557 }),
1558 None => Err(WorkGraphError::not_found(
1559 item.realm_id,
1560 item.namespace,
1561 item.id,
1562 )),
1563 };
1564 }
1565 insert_event_tx(&tx, &event)?;
1566 tx.commit()
1567 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1568 Ok(item)
1569 })
1570 }
1571
1572 async fn get_item(
1573 &self,
1574 realm_id: &str,
1575 namespace: &WorkNamespace,
1576 id: &WorkItemId,
1577 ) -> Result<Option<WorkItem>, WorkGraphError> {
1578 self.with_connection(|conn| select_item(conn, realm_id, namespace, id))
1579 }
1580
1581 async fn list_items(&self, filter: WorkItemFilter) -> Result<Vec<WorkItem>, WorkGraphError> {
1582 self.with_connection(|conn| list_sqlite_items(conn, &filter))
1583 }
1584
1585 async fn insert_execution_binding(
1586 &self,
1587 commit: crate::WorkExecutionBindCommit,
1588 expected_item_revision: u64,
1589 event: WorkGraphEvent,
1590 ) -> Result<WorkExecutionBinding, WorkGraphError> {
1591 let (binding, effect) = commit.into_parts();
1592 binding.validate()?;
1593 crate::WorkExecutionMachine::validate_projection(&binding)?;
1594 let (expected_state, expected_effect) =
1595 crate::WorkExecutionMachine::bind(&binding.binding_id, binding.target.run_id())?;
1596 if binding.machine_state != expected_state || effect != expected_effect {
1597 return Err(WorkGraphError::InvalidInput(format!(
1598 "work execution binding {} lacks canonical bind authority",
1599 binding.binding_id
1600 )));
1601 }
1602 self.with_connection(|conn| {
1603 let tx = conn
1604 .transaction_with_behavior(TransactionBehavior::Immediate)
1605 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1606 if let Some(existing) = select_execution_binding(
1607 &tx,
1608 &binding.work_ref.realm_id,
1609 &binding.work_ref.namespace,
1610 &binding.binding_id,
1611 )? {
1612 return if existing == binding {
1613 Ok(existing)
1614 } else {
1615 Err(WorkGraphError::Conflict(format!(
1616 "work execution binding {} already exists with different content",
1617 binding.binding_id
1618 )))
1619 };
1620 }
1621 let items = select_item(
1622 &tx,
1623 &binding.work_ref.realm_id,
1624 &binding.work_ref.namespace,
1625 &binding.work_ref.item_id,
1626 )?
1627 .into_iter()
1628 .collect::<Vec<_>>();
1629 let bindings = list_sqlite_execution_bindings(
1630 &tx,
1631 &WorkExecutionBindingFilter {
1632 realm_id: Some(binding.work_ref.realm_id.clone()),
1633 namespace: Some(binding.work_ref.namespace.clone()),
1634 item_id: Some(binding.work_ref.item_id.clone()),
1635 current_only: false,
1636 limit: None,
1637 },
1638 )?;
1639 validate_execution_binding_insert(
1640 &binding,
1641 expected_item_revision,
1642 items.iter(),
1643 bindings.iter(),
1644 )?;
1645 insert_execution_binding_tx(&tx, &binding)?;
1646 insert_event_tx(&tx, &event)?;
1647 tx.commit()
1648 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1649 Ok(binding)
1650 })
1651 }
1652
1653 async fn get_execution_binding(
1654 &self,
1655 realm_id: &str,
1656 namespace: &WorkNamespace,
1657 binding_id: &WorkExecutionBindingId,
1658 ) -> Result<Option<WorkExecutionBinding>, WorkGraphError> {
1659 self.with_connection(|conn| select_execution_binding(conn, realm_id, namespace, binding_id))
1660 }
1661
1662 async fn get_execution_binding_by_target_run(
1663 &self,
1664 realm_id: &str,
1665 run_id: &str,
1666 ) -> Result<Option<WorkExecutionBinding>, WorkGraphError> {
1667 self.with_connection(|conn| {
1668 conn.query_row(
1669 "SELECT binding_json FROM workgraph_execution_bindings
1670 WHERE realm_id = ?1 AND target_run_id = ?2",
1671 params![realm_id, run_id],
1672 |row| row_json(row, 0),
1673 )
1674 .optional()
1675 .map_err(|error| WorkGraphError::Store(error.to_string()))
1676 })
1677 }
1678
1679 async fn update_execution_binding_cas(
1680 &self,
1681 commit: crate::WorkExecutionObservationCommit,
1682 expected_previous_revision: u64,
1683 event: WorkGraphEvent,
1684 ) -> Result<WorkExecutionBinding, WorkGraphError> {
1685 let (previous, observation, binding, effect) = commit.into_parts();
1686 crate::WorkExecutionMachine::validate_projection(&binding)?;
1687 self.with_connection(|conn| {
1688 let tx = conn
1689 .transaction_with_behavior(TransactionBehavior::Immediate)
1690 .map_err(|error| WorkGraphError::Store(error.to_string()))?;
1691 let current = select_execution_binding(
1692 &tx,
1693 &binding.work_ref.realm_id,
1694 &binding.work_ref.namespace,
1695 &binding.binding_id,
1696 )?
1697 .ok_or_else(|| {
1698 WorkGraphError::Conflict(format!(
1699 "work execution binding {} does not exist",
1700 binding.binding_id
1701 ))
1702 })?;
1703 if current.machine_state.revision != expected_previous_revision {
1704 return Err(WorkGraphError::Conflict(format!(
1705 "stale work execution revision for {}: expected {}, actual {}",
1706 binding.binding_id, expected_previous_revision, current.machine_state.revision
1707 )));
1708 }
1709 if current != previous {
1710 return Err(WorkGraphError::Conflict(format!(
1711 "work execution transition authority for {} was minted from a different predecessor",
1712 binding.binding_id
1713 )));
1714 }
1715 let (expected_binding, expected_effect) = crate::WorkExecutionMachine::observe(
1716 current.clone(),
1717 expected_previous_revision,
1718 observation,
1719 )?;
1720 if binding != expected_binding || effect != expected_effect {
1721 return Err(WorkGraphError::Conflict(format!(
1722 "work execution transition authority for {} does not match the generated machine result",
1723 binding.binding_id
1724 )));
1725 }
1726 if !current.has_same_immutable_spec(&binding) {
1727 return Err(WorkGraphError::Conflict(format!(
1728 "immutable work execution specification changed for {}",
1729 binding.binding_id
1730 )));
1731 }
1732 let json = serde_json::to_string(&binding)
1733 .map_err(|error| WorkGraphError::Store(error.to_string()))?;
1734 let recovery_pending = execution_recovery_pending(&binding)?;
1735 let changed = tx
1736 .execute(
1737 "UPDATE workgraph_execution_bindings
1738 SET revision = ?1, recovery_pending = ?2, binding_json = ?3
1739 WHERE realm_id = ?4 AND namespace = ?5 AND binding_id = ?6
1740 AND revision = ?7",
1741 params![
1742 binding.machine_state.revision,
1743 recovery_pending,
1744 json,
1745 binding.work_ref.realm_id,
1746 binding.work_ref.namespace.as_str(),
1747 binding.binding_id.as_str(),
1748 expected_previous_revision,
1749 ],
1750 )
1751 .map_err(|error| WorkGraphError::Store(error.to_string()))?;
1752 if changed == 0 {
1753 let current = select_execution_binding(
1754 &tx,
1755 &binding.work_ref.realm_id,
1756 &binding.work_ref.namespace,
1757 &binding.binding_id,
1758 )?;
1759 return match current {
1760 Some(current) => Err(WorkGraphError::Conflict(format!(
1761 "stale work execution revision for {}: expected {}, actual {}",
1762 binding.binding_id,
1763 expected_previous_revision,
1764 current.machine_state.revision
1765 ))),
1766 None => Err(WorkGraphError::Conflict(format!(
1767 "work execution binding {} does not exist",
1768 binding.binding_id
1769 ))),
1770 };
1771 }
1772 insert_event_tx(&tx, &event)?;
1773 tx.commit()
1774 .map_err(|error| WorkGraphError::Store(error.to_string()))?;
1775 Ok(binding)
1776 })
1777 }
1778
1779 async fn list_execution_bindings(
1780 &self,
1781 filter: WorkExecutionBindingFilter,
1782 ) -> Result<Vec<WorkExecutionBinding>, WorkGraphError> {
1783 self.with_connection(|conn| list_sqlite_execution_bindings(conn, &filter))
1784 }
1785
1786 async fn list_execution_bindings_for_recovery(
1787 &self,
1788 realm_id: &str,
1789 ) -> Result<Vec<WorkExecutionBinding>, WorkGraphError> {
1790 self.with_connection(|conn| {
1791 let mut statement = conn
1792 .prepare(
1793 "SELECT binding_json FROM workgraph_execution_bindings
1794 WHERE realm_id = ?1 AND recovery_pending = 1
1795 ORDER BY created_at_utc ASC, binding_id ASC",
1796 )
1797 .map_err(|error| WorkGraphError::Store(error.to_string()))?;
1798 let rows = statement
1799 .query_map([realm_id], |row| row_json::<WorkExecutionBinding>(row, 0))
1800 .map_err(|error| WorkGraphError::Store(error.to_string()))?;
1801 rows.map(|row| row.map_err(|error| WorkGraphError::Store(error.to_string())))
1802 .collect()
1803 })
1804 }
1805
1806 async fn insert_goal(
1807 &self,
1808 item: WorkItem,
1809 item_event: WorkGraphEvent,
1810 attention: WorkAttentionBinding,
1811 attention_event: WorkGraphEvent,
1812 ) -> Result<(WorkItem, WorkAttentionBinding), WorkGraphError> {
1813 WorkGraphMachine::validate_item_projection(&item)?;
1814 self.with_connection(|conn| {
1815 let tx = conn
1816 .transaction_with_behavior(TransactionBehavior::Immediate)
1817 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1818 if let Some(occupant) = active_target_occupant_tx(&tx, &attention)? {
1819 return Err(active_target_conflict(&attention, &occupant));
1820 }
1821 insert_item_tx(&tx, &item)?;
1822 insert_attention_tx(&tx, &attention)?;
1823 insert_event_tx(&tx, &item_event)?;
1824 insert_event_tx(&tx, &attention_event)?;
1825 tx.commit()
1826 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1827 Ok((item, attention))
1828 })
1829 }
1830
1831 async fn update_attention_cas(
1832 &self,
1833 attention: WorkAttentionBinding,
1834 expected_previous_revision: u64,
1835 event: WorkGraphEvent,
1836 ) -> Result<WorkAttentionBinding, WorkGraphError> {
1837 self.with_connection(|conn| {
1838 let tx = conn
1839 .transaction_with_behavior(TransactionBehavior::Immediate)
1840 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1841 let changed = update_attention_tx(&tx, &attention, expected_previous_revision)?;
1842 if changed == 0 {
1843 let actual = current_attention_revision_tx(
1844 &tx,
1845 &attention.work_ref.realm_id,
1846 &attention.work_ref.namespace,
1847 &attention.binding_id,
1848 )?;
1849 return match actual {
1850 Some(actual) => Err(WorkGraphError::StaleRevision {
1851 id: attention.work_ref.item_id,
1852 expected: expected_previous_revision,
1853 actual,
1854 }),
1855 None => Err(WorkGraphError::not_found(
1856 attention.work_ref.realm_id,
1857 attention.work_ref.namespace,
1858 attention.work_ref.item_id,
1859 )),
1860 };
1861 }
1862 if let Some(occupant) = active_target_occupant_tx(&tx, &attention)? {
1866 return Err(active_target_conflict(&attention, &occupant));
1867 }
1868 insert_event_tx(&tx, &event)?;
1869 tx.commit()
1870 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1871 Ok(attention)
1872 })
1873 }
1874
1875 async fn reassign_attention_cas(
1876 &self,
1877 previous: WorkAttentionBinding,
1878 expected_previous_revision: u64,
1879 previous_event: WorkGraphEvent,
1880 replacement: WorkAttentionBinding,
1881 replacement_event: WorkGraphEvent,
1882 ) -> Result<(WorkAttentionBinding, WorkAttentionBinding), WorkGraphError> {
1883 self.with_connection(|conn| {
1884 let tx = conn
1885 .transaction_with_behavior(TransactionBehavior::Immediate)
1886 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1887 let changed = update_attention_tx(&tx, &previous, expected_previous_revision)?;
1888 if changed == 0 {
1889 let actual = current_attention_revision_tx(
1890 &tx,
1891 &previous.work_ref.realm_id,
1892 &previous.work_ref.namespace,
1893 &previous.binding_id,
1894 )?;
1895 return match actual {
1896 Some(actual) => Err(WorkGraphError::StaleRevision {
1897 id: previous.work_ref.item_id,
1898 expected: expected_previous_revision,
1899 actual,
1900 }),
1901 None => Err(WorkGraphError::attention_not_found(
1902 previous.work_ref.realm_id,
1903 previous.work_ref.namespace,
1904 previous.binding_id,
1905 )),
1906 };
1907 }
1908 if let Some(occupant) = active_target_occupant_tx(&tx, &replacement)? {
1912 return Err(active_target_conflict(&replacement, &occupant));
1913 }
1914 insert_attention_tx(&tx, &replacement)?;
1915 insert_event_tx(&tx, &previous_event)?;
1916 insert_event_tx(&tx, &replacement_event)?;
1917 tx.commit()
1918 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1919 Ok((previous, replacement))
1920 })
1921 }
1922
1923 async fn update_item_and_attention_cas(
1924 &self,
1925 item: WorkItem,
1926 expected_previous_revision: u64,
1927 item_event: WorkGraphEvent,
1928 attention_updates: Vec<(WorkAttentionBinding, u64, WorkGraphEvent)>,
1929 ) -> Result<WorkItem, WorkGraphError> {
1930 WorkGraphMachine::validate_item_projection(&item)?;
1931 self.with_connection(|conn| {
1932 let tx = conn
1933 .transaction_with_behavior(TransactionBehavior::Immediate)
1934 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1935 let changed = update_item_tx(&tx, &item, expected_previous_revision)?;
1936 if changed == 0 {
1937 let actual = current_revision_tx(&tx, &item.realm_id, &item.namespace, &item.id)?;
1938 return match actual {
1939 Some(actual) => Err(WorkGraphError::StaleRevision {
1940 id: item.id,
1941 expected: expected_previous_revision,
1942 actual,
1943 }),
1944 None => Err(WorkGraphError::not_found(
1945 item.realm_id,
1946 item.namespace,
1947 item.id,
1948 )),
1949 };
1950 }
1951 insert_event_tx(&tx, &item_event)?;
1952 for (attention, expected_revision, event) in &attention_updates {
1953 let changed = update_attention_tx(&tx, attention, *expected_revision)?;
1954 if changed == 0 {
1955 let actual = current_attention_revision_tx(
1956 &tx,
1957 &attention.work_ref.realm_id,
1958 &attention.work_ref.namespace,
1959 &attention.binding_id,
1960 )?;
1961 return match actual {
1962 Some(actual) => Err(WorkGraphError::StaleRevision {
1963 id: attention.work_ref.item_id.clone(),
1964 expected: *expected_revision,
1965 actual,
1966 }),
1967 None => Err(WorkGraphError::not_found(
1968 attention.work_ref.realm_id.clone(),
1969 attention.work_ref.namespace.clone(),
1970 attention.work_ref.item_id.clone(),
1971 )),
1972 };
1973 }
1974 if let Some(occupant) = active_target_occupant_tx(&tx, attention)? {
1977 return Err(active_target_conflict(attention, &occupant));
1978 }
1979 insert_event_tx(&tx, event)?;
1980 }
1981 tx.commit()
1982 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
1983 Ok(item)
1984 })
1985 }
1986
1987 async fn get_attention(
1988 &self,
1989 realm_id: &str,
1990 namespace: &WorkNamespace,
1991 binding_id: &WorkAttentionBindingId,
1992 ) -> Result<Option<WorkAttentionBinding>, WorkGraphError> {
1993 self.with_connection(|conn| select_attention(conn, realm_id, namespace, binding_id))
1994 }
1995
1996 async fn list_attention(
1997 &self,
1998 filter: AttentionListRequest,
1999 ) -> Result<Vec<WorkAttentionBinding>, WorkGraphError> {
2000 self.with_connection(|conn| list_sqlite_attention(conn, &filter, None))
2001 }
2002
2003 async fn list_attention_bounded(
2004 &self,
2005 filter: AttentionListRequest,
2006 limit: usize,
2007 ) -> Result<Vec<WorkAttentionBinding>, WorkGraphError> {
2008 self.with_connection(|conn| list_sqlite_attention(conn, &filter, Some(limit)))
2009 }
2010
2011 async fn prune_terminal_attention(
2012 &self,
2013 filter: AttentionPruneRequest,
2014 ) -> Result<u64, WorkGraphError> {
2015 self.with_connection(|conn| {
2016 let tx = conn
2017 .transaction_with_behavior(TransactionBehavior::Immediate)
2018 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2019 let candidates: Vec<(String, String, String)> = {
2023 let mut stmt = tx
2024 .prepare(
2025 "SELECT realm_id, namespace, binding_id, attention_json
2026 FROM workgraph_attention
2027 WHERE status IN ('superseded', 'stopped') OR status IS NULL",
2028 )
2029 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2030 let rows = stmt
2031 .query_map([], |row| {
2032 Ok((
2033 row.get::<_, String>(0)?,
2034 row.get::<_, String>(1)?,
2035 row.get::<_, String>(2)?,
2036 row_json::<WorkAttentionBinding>(row, 3)?,
2037 ))
2038 })
2039 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2040 let mut candidates = Vec::new();
2041 for row in rows {
2042 let (realm_id, namespace, binding_id, binding) =
2043 row.map_err(|err| WorkGraphError::Store(err.to_string()))?;
2044 let in_scope = filter
2045 .realm_id
2046 .as_ref()
2047 .is_none_or(|realm| &binding.work_ref.realm_id == realm)
2048 && filter
2049 .namespace
2050 .as_ref()
2051 .is_none_or(|ns| &binding.work_ref.namespace == ns)
2052 && filter
2053 .updated_before
2054 .is_none_or(|updated_before| binding.updated_at < updated_before);
2055 if in_scope && binding.status.is_terminal() {
2056 candidates.push((realm_id, namespace, binding_id));
2057 }
2058 }
2059 candidates
2060 };
2061 let mut pruned = 0u64;
2062 for (realm_id, namespace, binding_id) in candidates {
2063 pruned += tx
2064 .execute(
2065 "DELETE FROM workgraph_attention
2066 WHERE realm_id = ?1 AND namespace = ?2 AND binding_id = ?3",
2067 params![realm_id, namespace, binding_id],
2068 )
2069 .map_err(|err| WorkGraphError::Store(err.to_string()))?
2070 as u64;
2071 }
2072 tx.commit()
2073 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2074 Ok(pruned)
2075 })
2076 }
2077
2078 async fn insert_edge(
2079 &self,
2080 edge: WorkEdge,
2081 event: WorkGraphEvent,
2082 ) -> Result<WorkEdge, WorkGraphError> {
2083 self.with_connection(|conn| {
2084 let tx = conn
2085 .transaction_with_behavior(TransactionBehavior::Immediate)
2086 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2087 insert_edge_tx(&tx, &edge)?;
2088 insert_event_tx(&tx, &event)?;
2089 tx.commit()
2090 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2091 Ok(edge)
2092 })
2093 }
2094
2095 async fn insert_edge_validated(
2096 &self,
2097 edge: WorkEdge,
2098 event: WorkGraphEvent,
2099 ) -> Result<WorkEdge, WorkGraphError> {
2100 self.with_connection(|conn| {
2101 let tx = conn
2102 .transaction_with_behavior(TransactionBehavior::Immediate)
2103 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2104 let existing_edges = list_sqlite_edges(&tx, &edge.realm_id, &edge.namespace, None)?;
2105 let existing_items = list_sqlite_items(
2106 &tx,
2107 &WorkItemFilter {
2108 realm_id: Some(edge.realm_id.clone()),
2109 namespace: Some(edge.namespace.clone()),
2110 include_terminal: true,
2111 ..WorkItemFilter::default()
2112 },
2113 )?;
2114 WorkGraphMachine::validate_link(&edge, &existing_items, &existing_edges)?;
2115 insert_edge_tx(&tx, &edge)?;
2116 insert_event_tx(&tx, &event)?;
2117 tx.commit()
2118 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2119 Ok(edge)
2120 })
2121 }
2122
2123 async fn list_edges(
2124 &self,
2125 realm_id: &str,
2126 namespace: &WorkNamespace,
2127 ) -> Result<Vec<WorkEdge>, WorkGraphError> {
2128 self.with_connection(|conn| list_sqlite_edges(conn, realm_id, namespace, None))
2129 }
2130
2131 async fn list_edges_bounded(
2132 &self,
2133 realm_id: &str,
2134 namespace: &WorkNamespace,
2135 limit: usize,
2136 ) -> Result<Vec<WorkEdge>, WorkGraphError> {
2137 self.with_connection(|conn| list_sqlite_edges(conn, realm_id, namespace, Some(limit)))
2138 }
2139
2140 async fn list_events(
2141 &self,
2142 filter: WorkGraphEventFilter,
2143 ) -> Result<Vec<WorkGraphEvent>, WorkGraphError> {
2144 self.with_connection(|conn| list_sqlite_events(conn, &filter))
2145 }
2146
2147 async fn list_public_events(
2148 &self,
2149 filter: WorkGraphEventFilter,
2150 ) -> Result<Vec<WorkGraphEvent>, WorkGraphError> {
2151 self.with_connection(|conn| list_sqlite_public_events(conn, &filter))
2152 }
2153
2154 async fn latest_event_seq(
2155 &self,
2156 filter: WorkGraphEventFilter,
2157 ) -> Result<Option<i64>, WorkGraphError> {
2158 self.with_connection(|conn| latest_sqlite_event_seq(conn, &filter))
2159 }
2160}
2161
2162#[cfg(not(target_arch = "wasm32"))]
2163fn build_released_0_8_15_workgraph_schema(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2164 migration_0001_workgraph_schema(tx)?;
2165 migration_0002_attention_query_columns(tx)
2166}
2167
2168#[cfg(not(target_arch = "wasm32"))]
2169const RELEASED_0_8_15_WORKGRAPH_OBJECTS: &[meerkat_sqlite::SchemaObject] = &[
2170 meerkat_sqlite::SchemaObject {
2171 kind: meerkat_sqlite::SchemaObjectKind::Table,
2172 name: "workgraph_items",
2173 },
2174 meerkat_sqlite::SchemaObject {
2175 kind: meerkat_sqlite::SchemaObjectKind::Index,
2176 name: "idx_workgraph_items_realm_namespace_updated",
2177 },
2178 meerkat_sqlite::SchemaObject {
2179 kind: meerkat_sqlite::SchemaObjectKind::Table,
2180 name: "workgraph_attention",
2181 },
2182 meerkat_sqlite::SchemaObject {
2183 kind: meerkat_sqlite::SchemaObjectKind::Index,
2184 name: "idx_workgraph_attention_realm_namespace_updated",
2185 },
2186 meerkat_sqlite::SchemaObject {
2187 kind: meerkat_sqlite::SchemaObjectKind::Index,
2188 name: "idx_workgraph_attention_scope_status",
2189 },
2190 meerkat_sqlite::SchemaObject {
2191 kind: meerkat_sqlite::SchemaObjectKind::Table,
2192 name: "workgraph_edges",
2193 },
2194 meerkat_sqlite::SchemaObject {
2195 kind: meerkat_sqlite::SchemaObjectKind::Table,
2196 name: "workgraph_events",
2197 },
2198 meerkat_sqlite::SchemaObject {
2199 kind: meerkat_sqlite::SchemaObjectKind::Index,
2200 name: "idx_workgraph_events_realm_namespace_seq",
2201 },
2202];
2203
2204#[cfg(not(target_arch = "wasm32"))]
2205fn verify_released_0_8_15_workgraph_schema(conn: &Connection) -> Result<(), String> {
2206 meerkat_sqlite::verify_released_schema_fingerprint(
2207 conn,
2208 &WORKGRAPH_DOMAIN,
2209 RELEASED_0_8_15_WORKGRAPH_OBJECTS,
2210 build_released_0_8_15_workgraph_schema,
2211 )
2212}
2213
2214#[cfg(not(target_arch = "wasm32"))]
2215#[cfg(not(target_arch = "wasm32"))]
2222pub const WORKGRAPH_DOMAIN: meerkat_sqlite::SchemaDomain = meerkat_sqlite::SchemaDomain {
2223 name: "workgraph",
2224 migrations: &[
2225 meerkat_sqlite::Migration {
2226 version: 1,
2227 name: "base-schema",
2228 apply: migration_0001_workgraph_schema,
2229 },
2230 meerkat_sqlite::Migration {
2231 version: 2,
2232 name: "attention-query-columns",
2233 apply: migration_0002_attention_query_columns,
2234 },
2235 meerkat_sqlite::Migration {
2236 version: 3,
2237 name: "execution-bindings",
2238 apply: migration_0003_execution_bindings,
2239 },
2240 ],
2241 initialize_current: initialize_current_workgraph_schema,
2242 allowed_existing_versions: &[2, 3],
2243 released_predecessors: &[meerkat_sqlite::SchemaPredecessor {
2244 version: 2,
2245 verify: verify_released_0_8_15_workgraph_schema,
2246 }],
2247 owned_objects: &[
2248 meerkat_sqlite::SchemaObject {
2249 kind: meerkat_sqlite::SchemaObjectKind::Table,
2250 name: "workgraph_items",
2251 },
2252 meerkat_sqlite::SchemaObject {
2253 kind: meerkat_sqlite::SchemaObjectKind::Index,
2254 name: "idx_workgraph_items_realm_namespace_updated",
2255 },
2256 meerkat_sqlite::SchemaObject {
2257 kind: meerkat_sqlite::SchemaObjectKind::Table,
2258 name: "workgraph_attention",
2259 },
2260 meerkat_sqlite::SchemaObject {
2261 kind: meerkat_sqlite::SchemaObjectKind::Index,
2262 name: "idx_workgraph_attention_realm_namespace_updated",
2263 },
2264 meerkat_sqlite::SchemaObject {
2265 kind: meerkat_sqlite::SchemaObjectKind::Index,
2266 name: "idx_workgraph_attention_scope_status",
2267 },
2268 meerkat_sqlite::SchemaObject {
2269 kind: meerkat_sqlite::SchemaObjectKind::Table,
2270 name: "workgraph_edges",
2271 },
2272 meerkat_sqlite::SchemaObject {
2273 kind: meerkat_sqlite::SchemaObjectKind::Table,
2274 name: "workgraph_execution_bindings",
2275 },
2276 meerkat_sqlite::SchemaObject {
2277 kind: meerkat_sqlite::SchemaObjectKind::Index,
2278 name: "idx_workgraph_execution_bindings_item",
2279 },
2280 meerkat_sqlite::SchemaObject {
2281 kind: meerkat_sqlite::SchemaObjectKind::Index,
2282 name: "idx_workgraph_execution_bindings_root",
2283 },
2284 meerkat_sqlite::SchemaObject {
2285 kind: meerkat_sqlite::SchemaObjectKind::Index,
2286 name: "idx_workgraph_execution_bindings_supersedes",
2287 },
2288 meerkat_sqlite::SchemaObject {
2289 kind: meerkat_sqlite::SchemaObjectKind::Index,
2290 name: "idx_workgraph_execution_bindings_target_run",
2291 },
2292 meerkat_sqlite::SchemaObject {
2293 kind: meerkat_sqlite::SchemaObjectKind::Index,
2294 name: "idx_workgraph_execution_bindings_recovery",
2295 },
2296 meerkat_sqlite::SchemaObject {
2297 kind: meerkat_sqlite::SchemaObjectKind::Table,
2298 name: "workgraph_events",
2299 },
2300 meerkat_sqlite::SchemaObject {
2301 kind: meerkat_sqlite::SchemaObjectKind::Index,
2302 name: "idx_workgraph_events_realm_namespace_seq",
2303 },
2304 ],
2305 retired_objects: &[],
2306};
2307
2308#[cfg(not(target_arch = "wasm32"))]
2309fn initialize_current_workgraph_schema(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2310 migration_0001_workgraph_schema(tx)?;
2311 migration_0002_attention_query_columns(tx)?;
2312 migration_0003_execution_bindings(tx)
2313}
2314
2315#[cfg(not(target_arch = "wasm32"))]
2316fn migration_0003_execution_bindings(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2317 tx.execute_batch(
2318 r"
2319 CREATE TABLE IF NOT EXISTS workgraph_execution_bindings (
2320 realm_id TEXT NOT NULL,
2321 namespace TEXT NOT NULL,
2322 binding_id TEXT NOT NULL,
2323 item_id TEXT NOT NULL,
2324 supersedes_binding_id TEXT,
2325 idempotency_key TEXT NOT NULL,
2326 target_run_id TEXT NOT NULL,
2327 revision INTEGER NOT NULL,
2328 recovery_pending INTEGER NOT NULL CHECK (recovery_pending IN (0, 1)),
2329 created_at_utc TEXT NOT NULL,
2330 binding_json TEXT NOT NULL,
2331 PRIMARY KEY (realm_id, namespace, binding_id),
2332 UNIQUE (realm_id, namespace, item_id, idempotency_key),
2333 UNIQUE (realm_id, namespace, target_run_id)
2334 );
2335 CREATE INDEX IF NOT EXISTS idx_workgraph_execution_bindings_item
2336 ON workgraph_execution_bindings
2337 (realm_id, namespace, item_id, created_at_utc, binding_id);
2338 CREATE UNIQUE INDEX IF NOT EXISTS idx_workgraph_execution_bindings_root
2339 ON workgraph_execution_bindings (realm_id, namespace, item_id)
2340 WHERE supersedes_binding_id IS NULL;
2341 CREATE UNIQUE INDEX IF NOT EXISTS idx_workgraph_execution_bindings_supersedes
2342 ON workgraph_execution_bindings (realm_id, namespace, supersedes_binding_id)
2343 WHERE supersedes_binding_id IS NOT NULL;
2344 CREATE UNIQUE INDEX IF NOT EXISTS idx_workgraph_execution_bindings_target_run
2345 ON workgraph_execution_bindings (target_run_id);
2346 CREATE INDEX IF NOT EXISTS idx_workgraph_execution_bindings_recovery
2347 ON workgraph_execution_bindings
2348 (realm_id, recovery_pending, created_at_utc, binding_id);
2349 ",
2350 )
2351}
2352
2353#[cfg(not(target_arch = "wasm32"))]
2354fn migration_0001_workgraph_schema(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2355 tx.execute_batch(
2356 r"
2357 CREATE TABLE IF NOT EXISTS workgraph_items (
2358 realm_id TEXT NOT NULL,
2359 namespace TEXT NOT NULL,
2360 item_id TEXT NOT NULL,
2361 revision INTEGER NOT NULL,
2362 updated_at_utc TEXT NOT NULL,
2363 item_json TEXT NOT NULL,
2364 PRIMARY KEY (realm_id, namespace, item_id)
2365 );
2366 CREATE INDEX IF NOT EXISTS idx_workgraph_items_realm_namespace_updated
2367 ON workgraph_items (realm_id, namespace, updated_at_utc);
2368
2369 CREATE TABLE IF NOT EXISTS workgraph_attention (
2370 realm_id TEXT NOT NULL,
2371 namespace TEXT NOT NULL,
2372 binding_id TEXT NOT NULL,
2373 revision INTEGER NOT NULL,
2374 updated_at_utc TEXT NOT NULL,
2375 attention_json TEXT NOT NULL,
2376 PRIMARY KEY (realm_id, namespace, binding_id)
2377 );
2378 CREATE INDEX IF NOT EXISTS idx_workgraph_attention_realm_namespace_updated
2379 ON workgraph_attention (realm_id, namespace, updated_at_utc);
2380
2381 CREATE TABLE IF NOT EXISTS workgraph_edges (
2382 realm_id TEXT NOT NULL,
2383 namespace TEXT NOT NULL,
2384 edge_kind TEXT NOT NULL,
2385 from_id TEXT NOT NULL,
2386 to_id TEXT NOT NULL,
2387 edge_json TEXT NOT NULL,
2388 PRIMARY KEY (realm_id, namespace, edge_kind, from_id, to_id)
2389 );
2390
2391 CREATE TABLE IF NOT EXISTS workgraph_events (
2392 seq INTEGER PRIMARY KEY AUTOINCREMENT,
2393 realm_id TEXT NOT NULL,
2394 namespace TEXT NOT NULL,
2395 item_id TEXT,
2396 event_kind TEXT NOT NULL,
2397 at_utc TEXT NOT NULL,
2398 event_json TEXT NOT NULL
2399 );
2400 CREATE INDEX IF NOT EXISTS idx_workgraph_events_realm_namespace_seq
2401 ON workgraph_events (realm_id, namespace, seq);
2402 ",
2403 )
2404}
2405
2406#[cfg(not(target_arch = "wasm32"))]
2407fn insert_item_tx(tx: &Transaction<'_>, item: &WorkItem) -> Result<(), WorkGraphError> {
2408 let json = serde_json::to_string(item).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2409 tx.execute(
2410 "INSERT INTO workgraph_items (realm_id, namespace, item_id, revision, updated_at_utc, item_json)
2411 VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
2412 params![
2413 item.realm_id,
2414 item.namespace.as_str(),
2415 item.id.as_str(),
2416 item.revision,
2417 item.updated_at.to_rfc3339(),
2418 json,
2419 ],
2420 )
2421 .map_err(|err| map_sqlite_insert_item_error(err, item))?;
2422 Ok(())
2423}
2424
2425#[cfg(not(target_arch = "wasm32"))]
2426fn update_item_tx(
2427 tx: &Transaction<'_>,
2428 item: &WorkItem,
2429 expected_previous_revision: u64,
2430) -> Result<usize, WorkGraphError> {
2431 let json = serde_json::to_string(item).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2432 tx.execute(
2433 "UPDATE workgraph_items
2434 SET revision = ?4, updated_at_utc = ?5, item_json = ?6
2435 WHERE realm_id = ?1 AND namespace = ?2 AND item_id = ?3 AND revision = ?7",
2436 params![
2437 item.realm_id,
2438 item.namespace.as_str(),
2439 item.id.as_str(),
2440 item.revision,
2441 item.updated_at.to_rfc3339(),
2442 json,
2443 expected_previous_revision,
2444 ],
2445 )
2446 .map_err(|err| WorkGraphError::Store(err.to_string()))
2447}
2448
2449#[cfg(not(target_arch = "wasm32"))]
2450fn upsert_item_tx(tx: &Transaction<'_>, item: &WorkItem) -> Result<(), WorkGraphError> {
2451 let json = serde_json::to_string(item).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2452 tx.execute(
2453 "INSERT INTO workgraph_items
2454 (realm_id, namespace, item_id, revision, updated_at_utc, item_json)
2455 VALUES (?1, ?2, ?3, ?4, ?5, ?6)
2456 ON CONFLICT(realm_id, namespace, item_id) DO UPDATE SET
2457 revision = excluded.revision,
2458 updated_at_utc = excluded.updated_at_utc,
2459 item_json = excluded.item_json",
2460 params![
2461 item.realm_id,
2462 item.namespace.as_str(),
2463 item.id.as_str(),
2464 item.revision,
2465 item.updated_at.to_rfc3339(),
2466 json,
2467 ],
2468 )
2469 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2470 Ok(())
2471}
2472
2473#[cfg(not(target_arch = "wasm32"))]
2474fn map_sqlite_insert_item_error(err: Error, item: &WorkItem) -> WorkGraphError {
2475 if sqlite_constraint_violation(&err) {
2476 return WorkGraphError::Conflict(format!("work item {} already exists", item.id));
2477 }
2478 WorkGraphError::Store(err.to_string())
2479}
2480
2481#[cfg(not(target_arch = "wasm32"))]
2482fn map_sqlite_insert_attention_error(
2483 err: Error,
2484 attention: &WorkAttentionBinding,
2485) -> WorkGraphError {
2486 if sqlite_constraint_violation(&err) {
2487 return WorkGraphError::Conflict(format!(
2488 "work attention binding {} already exists",
2489 attention.binding_id
2490 ));
2491 }
2492 WorkGraphError::Store(err.to_string())
2493}
2494
2495#[cfg(not(target_arch = "wasm32"))]
2496fn sqlite_constraint_violation(err: &Error) -> bool {
2497 matches!(
2498 err,
2499 Error::SqliteFailure(sqlite_error, _)
2500 if sqlite_error.code == ErrorCode::ConstraintViolation
2501 )
2502}
2503
2504#[cfg(not(target_arch = "wasm32"))]
2505fn current_revision_tx(
2506 tx: &Transaction<'_>,
2507 realm_id: &str,
2508 namespace: &WorkNamespace,
2509 id: &WorkItemId,
2510) -> Result<Option<u64>, WorkGraphError> {
2511 tx.query_row(
2512 "SELECT revision FROM workgraph_items WHERE realm_id = ?1 AND namespace = ?2 AND item_id = ?3",
2513 params![realm_id, namespace.as_str(), id.as_str()],
2514 |row| row.get::<_, u64>(0),
2515 )
2516 .optional()
2517 .map_err(|err| WorkGraphError::Store(err.to_string()))
2518}
2519
2520#[cfg(not(target_arch = "wasm32"))]
2521fn insert_attention_tx(
2522 tx: &Transaction<'_>,
2523 attention: &WorkAttentionBinding,
2524) -> Result<(), WorkGraphError> {
2525 let json =
2526 serde_json::to_string(attention).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2527 tx.execute(
2528 "INSERT INTO workgraph_attention
2529 (realm_id, namespace, binding_id, revision, updated_at_utc, attention_json,
2530 status, target_key)
2531 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
2532 params![
2533 attention.work_ref.realm_id,
2534 attention.work_ref.namespace.as_str(),
2535 attention.binding_id.as_str(),
2536 attention.machine_state.revision,
2537 attention.updated_at.to_rfc3339(),
2538 json,
2539 attention.status.status_key(),
2540 attention.target.target_key(),
2541 ],
2542 )
2543 .map_err(|err| map_sqlite_insert_attention_error(err, attention))?;
2544 Ok(())
2545}
2546
2547#[cfg(not(target_arch = "wasm32"))]
2548fn update_attention_tx(
2549 tx: &Transaction<'_>,
2550 attention: &WorkAttentionBinding,
2551 expected_previous_revision: u64,
2552) -> Result<usize, WorkGraphError> {
2553 let json =
2554 serde_json::to_string(attention).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2555 tx.execute(
2556 "UPDATE workgraph_attention
2557 SET revision = ?4, updated_at_utc = ?5, attention_json = ?6,
2558 status = ?8, target_key = ?9
2559 WHERE realm_id = ?1 AND namespace = ?2 AND binding_id = ?3 AND revision = ?7",
2560 params![
2561 attention.work_ref.realm_id,
2562 attention.work_ref.namespace.as_str(),
2563 attention.binding_id.as_str(),
2564 attention.machine_state.revision,
2565 attention.updated_at.to_rfc3339(),
2566 json,
2567 expected_previous_revision,
2568 attention.status.status_key(),
2569 attention.target.target_key(),
2570 ],
2571 )
2572 .map_err(|err| WorkGraphError::Store(err.to_string()))
2573}
2574
2575#[cfg(not(target_arch = "wasm32"))]
2582fn migration_0002_attention_query_columns(tx: &Transaction<'_>) -> Result<(), rusqlite::Error> {
2583 let existing: Vec<String> = tx
2584 .prepare("PRAGMA table_info(workgraph_attention)")?
2585 .query_map([], |row| row.get::<_, String>(1))?
2586 .collect::<Result<_, _>>()?;
2587 if !existing.iter().any(|name| name == "status") {
2588 tx.execute("ALTER TABLE workgraph_attention ADD COLUMN status TEXT", [])?;
2589 }
2590 if !existing.iter().any(|name| name == "target_key") {
2591 tx.execute(
2592 "ALTER TABLE workgraph_attention ADD COLUMN target_key TEXT",
2593 [],
2594 )?;
2595 }
2596 let backfill: Vec<(String, String, String, WorkAttentionBinding)> = {
2597 let mut stmt = tx.prepare(
2598 "SELECT realm_id, namespace, binding_id, attention_json
2599 FROM workgraph_attention
2600 WHERE status IS NULL OR target_key IS NULL",
2601 )?;
2602 let rows = stmt.query_map([], |row| {
2603 Ok((
2604 row.get::<_, String>(0)?,
2605 row.get::<_, String>(1)?,
2606 row.get::<_, String>(2)?,
2607 row_json::<WorkAttentionBinding>(row, 3)?,
2608 ))
2609 })?;
2610 rows.collect::<Result<_, _>>()?
2611 };
2612 for (realm_id, namespace, binding_id, binding) in backfill {
2613 tx.execute(
2614 "UPDATE workgraph_attention
2615 SET status = ?4, target_key = ?5
2616 WHERE realm_id = ?1 AND namespace = ?2 AND binding_id = ?3",
2617 params![
2618 realm_id,
2619 namespace,
2620 binding_id,
2621 binding.status.status_key(),
2622 binding.target.target_key(),
2623 ],
2624 )?;
2625 }
2626 tx.execute(
2627 "CREATE INDEX IF NOT EXISTS idx_workgraph_attention_scope_status
2628 ON workgraph_attention (realm_id, namespace, status, target_key)",
2629 [],
2630 )?;
2631 Ok(())
2632}
2633
2634#[cfg(not(target_arch = "wasm32"))]
2635fn pre_0_8_10_attention_import_error(
2636 binding_id: &str,
2637 detail: impl std::fmt::Display,
2638) -> rusqlite::Error {
2639 rusqlite::Error::ToSqlConversionFailure(Box::new(std::io::Error::new(
2640 std::io::ErrorKind::InvalidData,
2641 format!("pre-v0.8.10 workgraph attention row `{binding_id}`: {detail}"),
2642 )))
2643}
2644
2645#[cfg(not(target_arch = "wasm32"))]
2655pub fn prepare_pre_0_8_10_workgraph_attention(
2656 tx: &Transaction<'_>,
2657) -> Result<meerkat_sqlite::MaintenancePrepareReport, rusqlite::Error> {
2658 let columns = tx
2659 .prepare("PRAGMA table_info(workgraph_attention)")?
2660 .query_map([], |row| row.get::<_, String>(1))?
2661 .collect::<Result<Vec<_>, _>>()?;
2662 let has_status = columns.iter().any(|name| name == "status");
2663 let has_target_key = columns.iter().any(|name| name == "target_key");
2664 match (has_status, has_target_key) {
2665 (false, false) => {
2666 return Ok(meerkat_sqlite::MaintenancePrepareReport::default());
2667 }
2668 (true, true) => {}
2669 _ => {
2670 return Err(pre_0_8_10_attention_import_error(
2671 "<catalog>",
2672 "status and target_key projection columns are not an exact pair",
2673 ));
2674 }
2675 }
2676
2677 struct ProjectionRepair {
2678 realm_id: String,
2679 namespace: String,
2680 binding_id: String,
2681 source_status: Option<String>,
2682 source_target_key: Option<String>,
2683 expected_status: String,
2684 expected_target_key: String,
2685 }
2686
2687 let repairs = {
2688 let mut statement = tx.prepare(
2689 "SELECT realm_id, namespace, binding_id, attention_json, status, target_key
2690 FROM workgraph_attention
2691 ORDER BY realm_id, namespace, binding_id",
2692 )?;
2693 let rows = statement.query_map([], |row| {
2694 Ok((
2695 row.get::<_, String>(0)?,
2696 row.get::<_, String>(1)?,
2697 row.get::<_, String>(2)?,
2698 row.get::<_, String>(3)?,
2699 row.get::<_, Option<String>>(4)?,
2700 row.get::<_, Option<String>>(5)?,
2701 ))
2702 })?;
2703 let mut repairs = Vec::new();
2704 for row in rows {
2705 let (realm_id, namespace, binding_id, attention_json, status, target_key) = row?;
2706 let binding: WorkAttentionBinding = serde_json::from_str(&attention_json)
2707 .map_err(|error| pre_0_8_10_attention_import_error(&binding_id, error))?;
2708 let expected_status = binding.status.status_key().to_string();
2709 let expected_target_key = binding.target.target_key();
2710 if status
2711 .as_deref()
2712 .is_some_and(|value| value != expected_status)
2713 {
2714 return Err(pre_0_8_10_attention_import_error(
2715 &binding_id,
2716 format!(
2717 "status projection `{}` disagrees with typed authority `{expected_status}`",
2718 status.as_deref().unwrap_or_default()
2719 ),
2720 ));
2721 }
2722 if target_key
2723 .as_deref()
2724 .is_some_and(|value| value != expected_target_key)
2725 {
2726 return Err(pre_0_8_10_attention_import_error(
2727 &binding_id,
2728 format!(
2729 "target_key projection `{}` disagrees with typed authority `{expected_target_key}`",
2730 target_key.as_deref().unwrap_or_default()
2731 ),
2732 ));
2733 }
2734 if status.is_none() || target_key.is_none() {
2735 repairs.push(ProjectionRepair {
2736 realm_id,
2737 namespace,
2738 binding_id,
2739 source_status: status,
2740 source_target_key: target_key,
2741 expected_status,
2742 expected_target_key,
2743 });
2744 }
2745 }
2746 repairs
2747 };
2748
2749 let changed = repairs.len();
2750 for repair in repairs {
2751 let updated = tx.execute(
2752 "UPDATE workgraph_attention
2753 SET status = ?4, target_key = ?5
2754 WHERE realm_id = ?1 AND namespace = ?2 AND binding_id = ?3
2755 AND status IS ?6 AND target_key IS ?7",
2756 params![
2757 repair.realm_id,
2758 repair.namespace,
2759 repair.binding_id,
2760 repair.expected_status,
2761 repair.expected_target_key,
2762 repair.source_status,
2763 repair.source_target_key,
2764 ],
2765 )?;
2766 if updated != 1 {
2767 return Err(pre_0_8_10_attention_import_error(
2768 &repair.binding_id,
2769 "source projection changed inside the maintenance transaction",
2770 ));
2771 }
2772 }
2773
2774 Ok(meerkat_sqlite::MaintenancePrepareReport { changed })
2775}
2776
2777#[cfg(not(target_arch = "wasm32"))]
2783fn active_target_occupant_tx(
2784 tx: &Transaction<'_>,
2785 candidate: &WorkAttentionBinding,
2786) -> Result<Option<WorkAttentionBindingId>, WorkGraphError> {
2787 if !matches!(candidate.status, WorkAttentionStatus::Active) {
2788 return Ok(None);
2789 }
2790 let target_key = candidate.target.target_key();
2791 let mut stmt = tx
2792 .prepare(
2793 "SELECT binding_id, attention_json FROM workgraph_attention
2794 WHERE realm_id = ?1 AND namespace = ?2 AND binding_id != ?3
2795 AND (status = 'active' OR status IS NULL)
2796 AND (target_key = ?4 OR target_key IS NULL)",
2797 )
2798 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2799 let rows = stmt
2800 .query_map(
2801 params![
2802 candidate.work_ref.realm_id,
2803 candidate.work_ref.namespace.as_str(),
2804 candidate.binding_id.as_str(),
2805 target_key,
2806 ],
2807 |row| {
2808 Ok((
2809 row.get::<_, String>(0)?,
2810 row_json::<WorkAttentionBinding>(row, 1)?,
2811 ))
2812 },
2813 )
2814 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2815 for row in rows {
2816 let (_, binding) = row.map_err(|err| WorkGraphError::Store(err.to_string()))?;
2817 if matches!(binding.status, WorkAttentionStatus::Active)
2818 && binding.target.target_key() == target_key
2819 {
2820 return Ok(Some(binding.binding_id));
2821 }
2822 }
2823 Ok(None)
2824}
2825
2826fn active_target_conflict(
2829 candidate: &WorkAttentionBinding,
2830 occupant: &WorkAttentionBindingId,
2831) -> WorkGraphError {
2832 WorkGraphError::Conflict(format!(
2833 "active attention binding {occupant} already targets {} in {}/{}",
2834 candidate.target.target_key(),
2835 candidate.work_ref.realm_id,
2836 candidate.work_ref.namespace.as_str(),
2837 ))
2838}
2839
2840fn active_target_occupant_in<'a>(
2843 bindings: impl Iterator<Item = &'a WorkAttentionBinding>,
2844 candidate: &WorkAttentionBinding,
2845) -> Option<WorkAttentionBindingId> {
2846 if !matches!(candidate.status, WorkAttentionStatus::Active) {
2847 return None;
2848 }
2849 let target_key = candidate.target.target_key();
2850 bindings
2851 .filter(|binding| {
2852 binding.binding_id != candidate.binding_id
2853 && binding.work_ref.realm_id == candidate.work_ref.realm_id
2854 && binding.work_ref.namespace == candidate.work_ref.namespace
2855 && matches!(binding.status, WorkAttentionStatus::Active)
2856 && binding.target.target_key() == target_key
2857 })
2858 .map(|binding| binding.binding_id.clone())
2859 .next()
2860}
2861
2862#[cfg(not(target_arch = "wasm32"))]
2863fn upsert_attention_tx(
2864 tx: &Transaction<'_>,
2865 attention: &WorkAttentionBinding,
2866) -> Result<(), WorkGraphError> {
2867 let json =
2868 serde_json::to_string(attention).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2869 tx.execute(
2870 "INSERT INTO workgraph_attention
2871 (realm_id, namespace, binding_id, revision, updated_at_utc, attention_json,
2872 status, target_key)
2873 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
2874 ON CONFLICT(realm_id, namespace, binding_id) DO UPDATE SET
2875 revision = excluded.revision,
2876 updated_at_utc = excluded.updated_at_utc,
2877 attention_json = excluded.attention_json,
2878 status = excluded.status,
2879 target_key = excluded.target_key",
2880 params![
2881 attention.work_ref.realm_id,
2882 attention.work_ref.namespace.as_str(),
2883 attention.binding_id.as_str(),
2884 attention.machine_state.revision,
2885 attention.updated_at.to_rfc3339(),
2886 json,
2887 attention.status.status_key(),
2888 attention.target.target_key(),
2889 ],
2890 )
2891 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2892 Ok(())
2893}
2894
2895#[cfg(not(target_arch = "wasm32"))]
2896fn current_attention_revision_tx(
2897 tx: &Transaction<'_>,
2898 realm_id: &str,
2899 namespace: &WorkNamespace,
2900 binding_id: &WorkAttentionBindingId,
2901) -> Result<Option<u64>, WorkGraphError> {
2902 tx.query_row(
2903 "SELECT revision FROM workgraph_attention
2904 WHERE realm_id = ?1 AND namespace = ?2 AND binding_id = ?3",
2905 params![realm_id, namespace.as_str(), binding_id.as_str()],
2906 |row| row.get::<_, u64>(0),
2907 )
2908 .optional()
2909 .map_err(|err| WorkGraphError::Store(err.to_string()))
2910}
2911
2912#[cfg(not(target_arch = "wasm32"))]
2913fn insert_edge_tx(tx: &Transaction<'_>, edge: &WorkEdge) -> Result<(), WorkGraphError> {
2914 let json = serde_json::to_string(edge).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2915 tx.execute(
2916 "INSERT INTO workgraph_edges
2917 (realm_id, namespace, edge_kind, from_id, to_id, edge_json)
2918 VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
2919 params![
2920 edge.realm_id,
2921 edge.namespace.as_str(),
2922 format!("{:?}", edge.kind),
2923 edge.from_id.as_str(),
2924 edge.to_id.as_str(),
2925 json,
2926 ],
2927 )
2928 .map_err(|err| map_sqlite_insert_edge_error(err, edge))?;
2929 Ok(())
2930}
2931
2932fn duplicate_edge_error(edge: &WorkEdge) -> WorkGraphError {
2933 WorkGraphError::Conflict(format!(
2934 "work edge {:?} {} -> {} already exists",
2935 edge.kind, edge.from_id, edge.to_id
2936 ))
2937}
2938
2939#[cfg(not(target_arch = "wasm32"))]
2940fn map_sqlite_insert_edge_error(err: rusqlite::Error, edge: &WorkEdge) -> WorkGraphError {
2941 match err {
2942 rusqlite::Error::SqliteFailure(failure, _)
2943 if failure.code == ErrorCode::ConstraintViolation =>
2944 {
2945 duplicate_edge_error(edge)
2946 }
2947 err => WorkGraphError::Store(err.to_string()),
2948 }
2949}
2950
2951#[cfg(not(target_arch = "wasm32"))]
2952fn insert_event_tx(tx: &Transaction<'_>, event: &WorkGraphEvent) -> Result<(), WorkGraphError> {
2953 let json =
2954 serde_json::to_string(event).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2955 tx.execute(
2956 "INSERT INTO workgraph_events
2957 (realm_id, namespace, item_id, event_kind, at_utc, event_json)
2958 VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
2959 params![
2960 event.realm_id,
2961 event.namespace.as_str(),
2962 event.item_id.as_ref().map(WorkItemId::as_str),
2963 format!("{:?}", event.kind),
2964 event.at.to_rfc3339(),
2965 json,
2966 ],
2967 )
2968 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
2969 Ok(())
2970}
2971
2972#[cfg(not(target_arch = "wasm32"))]
2973fn select_item(
2974 conn: &Connection,
2975 realm_id: &str,
2976 namespace: &WorkNamespace,
2977 id: &WorkItemId,
2978) -> Result<Option<WorkItem>, WorkGraphError> {
2979 conn.query_row(
2980 "SELECT item_json FROM workgraph_items WHERE realm_id = ?1 AND namespace = ?2 AND item_id = ?3",
2981 params![realm_id, namespace.as_str(), id.as_str()],
2982 |row| row_json(row, 0),
2983 )
2984 .optional()
2985 .map_err(|err| WorkGraphError::Store(err.to_string()))
2986}
2987
2988#[cfg(not(target_arch = "wasm32"))]
2989fn insert_execution_binding_tx(
2990 tx: &Transaction<'_>,
2991 binding: &WorkExecutionBinding,
2992) -> Result<(), WorkGraphError> {
2993 let json =
2994 serde_json::to_string(binding).map_err(|err| WorkGraphError::Store(err.to_string()))?;
2995 let recovery_pending = execution_recovery_pending(binding)?;
2996 tx.execute(
2997 "INSERT INTO workgraph_execution_bindings
2998 (realm_id, namespace, binding_id, item_id, supersedes_binding_id,
2999 idempotency_key, target_run_id, revision, recovery_pending,
3000 created_at_utc, binding_json)
3001 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)",
3002 params![
3003 binding.work_ref.realm_id,
3004 binding.work_ref.namespace.as_str(),
3005 binding.binding_id.as_str(),
3006 binding.work_ref.item_id.as_str(),
3007 binding
3008 .supersedes
3009 .as_ref()
3010 .map(WorkExecutionBindingId::as_str),
3011 binding.idempotency_key,
3012 binding.target.run_id(),
3013 binding.machine_state.revision,
3014 recovery_pending,
3015 binding.created_at.to_rfc3339(),
3016 json,
3017 ],
3018 )
3019 .map_err(|error| {
3020 if sqlite_constraint_violation(&error) {
3021 WorkGraphError::Conflict(format!(
3022 "work execution binding {} conflicts with the existing execution chain",
3023 binding.binding_id
3024 ))
3025 } else {
3026 WorkGraphError::Store(error.to_string())
3027 }
3028 })?;
3029 Ok(())
3030}
3031
3032#[cfg(not(target_arch = "wasm32"))]
3033fn upsert_execution_binding_tx(
3034 tx: &Transaction<'_>,
3035 binding: &WorkExecutionBinding,
3036) -> Result<(), WorkGraphError> {
3037 let json =
3038 serde_json::to_string(binding).map_err(|error| WorkGraphError::Store(error.to_string()))?;
3039 let recovery_pending = execution_recovery_pending(binding)?;
3040 tx.execute(
3041 "INSERT INTO workgraph_execution_bindings
3042 (realm_id, namespace, binding_id, item_id, supersedes_binding_id,
3043 idempotency_key, target_run_id, revision, recovery_pending,
3044 created_at_utc, binding_json)
3045 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11)
3046 ON CONFLICT(realm_id, namespace, binding_id) DO UPDATE SET
3047 revision = excluded.revision,
3048 recovery_pending = excluded.recovery_pending,
3049 binding_json = excluded.binding_json",
3050 params![
3051 binding.work_ref.realm_id,
3052 binding.work_ref.namespace.as_str(),
3053 binding.binding_id.as_str(),
3054 binding.work_ref.item_id.as_str(),
3055 binding
3056 .supersedes
3057 .as_ref()
3058 .map(WorkExecutionBindingId::as_str),
3059 binding.idempotency_key,
3060 binding.target.run_id(),
3061 binding.machine_state.revision,
3062 recovery_pending,
3063 binding.created_at.to_rfc3339(),
3064 json,
3065 ],
3066 )
3067 .map_err(|error| WorkGraphError::Store(error.to_string()))?;
3068 Ok(())
3069}
3070
3071#[cfg(not(target_arch = "wasm32"))]
3072fn execution_recovery_pending(binding: &WorkExecutionBinding) -> Result<i64, WorkGraphError> {
3073 Ok(i64::from(!crate::WorkExecutionMachine::retry_eligible(
3074 binding,
3075 )?))
3076}
3077
3078#[cfg(not(target_arch = "wasm32"))]
3079fn select_execution_binding(
3080 conn: &Connection,
3081 realm_id: &str,
3082 namespace: &WorkNamespace,
3083 binding_id: &WorkExecutionBindingId,
3084) -> Result<Option<WorkExecutionBinding>, WorkGraphError> {
3085 conn.query_row(
3086 "SELECT binding_json FROM workgraph_execution_bindings
3087 WHERE realm_id = ?1 AND namespace = ?2 AND binding_id = ?3",
3088 params![realm_id, namespace.as_str(), binding_id.as_str()],
3089 |row| row_json(row, 0),
3090 )
3091 .optional()
3092 .map_err(|err| WorkGraphError::Store(err.to_string()))
3093}
3094
3095#[cfg(not(target_arch = "wasm32"))]
3096fn list_sqlite_execution_bindings(
3097 conn: &Connection,
3098 filter: &WorkExecutionBindingFilter,
3099) -> Result<Vec<WorkExecutionBinding>, WorkGraphError> {
3100 let mut stmt = conn
3101 .prepare(
3102 "SELECT binding_json FROM workgraph_execution_bindings
3103 ORDER BY created_at_utc ASC, binding_id ASC",
3104 )
3105 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3106 let rows = stmt
3107 .query_map([], |row| row_json::<WorkExecutionBinding>(row, 0))
3108 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3109 let mut all = Vec::new();
3110 for row in rows {
3111 all.push(row.map_err(|err| WorkGraphError::Store(err.to_string()))?);
3112 }
3113 let superseded = all
3114 .iter()
3115 .filter_map(|binding| {
3116 binding.supersedes.clone().map(|supersedes| {
3117 (
3118 binding.work_ref.realm_id.clone(),
3119 binding.work_ref.namespace.clone(),
3120 supersedes,
3121 )
3122 })
3123 })
3124 .collect::<std::collections::BTreeSet<_>>();
3125 let mut bindings = all
3126 .into_iter()
3127 .filter(|binding| execution_binding_matches_filter(binding, filter, &superseded))
3128 .collect::<Vec<_>>();
3129 if let Some(limit) = filter.limit {
3130 bindings.truncate(limit);
3131 }
3132 Ok(bindings)
3133}
3134
3135#[cfg(not(target_arch = "wasm32"))]
3136fn list_sqlite_items(
3137 conn: &Connection,
3138 filter: &WorkItemFilter,
3139) -> Result<Vec<WorkItem>, WorkGraphError> {
3140 let mut stmt = conn
3141 .prepare("SELECT item_json FROM workgraph_items ORDER BY updated_at_utc ASC, item_id ASC")
3142 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3143 let rows = stmt
3144 .query_map([], |row| row_json::<WorkItem>(row, 0))
3145 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3146 let mut items = Vec::new();
3147 for row in rows {
3148 let item = row.map_err(|err| WorkGraphError::Store(err.to_string()))?;
3149 if item_matches_filter(&item, filter) {
3150 items.push(item);
3151 if filter.limit.is_some_and(|limit| items.len() >= limit) {
3152 break;
3153 }
3154 }
3155 }
3156 Ok(items)
3157}
3158
3159#[cfg(not(target_arch = "wasm32"))]
3160fn select_attention(
3161 conn: &Connection,
3162 realm_id: &str,
3163 namespace: &WorkNamespace,
3164 binding_id: &WorkAttentionBindingId,
3165) -> Result<Option<WorkAttentionBinding>, WorkGraphError> {
3166 conn.query_row(
3167 "SELECT attention_json FROM workgraph_attention
3168 WHERE realm_id = ?1 AND namespace = ?2 AND binding_id = ?3",
3169 params![realm_id, namespace.as_str(), binding_id.as_str()],
3170 |row| row_json(row, 0),
3171 )
3172 .optional()
3173 .map_err(|err| WorkGraphError::Store(err.to_string()))
3174}
3175
3176#[cfg(not(target_arch = "wasm32"))]
3177fn list_sqlite_attention(
3178 conn: &Connection,
3179 filter: &AttentionListRequest,
3180 limit: Option<usize>,
3181) -> Result<Vec<WorkAttentionBinding>, WorkGraphError> {
3182 if limit == Some(0) {
3183 return Ok(Vec::new());
3184 }
3185 let mut clauses: Vec<String> = Vec::new();
3190 let mut params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
3191 if let Some(realm_id) = &filter.realm_id {
3192 params.push(Box::new(realm_id.clone()));
3193 clauses.push(format!("realm_id = ?{}", params.len()));
3194 }
3195 if let Some(namespace) = &filter.namespace {
3196 params.push(Box::new(namespace.as_str().to_string()));
3197 clauses.push(format!("namespace = ?{}", params.len()));
3198 }
3199 if let Some(status) = &filter.status {
3200 params.push(Box::new(status.status_key().to_string()));
3201 clauses.push(format!("(status = ?{} OR status IS NULL)", params.len()));
3202 }
3203 if let Some(target) = &filter.target {
3204 params.push(Box::new(target.target_key()));
3205 clauses.push(format!(
3206 "(target_key = ?{} OR target_key IS NULL)",
3207 params.len()
3208 ));
3209 }
3210 let where_clause = if clauses.is_empty() {
3211 String::new()
3212 } else {
3213 format!(" WHERE {}", clauses.join(" AND "))
3214 };
3215 let sql = format!(
3216 "SELECT attention_json FROM workgraph_attention{where_clause}
3217 ORDER BY updated_at_utc ASC, binding_id ASC"
3218 );
3219 let mut stmt = conn
3220 .prepare(&sql)
3221 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3222 let rows = stmt
3223 .query_map(rusqlite::params_from_iter(params.iter()), |row| {
3224 row_json::<WorkAttentionBinding>(row, 0)
3225 })
3226 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3227 let mut bindings = Vec::new();
3228 for row in rows {
3229 let binding = row.map_err(|err| WorkGraphError::Store(err.to_string()))?;
3230 if attention_matches_filter(&binding, filter) {
3231 bindings.push(binding);
3232 if limit.is_some_and(|limit| bindings.len() >= limit) {
3233 break;
3234 }
3235 }
3236 }
3237 Ok(bindings)
3238}
3239
3240#[cfg(not(target_arch = "wasm32"))]
3241fn list_sqlite_edges(
3242 conn: &Connection,
3243 realm_id: &str,
3244 namespace: &WorkNamespace,
3245 limit: Option<usize>,
3246) -> Result<Vec<WorkEdge>, WorkGraphError> {
3247 if limit == Some(0) {
3248 return Ok(Vec::new());
3249 }
3250 let mut stmt = conn
3251 .prepare(
3252 "SELECT edge_json FROM workgraph_edges
3253 WHERE realm_id = ?1 AND namespace = ?2
3254 ORDER BY edge_kind ASC, from_id ASC, to_id ASC",
3255 )
3256 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3257 let rows = stmt
3258 .query_map(params![realm_id, namespace.as_str()], |row| {
3259 row_json::<WorkEdge>(row, 0)
3260 })
3261 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3262 let mut edges = Vec::new();
3263 for row in rows {
3264 edges.push(row.map_err(|err| WorkGraphError::Store(err.to_string()))?);
3265 if limit.is_some_and(|limit| edges.len() >= limit) {
3266 break;
3267 }
3268 }
3269 Ok(edges)
3270}
3271
3272#[cfg(not(target_arch = "wasm32"))]
3273fn list_sqlite_events(
3274 conn: &Connection,
3275 filter: &WorkGraphEventFilter,
3276) -> Result<Vec<WorkGraphEvent>, WorkGraphError> {
3277 let mut stmt = conn
3278 .prepare("SELECT seq, event_json FROM workgraph_events ORDER BY seq ASC")
3279 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3280 let rows = stmt
3281 .query_map([], |row| {
3282 let seq = row.get::<_, i64>(0)?;
3283 let mut event = row_json::<WorkGraphEvent>(row, 1)?;
3284 event.seq = Some(seq);
3285 Ok(event)
3286 })
3287 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3288 let mut events = Vec::new();
3289 for row in rows {
3290 let event = row.map_err(|err| WorkGraphError::Store(err.to_string()))?;
3291 if event_matches_filter(&event, filter) {
3292 events.push(event);
3293 if filter.limit.is_some_and(|limit| events.len() >= limit) {
3294 break;
3295 }
3296 }
3297 }
3298 Ok(events)
3299}
3300
3301#[cfg(not(target_arch = "wasm32"))]
3302fn list_sqlite_public_events(
3303 conn: &Connection,
3304 filter: &WorkGraphEventFilter,
3305) -> Result<Vec<WorkGraphEvent>, WorkGraphError> {
3306 let limit = filter.limit.unwrap_or(usize::MAX);
3307 if limit == 0 {
3308 return Ok(Vec::new());
3309 }
3310
3311 let mut clauses = vec![
3312 "event_kind != 'ExecutionBound'".to_string(),
3313 "event_kind != 'ExecutionTransitioned'".to_string(),
3314 ];
3315 let mut params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
3316 if let Some(realm_id) = &filter.realm_id {
3317 params.push(Box::new(realm_id.clone()));
3318 clauses.push(format!("realm_id = ?{}", params.len()));
3319 }
3320 if !filter.all_namespaces
3321 && let Some(namespace) = &filter.namespace
3322 {
3323 params.push(Box::new(namespace.as_str().to_string()));
3324 clauses.push(format!("namespace = ?{}", params.len()));
3325 }
3326 if let Some(after_seq) = filter.after_seq {
3327 params.push(Box::new(after_seq));
3328 clauses.push(format!("seq > ?{}", params.len()));
3329 }
3330 params.push(Box::new(i64::try_from(limit).unwrap_or(i64::MAX)));
3331 let sql = format!(
3332 "SELECT seq, event_json FROM workgraph_events WHERE {} ORDER BY seq ASC LIMIT ?{}",
3333 clauses.join(" AND "),
3334 params.len()
3335 );
3336 let mut stmt = conn
3337 .prepare(&sql)
3338 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3339 let rows = stmt
3340 .query_map(rusqlite::params_from_iter(params.iter()), |row| {
3341 let seq = row.get::<_, i64>(0)?;
3342 let mut event = row_json::<WorkGraphEvent>(row, 1)?;
3343 event.seq = Some(seq);
3344 Ok(event)
3345 })
3346 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3347 rows.collect::<Result<Vec<_>, _>>()
3348 .map_err(|err| WorkGraphError::Store(err.to_string()))
3349}
3350
3351#[cfg(not(target_arch = "wasm32"))]
3352fn latest_sqlite_event_seq(
3353 conn: &Connection,
3354 filter: &WorkGraphEventFilter,
3355) -> Result<Option<i64>, WorkGraphError> {
3356 let mut clauses: Vec<String> = Vec::new();
3357 let mut params: Vec<Box<dyn rusqlite::ToSql>> = Vec::new();
3358 if let Some(realm_id) = &filter.realm_id {
3359 params.push(Box::new(realm_id.clone()));
3360 clauses.push(format!("realm_id = ?{}", params.len()));
3361 }
3362 if !filter.all_namespaces
3363 && let Some(namespace) = &filter.namespace
3364 {
3365 params.push(Box::new(namespace.as_str().to_string()));
3366 clauses.push(format!("namespace = ?{}", params.len()));
3367 }
3368 if let Some(after_seq) = filter.after_seq {
3369 params.push(Box::new(after_seq));
3370 clauses.push(format!("seq > ?{}", params.len()));
3371 }
3372 let where_clause = if clauses.is_empty() {
3373 String::new()
3374 } else {
3375 format!(" WHERE {}", clauses.join(" AND "))
3376 };
3377 conn.query_row(
3378 &format!("SELECT MAX(seq) FROM workgraph_events{where_clause}"),
3379 rusqlite::params_from_iter(params.iter()),
3380 |row| row.get::<_, Option<i64>>(0),
3381 )
3382 .map_err(|error| WorkGraphError::Store(error.to_string()))
3383}
3384
3385#[cfg(not(target_arch = "wasm32"))]
3386fn replay_event_tx(tx: &Transaction<'_>, event: &WorkGraphEvent) -> Result<(), WorkGraphError> {
3387 match event.kind {
3388 WorkGraphEventKind::Linked => {
3389 let edge = payload_field::<WorkEdge>(event, "edge")?;
3390 insert_edge_tx(tx, &edge)
3391 }
3392 WorkGraphEventKind::AttentionCreated | WorkGraphEventKind::AttentionUpdated => {
3393 let attention = payload_field::<WorkAttentionBinding>(event, "attention")?;
3394 upsert_attention_tx(tx, &attention)
3395 }
3396 WorkGraphEventKind::ExecutionBound => {
3397 let binding = payload_field::<WorkExecutionBinding>(event, "execution_binding")?;
3398 let commit = crate::WorkExecutionMachine::prepare_bind(binding.clone())?;
3399 if commit.binding() != &binding {
3400 return Err(WorkGraphError::Store(format!(
3401 "execution bind event for {} changed during authority validation",
3402 binding.binding_id
3403 )));
3404 }
3405 validate_execution_event_scope(event, &binding)?;
3406 let item = select_item(
3407 tx,
3408 &binding.work_ref.realm_id,
3409 &binding.work_ref.namespace,
3410 &binding.work_ref.item_id,
3411 )?
3412 .ok_or_else(|| {
3413 WorkGraphError::Store(format!(
3414 "execution bind for {} references a missing work item",
3415 binding.binding_id
3416 ))
3417 })?;
3418 let bindings = list_sqlite_execution_bindings(
3419 tx,
3420 &WorkExecutionBindingFilter {
3421 realm_id: Some(binding.work_ref.realm_id.clone()),
3422 namespace: Some(binding.work_ref.namespace.clone()),
3423 item_id: Some(binding.work_ref.item_id.clone()),
3424 current_only: false,
3425 limit: None,
3426 },
3427 )?;
3428 validate_execution_binding_insert(
3429 &binding,
3430 item.revision,
3431 std::iter::once(&item),
3432 bindings.iter(),
3433 )?;
3434 insert_execution_binding_tx(tx, &binding)
3435 }
3436 WorkGraphEventKind::ExecutionTransitioned => {
3437 let binding = payload_field::<WorkExecutionBinding>(event, "execution_binding")?;
3438 let observation =
3439 payload_field::<crate::WorkExecutionObservation>(event, "observation")?;
3440 validate_execution_event_scope(event, &binding)?;
3441 let current = select_execution_binding(
3442 tx,
3443 &binding.work_ref.realm_id,
3444 &binding.work_ref.namespace,
3445 &binding.binding_id,
3446 )?
3447 .ok_or_else(|| {
3448 WorkGraphError::Store(format!(
3449 "execution transition for {} precedes its bind event",
3450 binding.binding_id
3451 ))
3452 })?;
3453 let commit = crate::WorkExecutionMachine::prepare_observation(
3454 current.clone(),
3455 current.machine_state.revision,
3456 observation,
3457 )?;
3458 if commit.binding() != &binding {
3459 return Err(WorkGraphError::Store(format!(
3460 "execution transition event for {} is not the exact generated machine result",
3461 binding.binding_id
3462 )));
3463 }
3464 upsert_execution_binding_tx(tx, &binding)
3465 }
3466 WorkGraphEventKind::Created
3467 | WorkGraphEventKind::Updated
3468 | WorkGraphEventKind::Claimed
3469 | WorkGraphEventKind::Released
3470 | WorkGraphEventKind::Blocked
3471 | WorkGraphEventKind::Closed
3472 | WorkGraphEventKind::EvidenceAdded => {
3473 let item = payload_field::<WorkItem>(event, "item")?;
3474 upsert_item_tx(tx, &item)
3475 }
3476 }
3477}
3478
3479#[cfg(not(target_arch = "wasm32"))]
3480fn validate_execution_event_scope(
3481 event: &WorkGraphEvent,
3482 binding: &WorkExecutionBinding,
3483) -> Result<(), WorkGraphError> {
3484 if event.realm_id != binding.work_ref.realm_id
3485 || event.namespace != binding.work_ref.namespace
3486 || event.item_id.as_ref() != Some(&binding.work_ref.item_id)
3487 {
3488 return Err(WorkGraphError::Store(format!(
3489 "execution event scope does not match binding {}",
3490 binding.binding_id
3491 )));
3492 }
3493 Ok(())
3494}
3495
3496#[cfg(not(target_arch = "wasm32"))]
3497fn normalize_attention_for_terminal_items_tx(tx: &Transaction<'_>) -> Result<(), WorkGraphError> {
3498 let bindings = {
3499 let mut stmt = tx
3500 .prepare("SELECT attention_json FROM workgraph_attention")
3501 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3502 let rows = stmt
3503 .query_map([], |row| row_json::<WorkAttentionBinding>(row, 0))
3504 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3505 let mut bindings = Vec::new();
3506 for row in rows {
3507 bindings.push(row.map_err(|err| WorkGraphError::Store(err.to_string()))?);
3508 }
3509 bindings
3510 };
3511
3512 for binding in bindings {
3513 if matches!(
3514 binding.status,
3515 WorkAttentionStatus::Stopped | WorkAttentionStatus::Superseded
3516 ) {
3517 continue;
3518 }
3519 let item = tx
3520 .query_row(
3521 "SELECT item_json FROM workgraph_items
3522 WHERE realm_id = ?1 AND namespace = ?2 AND item_id = ?3",
3523 params![
3524 binding.work_ref.realm_id,
3525 binding.work_ref.namespace.as_str(),
3526 binding.work_ref.item_id.as_str(),
3527 ],
3528 |row| row_json::<WorkItem>(row, 0),
3529 )
3530 .optional()
3531 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
3532 let Some(item) = item else {
3533 continue;
3534 };
3535 if WorkGraphMachine::classify_terminality(&item)? {
3538 let expected_revision = binding.machine_state.revision;
3539 let stopped = WorkAttentionMachine::stop(binding, expected_revision, item.updated_at)?;
3540 upsert_attention_tx(tx, &stopped)?;
3541 }
3542 }
3543 Ok(())
3544}
3545
3546#[cfg(not(target_arch = "wasm32"))]
3547fn payload_field<T: serde::de::DeserializeOwned>(
3548 event: &WorkGraphEvent,
3549 field: &str,
3550) -> Result<T, WorkGraphError> {
3551 let value = event.payload.get(field).ok_or_else(|| {
3552 WorkGraphError::Store(format!(
3553 "workgraph event {:?} missing payload field `{field}`",
3554 event.kind
3555 ))
3556 })?;
3557 serde_json::from_value(value.clone()).map_err(|err| WorkGraphError::Store(err.to_string()))
3558}
3559
3560#[cfg(not(target_arch = "wasm32"))]
3561fn row_json<T: serde::de::DeserializeOwned>(
3562 row: &rusqlite::Row<'_>,
3563 index: usize,
3564) -> rusqlite::Result<T> {
3565 let json = row.get::<_, String>(index)?;
3566 serde_json::from_str(&json).map_err(|err| {
3567 rusqlite::Error::FromSqlConversionFailure(index, rusqlite::types::Type::Text, Box::new(err))
3568 })
3569}
3570
3571#[cfg(test)]
3572#[allow(clippy::expect_used, clippy::unwrap_used)]
3573mod tests {
3574 use std::collections::BTreeSet;
3575
3576 use chrono::Utc;
3577 use serde_json::json;
3578
3579 use crate::types::WorkEdge;
3580 use crate::{
3581 AttentionDelegatedAuthority, AttentionProjectionPolicy, CreateWorkItemRequest,
3582 GoalAttentionTarget, GoalCreateRequest, GoalRequestCloseRequest, GoalTerminalStatus,
3583 LinkWorkItemsRequest, MemoryWorkGraphStore, WorkAttentionMode, WorkAttentionStatus,
3584 WorkCompletionPolicy, WorkEdgeKind, WorkExecutionBinding, WorkExecutionBindingId,
3585 WorkExecutionMachine, WorkExecutionObservation, WorkExecutionTarget, WorkGraphError,
3586 WorkGraphEvent, WorkGraphEventFilter, WorkGraphEventKind, WorkGraphService, WorkGraphStore,
3587 WorkItemFilter, WorkItemId, WorkItemRef, WorkNamespace,
3588 };
3589
3590 fn test_edge() -> WorkEdge {
3591 WorkEdge {
3592 realm_id: "realm".to_string(),
3593 namespace: WorkNamespace::default(),
3594 kind: WorkEdgeKind::Blocks,
3595 from_id: WorkItemId::generated(),
3596 to_id: WorkItemId::generated(),
3597 created_at: Utc::now(),
3598 }
3599 }
3600
3601 fn link_event(edge: &WorkEdge) -> WorkGraphEvent {
3602 WorkGraphEvent::graph(
3603 edge.realm_id.clone(),
3604 edge.namespace.clone(),
3605 WorkGraphEventKind::Linked,
3606 edge.created_at,
3607 json!({ "edge": edge }),
3608 )
3609 }
3610
3611 async fn stale_execution_commit_is_refused(store: std::sync::Arc<dyn WorkGraphStore>) {
3612 let service =
3613 WorkGraphService::with_scope(store.clone(), "realm", WorkNamespace::default());
3614 let item = service
3615 .create(CreateWorkItemRequest {
3616 title: "immutable execution".to_string(),
3617 ..Default::default()
3618 })
3619 .await
3620 .expect("item");
3621 let binding_id = WorkExecutionBindingId::new("execution-immutable").expect("binding id");
3622 let target = WorkExecutionTarget::mob_flow(
3623 "mob",
3624 "flow",
3625 format!("sha256:{}", "c".repeat(64)),
3626 "46371bce-c308-58a4-bf0b-0a262de45c12",
3627 crate::WorkExecutionAuthority::TargetOwner,
3628 json!({}),
3629 )
3630 .expect("target");
3631 let (machine_state, _) =
3632 WorkExecutionMachine::bind(&binding_id, target.run_id()).expect("bind machine");
3633 let bound = service
3634 .bind_execution(
3635 WorkExecutionBinding {
3636 binding_id,
3637 work_ref: WorkItemRef {
3638 realm_id: item.realm_id.clone(),
3639 namespace: item.namespace.clone(),
3640 item_id: item.id.clone(),
3641 },
3642 target,
3643 idempotency_key: "original-key".to_string(),
3644 correlation_id: "f9ae62da-662f-5c50-940e-442c529d8e1d".to_string(),
3645 supersedes: None,
3646 machine_state,
3647 created_at: Utc::now(),
3648 },
3649 item.revision,
3650 )
3651 .await
3652 .expect("bind execution")
3653 .binding;
3654 let commit = WorkExecutionMachine::prepare_observation(
3655 bound.clone(),
3656 bound.machine_state.revision,
3657 WorkExecutionObservation::FlowRunning,
3658 )
3659 .expect("machine-minted next-state authority");
3660 service
3661 .observe_execution(
3662 Some(bound.work_ref.realm_id.clone()),
3663 Some(bound.work_ref.namespace.clone()),
3664 bound.binding_id.clone(),
3665 bound.machine_state.revision,
3666 WorkExecutionObservation::FlowRunning,
3667 )
3668 .await
3669 .expect("commit competing transition");
3670 let event = WorkGraphEvent::item(
3671 item.realm_id,
3672 item.namespace,
3673 item.id,
3674 WorkGraphEventKind::ExecutionTransitioned,
3675 Utc::now(),
3676 json!({
3677 "execution_binding": commit.binding(),
3678 "observation": WorkExecutionObservation::FlowRunning,
3679 }),
3680 );
3681 let error = store
3682 .update_execution_binding_cas(commit, bound.machine_state.revision, event)
3683 .await
3684 .expect_err("store must reject a commit minted from a stale predecessor");
3685 assert!(matches!(error, WorkGraphError::Conflict(_)));
3686 }
3687
3688 async fn duplicate_execution_run_is_refused(store: std::sync::Arc<dyn WorkGraphStore>) {
3689 let service = WorkGraphService::with_scope(store, "realm", WorkNamespace::default());
3690 let first_item = service
3691 .create(CreateWorkItemRequest {
3692 title: "first execution".to_string(),
3693 ..Default::default()
3694 })
3695 .await
3696 .expect("first item");
3697 let second_item = service
3698 .create(CreateWorkItemRequest {
3699 title: "second execution".to_string(),
3700 ..Default::default()
3701 })
3702 .await
3703 .expect("second item");
3704 let run_id = "24d61f25-09db-5327-99e7-63d7390a1e95";
3705
3706 for (index, item) in [first_item, second_item].into_iter().enumerate() {
3707 let binding_id =
3708 WorkExecutionBindingId::new(format!("execution-run-{index}")).expect("binding id");
3709 let target = WorkExecutionTarget::mob_flow(
3710 "mob",
3711 "flow",
3712 format!("sha256:{}", "d".repeat(64)),
3713 run_id,
3714 crate::WorkExecutionAuthority::TargetOwner,
3715 json!({}),
3716 )
3717 .expect("target");
3718 let (machine_state, _) =
3719 WorkExecutionMachine::bind(&binding_id, target.run_id()).expect("bind machine");
3720 let result = service
3721 .bind_execution(
3722 WorkExecutionBinding {
3723 binding_id,
3724 work_ref: WorkItemRef {
3725 realm_id: item.realm_id,
3726 namespace: item.namespace,
3727 item_id: item.id,
3728 },
3729 target,
3730 idempotency_key: format!("run-key-{index}"),
3731 correlation_id: if index == 0 {
3732 "e25abdd9-29cf-56e3-9402-e86c78feec27".to_string()
3733 } else {
3734 "e8c85639-aa77-5d9b-ad77-b13e29675a21".to_string()
3735 },
3736 supersedes: None,
3737 machine_state,
3738 created_at: Utc::now(),
3739 },
3740 item.revision,
3741 )
3742 .await;
3743 if index == 0 {
3744 result.expect("first run binding");
3745 } else {
3746 assert!(matches!(result, Err(WorkGraphError::Conflict(_))));
3747 }
3748 }
3749 }
3750
3751 #[tokio::test]
3752 async fn memory_store_rejects_stale_execution_commit() {
3753 stale_execution_commit_is_refused(std::sync::Arc::new(MemoryWorkGraphStore::new())).await;
3754 }
3755
3756 #[cfg(not(target_arch = "wasm32"))]
3757 #[tokio::test]
3758 async fn sqlite_store_rejects_stale_execution_commit() {
3759 let dir = tempfile::tempdir().expect("tempdir");
3760 stale_execution_commit_is_refused(std::sync::Arc::new(
3761 crate::SqliteWorkGraphStore::open(dir.path().join("workgraph.sqlite3"))
3762 .expect("sqlite store"),
3763 ))
3764 .await;
3765 }
3766
3767 #[cfg(not(target_arch = "wasm32"))]
3768 #[tokio::test]
3769 async fn sqlite_public_event_limit_is_applied_after_internal_visibility_filter() {
3770 let dir = tempfile::tempdir().expect("tempdir");
3771 let store = crate::SqliteWorkGraphStore::open(dir.path().join("workgraph.sqlite3"))
3772 .expect("sqlite store");
3773 let namespace = WorkNamespace::default();
3774 let event = |kind| {
3775 WorkGraphEvent::graph(
3776 "realm".to_string(),
3777 namespace.clone(),
3778 kind,
3779 Utc::now(),
3780 json!({}),
3781 )
3782 };
3783 store
3784 .with_connection(|conn| {
3785 let tx = conn
3786 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)
3787 .map_err(|error| WorkGraphError::Store(error.to_string()))?;
3788 super::insert_event_tx(&tx, &event(WorkGraphEventKind::Created))?;
3789 for _ in 0..300 {
3790 super::insert_event_tx(&tx, &event(WorkGraphEventKind::ExecutionTransitioned))?;
3791 }
3792 super::insert_event_tx(&tx, &event(WorkGraphEventKind::EvidenceAdded))?;
3793 tx.commit()
3794 .map_err(|error| WorkGraphError::Store(error.to_string()))
3795 })
3796 .expect("insert event history");
3797
3798 let public = store
3799 .list_public_events(WorkGraphEventFilter {
3800 realm_id: Some("realm".to_string()),
3801 namespace: Some(namespace),
3802 after_seq: Some(1),
3803 limit: Some(1),
3804 ..WorkGraphEventFilter::default()
3805 })
3806 .await
3807 .expect("public event page");
3808 assert_eq!(public.len(), 1);
3809 assert_eq!(public[0].kind, WorkGraphEventKind::EvidenceAdded);
3810 assert_eq!(public[0].seq, Some(302));
3811 }
3812
3813 #[tokio::test]
3814 async fn memory_store_rejects_cross_item_run_reuse() {
3815 duplicate_execution_run_is_refused(std::sync::Arc::new(MemoryWorkGraphStore::new())).await;
3816 }
3817
3818 #[cfg(not(target_arch = "wasm32"))]
3819 #[tokio::test]
3820 async fn sqlite_store_rejects_cross_item_run_reuse() {
3821 let dir = tempfile::tempdir().expect("tempdir");
3822 duplicate_execution_run_is_refused(std::sync::Arc::new(
3823 crate::SqliteWorkGraphStore::open(dir.path().join("workgraph.sqlite3"))
3824 .expect("sqlite store"),
3825 ))
3826 .await;
3827 }
3828
3829 #[tokio::test]
3830 async fn memory_store_namespace_filters_do_not_leak() {
3831 let store = std::sync::Arc::new(MemoryWorkGraphStore::new());
3832 let default_service =
3833 WorkGraphService::with_scope(store.clone(), "realm", WorkNamespace::default());
3834 let other_service = WorkGraphService::with_scope(
3835 store.clone(),
3836 "realm",
3837 WorkNamespace::new("other").expect("namespace"),
3838 );
3839 default_service
3840 .create(CreateWorkItemRequest {
3841 realm_id: None,
3842 namespace: None,
3843 title: "default".to_string(),
3844 description: None,
3845 priority: Default::default(),
3846 completion_policy: Default::default(),
3847 labels: BTreeSet::new(),
3848 due_at: None,
3849 not_before: None,
3850 snoozed_until: None,
3851 external_refs: Vec::new(),
3852 evidence_refs: Vec::new(),
3853 status: None,
3854 })
3855 .await
3856 .expect("create default");
3857 other_service
3858 .create(CreateWorkItemRequest {
3859 realm_id: None,
3860 namespace: None,
3861 title: "other".to_string(),
3862 description: None,
3863 priority: Default::default(),
3864 completion_policy: Default::default(),
3865 labels: BTreeSet::new(),
3866 due_at: None,
3867 not_before: None,
3868 snoozed_until: None,
3869 external_refs: Vec::new(),
3870 evidence_refs: Vec::new(),
3871 status: None,
3872 })
3873 .await
3874 .expect("create other");
3875
3876 let items = store
3877 .list_items(WorkItemFilter {
3878 realm_id: Some("realm".to_string()),
3879 namespace: Some(WorkNamespace::default()),
3880 ..WorkItemFilter::default()
3881 })
3882 .await
3883 .expect("list");
3884 assert_eq!(items.len(), 1);
3885 assert_eq!(items[0].title, "default");
3886 }
3887
3888 #[tokio::test]
3889 async fn sqlite_rebuild_restores_execution_machine_state_from_events() {
3890 let temp = tempfile::tempdir().expect("tempdir");
3891 let store = std::sync::Arc::new(
3892 crate::SqliteWorkGraphStore::open(temp.path().join("workgraph.db"))
3893 .expect("sqlite store"),
3894 );
3895 let service =
3896 WorkGraphService::with_scope(store.clone(), "realm", WorkNamespace::default());
3897 let item = service
3898 .create(CreateWorkItemRequest {
3899 realm_id: None,
3900 namespace: None,
3901 title: "durable execution".to_string(),
3902 description: None,
3903 priority: Default::default(),
3904 completion_policy: Default::default(),
3905 labels: BTreeSet::new(),
3906 due_at: None,
3907 not_before: None,
3908 snoozed_until: None,
3909 external_refs: Vec::new(),
3910 evidence_refs: Vec::new(),
3911 status: None,
3912 })
3913 .await
3914 .expect("item");
3915 let binding_id = WorkExecutionBindingId::new("execution-sqlite").expect("binding id");
3916 let target = WorkExecutionTarget::mob_flow(
3917 "mob",
3918 "flow",
3919 format!("sha256:{}", "b".repeat(64)),
3920 "d8bb76bb-40e8-54f7-b859-d02827f7d296",
3921 crate::WorkExecutionAuthority::TargetOwner,
3922 json!({}),
3923 )
3924 .expect("target");
3925 let (machine_state, _) =
3926 WorkExecutionMachine::bind(&binding_id, target.run_id()).expect("machine bind");
3927 let bound = service
3928 .bind_execution(
3929 WorkExecutionBinding {
3930 binding_id,
3931 work_ref: WorkItemRef {
3932 realm_id: item.realm_id.clone(),
3933 namespace: item.namespace.clone(),
3934 item_id: item.id.clone(),
3935 },
3936 target,
3937 idempotency_key: "sqlite-key".to_string(),
3938 correlation_id: "6084cb0d-f5df-5814-aad9-c8c6c763ef54".to_string(),
3939 supersedes: None,
3940 machine_state,
3941 created_at: Utc::now(),
3942 },
3943 item.revision,
3944 )
3945 .await
3946 .expect("bind");
3947 let running = service
3948 .observe_execution(
3949 Some(item.realm_id.clone()),
3950 Some(item.namespace.clone()),
3951 bound.binding.binding_id,
3952 1,
3953 WorkExecutionObservation::FlowRunning,
3954 )
3955 .await
3956 .expect("running");
3957 assert_eq!(running.binding.machine_state.revision, 2);
3958
3959 store
3960 .rebuild_projection_from_events()
3961 .expect("rebuild projections");
3962 let restored = service
3963 .execution_binding(
3964 Some(item.realm_id),
3965 Some(item.namespace),
3966 running.binding.binding_id,
3967 )
3968 .await
3969 .expect("restored binding");
3970 assert_eq!(restored.machine_state.revision, 2);
3971 assert_eq!(
3972 service
3973 .execution_bindings_for_recovery(Some("realm".to_string()))
3974 .await
3975 .expect("active recovery queue")
3976 .len(),
3977 1
3978 );
3979 let failed = service
3980 .observe_execution(
3981 Some(restored.work_ref.realm_id.clone()),
3982 Some(restored.work_ref.namespace.clone()),
3983 restored.binding_id.clone(),
3984 restored.machine_state.revision,
3985 WorkExecutionObservation::FlowFailed {
3986 detail: Some("test failure".to_string()),
3987 },
3988 )
3989 .await
3990 .expect("observe failure");
3991 service
3992 .observe_execution(
3993 Some(failed.binding.work_ref.realm_id.clone()),
3994 Some(failed.binding.work_ref.namespace.clone()),
3995 failed.binding.binding_id,
3996 failed.binding.machine_state.revision,
3997 WorkExecutionObservation::FlowFailureEvidenceProjected,
3998 )
3999 .await
4000 .expect("terminal failure");
4001 assert!(
4002 service
4003 .execution_bindings_for_recovery(Some("realm".to_string()))
4004 .await
4005 .expect("terminal recovery queue")
4006 .is_empty()
4007 );
4008 }
4009
4010 #[tokio::test]
4011 async fn memory_store_duplicate_edge_does_not_append_event() {
4012 let store = MemoryWorkGraphStore::new();
4013 let edge = test_edge();
4014 store
4015 .insert_edge(edge.clone(), link_event(&edge))
4016 .await
4017 .expect("insert edge");
4018
4019 let error = store
4020 .insert_edge(edge.clone(), link_event(&edge))
4021 .await
4022 .expect_err("duplicate edge should fail");
4023 assert!(matches!(error, WorkGraphError::Conflict(_)));
4024
4025 let events = store
4026 .list_events(WorkGraphEventFilter {
4027 realm_id: Some(edge.realm_id),
4028 namespace: Some(edge.namespace),
4029 all_namespaces: false,
4030 after_seq: None,
4031 limit: None,
4032 })
4033 .await
4034 .expect("events");
4035 assert_eq!(events.len(), 1);
4036 }
4037
4038 #[cfg(not(target_arch = "wasm32"))]
4042 #[tokio::test]
4043 async fn sqlite_store_duplicate_item_insert_maps_to_conflict() {
4044 let dir = tempfile::tempdir().expect("tempdir");
4045 let path = dir.path().join("workgraph.sqlite3");
4046 let store = std::sync::Arc::new(crate::SqliteWorkGraphStore::open(&path).expect("open"));
4047 let service =
4048 WorkGraphService::with_scope(store.clone(), "realm", WorkNamespace::default());
4049 let item = service
4050 .create(CreateWorkItemRequest {
4051 realm_id: None,
4052 namespace: None,
4053 title: "unique item".to_string(),
4054 description: None,
4055 priority: Default::default(),
4056 completion_policy: Default::default(),
4057 labels: BTreeSet::new(),
4058 due_at: None,
4059 not_before: None,
4060 snoozed_until: None,
4061 external_refs: Vec::new(),
4062 evidence_refs: Vec::new(),
4063 status: None,
4064 })
4065 .await
4066 .expect("create");
4067
4068 let event = WorkGraphEvent::graph(
4069 item.realm_id.clone(),
4070 item.namespace.clone(),
4071 WorkGraphEventKind::Created,
4072 item.created_at,
4073 json!({ "item_id": item.id }),
4074 );
4075 let error = store
4076 .insert_item(item, event)
4077 .await
4078 .expect_err("duplicate item insert must fail");
4079 assert!(
4080 matches!(error, WorkGraphError::Conflict(_)),
4081 "duplicate item insert must map to Conflict, got: {error:?}"
4082 );
4083 }
4084
4085 #[cfg(not(target_arch = "wasm32"))]
4089 #[tokio::test]
4090 async fn sqlite_store_duplicate_attention_insert_maps_to_conflict() {
4091 let dir = tempfile::tempdir().expect("tempdir");
4092 let path = dir.path().join("workgraph.sqlite3");
4093 let store = std::sync::Arc::new(crate::SqliteWorkGraphStore::open(&path).expect("open"));
4094 let service =
4095 WorkGraphService::with_scope(store.clone(), "realm", WorkNamespace::default());
4096 let goal = service
4097 .create_goal(GoalCreateRequest {
4098 realm_id: None,
4099 namespace: None,
4100 title: "unique goal".to_string(),
4101 description: None,
4102 target: GoalAttentionTarget::Session {
4103 session_id: meerkat_core::SessionId::new(),
4104 },
4105 mode: WorkAttentionMode::Coordinate,
4106 completion_policy: WorkCompletionPolicy::SelfAttest,
4107 delegated_authority: AttentionDelegatedAuthority::AddEvidence,
4108 projection_policy: AttentionProjectionPolicy::default(),
4109 })
4110 .await
4111 .expect("create goal");
4112
4113 let mut fresh_item = goal.item.clone();
4114 fresh_item.id = WorkItemId::generated();
4115 let item_event = WorkGraphEvent::graph(
4116 fresh_item.realm_id.clone(),
4117 fresh_item.namespace.clone(),
4118 WorkGraphEventKind::Created,
4119 fresh_item.created_at,
4120 json!({ "item_id": fresh_item.id }),
4121 );
4122 let attention_event = WorkGraphEvent::graph(
4123 goal.attention.work_ref.realm_id.clone(),
4124 goal.attention.work_ref.namespace.clone(),
4125 WorkGraphEventKind::AttentionCreated,
4126 goal.attention.updated_at,
4127 json!({ "binding_id": goal.attention.binding_id }),
4128 );
4129 let error = store
4130 .insert_goal(fresh_item, item_event, goal.attention, attention_event)
4131 .await
4132 .expect_err("duplicate attention insert must fail");
4133 assert!(
4134 matches!(error, WorkGraphError::Conflict(_)),
4135 "duplicate attention insert must map to Conflict, got: {error:?}"
4136 );
4137 }
4138
4139 #[cfg(not(target_arch = "wasm32"))]
4140 #[tokio::test]
4141 async fn sqlite_persistence_survives_restart() {
4142 let dir = tempfile::tempdir().expect("tempdir");
4143 let path = dir.path().join("workgraph.sqlite3");
4144 let store = std::sync::Arc::new(crate::SqliteWorkGraphStore::open(&path).expect("open"));
4145 let service = WorkGraphService::with_scope(store, "realm", WorkNamespace::default());
4146 let item = service
4147 .create(CreateWorkItemRequest {
4148 realm_id: None,
4149 namespace: None,
4150 title: "persist me".to_string(),
4151 description: None,
4152 priority: Default::default(),
4153 completion_policy: Default::default(),
4154 labels: BTreeSet::new(),
4155 due_at: None,
4156 not_before: None,
4157 snoozed_until: None,
4158 external_refs: Vec::new(),
4159 evidence_refs: Vec::new(),
4160 status: None,
4161 })
4162 .await
4163 .expect("create");
4164
4165 let reopened = std::sync::Arc::new(crate::SqliteWorkGraphStore::open(&path).expect("open"));
4166 let service = WorkGraphService::with_scope(reopened, "realm", WorkNamespace::default());
4167 let fetched = service.get(None, None, item.id.clone()).await.expect("get");
4168 assert_eq!(fetched.title, "persist me");
4169 }
4170
4171 #[cfg(not(target_arch = "wasm32"))]
4172 #[tokio::test]
4173 async fn sqlite_item_without_machine_state_fails_closed_on_read() {
4174 let dir = tempfile::tempdir().expect("tempdir");
4175 let path = dir.path().join("workgraph.sqlite3");
4176 let store = std::sync::Arc::new(crate::SqliteWorkGraphStore::open(&path).expect("open"));
4177 let service =
4178 WorkGraphService::with_scope(store.clone(), "realm", WorkNamespace::default());
4179 let item = service
4180 .create(CreateWorkItemRequest {
4181 realm_id: None,
4182 namespace: None,
4183 title: "legacy item".to_string(),
4184 description: None,
4185 priority: Default::default(),
4186 completion_policy: Default::default(),
4187 labels: BTreeSet::new(),
4188 due_at: None,
4189 not_before: None,
4190 snoozed_until: None,
4191 external_refs: Vec::new(),
4192 evidence_refs: Vec::new(),
4193 status: None,
4194 })
4195 .await
4196 .expect("create");
4197
4198 store
4199 .with_connection(|conn| {
4200 let json: String = conn
4201 .query_row(
4202 "SELECT item_json FROM workgraph_items
4203 WHERE realm_id = ?1 AND namespace = ?2 AND item_id = ?3",
4204 rusqlite::params![
4205 &item.realm_id,
4206 item.namespace.as_str(),
4207 item.id.as_str()
4208 ],
4209 |row| row.get(0),
4210 )
4211 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
4212 let mut value = serde_json::from_str::<serde_json::Value>(&json)
4213 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
4214 value
4215 .as_object_mut()
4216 .expect("item json object")
4217 .remove("machine_state");
4218 conn.execute(
4219 "UPDATE workgraph_items
4220 SET item_json = ?4
4221 WHERE realm_id = ?1 AND namespace = ?2 AND item_id = ?3",
4222 rusqlite::params![
4223 &item.realm_id,
4224 item.namespace.as_str(),
4225 item.id.as_str(),
4226 serde_json::to_string(&value)
4227 .map_err(|err| WorkGraphError::Store(err.to_string()))?
4228 ],
4229 )
4230 .map_err(|err| WorkGraphError::Store(err.to_string()))?;
4231 Ok(())
4232 })
4233 .expect("strip machine state");
4234
4235 let reopened = std::sync::Arc::new(crate::SqliteWorkGraphStore::open(&path).expect("open"));
4240 let service = WorkGraphService::with_scope(reopened, "realm", WorkNamespace::default());
4241 let err = service
4242 .get(None, None, item.id)
4243 .await
4244 .expect_err("reading an item with no machine_state must fail closed");
4245 assert!(
4246 matches!(err, WorkGraphError::Store(_)),
4247 "expected a typed Store deserialization error, got: {err:?}"
4248 );
4249 }
4250
4251 #[cfg(not(target_arch = "wasm32"))]
4252 #[tokio::test]
4253 async fn sqlite_event_replay_rebuilds_projection() {
4254 let dir = tempfile::tempdir().expect("tempdir");
4255 let path = dir.path().join("workgraph.sqlite3");
4256 let store = std::sync::Arc::new(crate::SqliteWorkGraphStore::open(&path).expect("open"));
4257 let service =
4258 WorkGraphService::with_scope(store.clone(), "realm", WorkNamespace::default());
4259 let blocker = service
4260 .create(CreateWorkItemRequest {
4261 realm_id: None,
4262 namespace: None,
4263 title: "blocker".to_string(),
4264 description: None,
4265 priority: Default::default(),
4266 completion_policy: Default::default(),
4267 labels: BTreeSet::new(),
4268 due_at: None,
4269 not_before: None,
4270 snoozed_until: None,
4271 external_refs: Vec::new(),
4272 evidence_refs: Vec::new(),
4273 status: None,
4274 })
4275 .await
4276 .expect("create blocker");
4277 let blocked = service
4278 .create(CreateWorkItemRequest {
4279 realm_id: None,
4280 namespace: None,
4281 title: "blocked".to_string(),
4282 description: None,
4283 priority: Default::default(),
4284 completion_policy: Default::default(),
4285 labels: BTreeSet::new(),
4286 due_at: None,
4287 not_before: None,
4288 snoozed_until: None,
4289 external_refs: Vec::new(),
4290 evidence_refs: Vec::new(),
4291 status: None,
4292 })
4293 .await
4294 .expect("create blocked");
4295 service
4296 .link(LinkWorkItemsRequest {
4297 realm_id: None,
4298 namespace: None,
4299 kind: WorkEdgeKind::Blocks,
4300 from_id: blocker.id.clone(),
4301 to_id: blocked.id.clone(),
4302 })
4303 .await
4304 .expect("link");
4305
4306 store
4307 .with_connection(|conn| {
4308 conn.execute("DELETE FROM workgraph_items", [])
4309 .map_err(|err| crate::WorkGraphError::Store(err.to_string()))?;
4310 conn.execute("DELETE FROM workgraph_edges", [])
4311 .map_err(|err| crate::WorkGraphError::Store(err.to_string()))?;
4312 Ok(())
4313 })
4314 .expect("clear projection");
4315
4316 let empty_items = store
4317 .list_items(WorkItemFilter {
4318 realm_id: Some("realm".to_string()),
4319 namespace: Some(WorkNamespace::default()),
4320 ..WorkItemFilter::default()
4321 })
4322 .await
4323 .expect("empty list");
4324 assert!(empty_items.is_empty());
4325
4326 store
4327 .rebuild_projection_from_events()
4328 .expect("rebuild projection");
4329
4330 let rebuilt_items = store
4331 .list_items(WorkItemFilter {
4332 realm_id: Some("realm".to_string()),
4333 namespace: Some(WorkNamespace::default()),
4334 ..WorkItemFilter::default()
4335 })
4336 .await
4337 .expect("rebuilt list");
4338 assert_eq!(rebuilt_items.len(), 2);
4339 let rebuilt_edges = store
4340 .list_edges("realm", &WorkNamespace::default())
4341 .await
4342 .expect("rebuilt edges");
4343 assert_eq!(rebuilt_edges.len(), 1);
4344 }
4345
4346 #[cfg(not(target_arch = "wasm32"))]
4347 #[tokio::test]
4348 async fn sqlite_event_replay_stops_attention_for_terminal_goal_items() {
4349 let dir = tempfile::tempdir().expect("tempdir");
4350 let path = dir.path().join("workgraph.sqlite3");
4351 let store = std::sync::Arc::new(crate::SqliteWorkGraphStore::open(&path).expect("open"));
4352 let service =
4353 WorkGraphService::with_scope(store.clone(), "realm", WorkNamespace::default());
4354 let session_id = meerkat_core::SessionId::parse("019e63c2-0000-7000-8000-000000000045")
4355 .expect("session id");
4356 let goal = service
4357 .create_goal(GoalCreateRequest {
4358 realm_id: None,
4359 namespace: None,
4360 title: "terminal goal".to_string(),
4361 description: None,
4362 target: GoalAttentionTarget::Session { session_id },
4363 mode: WorkAttentionMode::Pursue,
4364 completion_policy: WorkCompletionPolicy::SelfAttest,
4365 delegated_authority: AttentionDelegatedAuthority::CloseIfPolicyAllows,
4366 projection_policy: AttentionProjectionPolicy::default(),
4367 })
4368 .await
4369 .expect("create goal");
4370 service
4371 .goal_request_close(GoalRequestCloseRequest {
4372 binding_id: goal.attention.binding_id.clone(),
4373 realm_id: None,
4374 namespace: None,
4375 expected_revision: goal.item.revision,
4376 status: GoalTerminalStatus::Completed,
4377 })
4378 .await
4379 .expect("close goal");
4380
4381 store
4382 .with_connection(|conn| {
4383 conn.execute("DELETE FROM workgraph_items", [])
4384 .map_err(|err| crate::WorkGraphError::Store(err.to_string()))?;
4385 conn.execute("DELETE FROM workgraph_attention", [])
4386 .map_err(|err| crate::WorkGraphError::Store(err.to_string()))?;
4387 Ok(())
4388 })
4389 .expect("clear projection");
4390
4391 store
4392 .rebuild_projection_from_events()
4393 .expect("rebuild projection");
4394
4395 let binding = store
4396 .get_attention(
4397 "realm",
4398 &WorkNamespace::default(),
4399 &goal.attention.binding_id,
4400 )
4401 .await
4402 .expect("read binding")
4403 .expect("rebuilt binding");
4404 assert_eq!(binding.status, WorkAttentionStatus::Stopped);
4405 }
4406
4407 #[cfg(not(target_arch = "wasm32"))]
4408 #[tokio::test]
4409 async fn sqlite_store_duplicate_edge_does_not_append_event() {
4410 let dir = tempfile::tempdir().expect("tempdir");
4411 let path = dir.path().join("workgraph.sqlite3");
4412 let store = crate::SqliteWorkGraphStore::open(&path).expect("open");
4413 let edge = test_edge();
4414 store
4415 .insert_edge(edge.clone(), link_event(&edge))
4416 .await
4417 .expect("insert edge");
4418
4419 let error = store
4420 .insert_edge(edge.clone(), link_event(&edge))
4421 .await
4422 .expect_err("duplicate edge should fail");
4423 assert!(matches!(error, WorkGraphError::Conflict(_)));
4424
4425 let events = store
4426 .list_events(WorkGraphEventFilter {
4427 realm_id: Some(edge.realm_id),
4428 namespace: Some(edge.namespace),
4429 all_namespaces: false,
4430 after_seq: None,
4431 limit: None,
4432 })
4433 .await
4434 .expect("events");
4435 assert_eq!(events.len(), 1);
4436 }
4437}
4438
4439#[cfg(all(test, not(target_arch = "wasm32")))]
4440#[allow(clippy::expect_used, clippy::unwrap_used)]
4441mod legacy_schema_tests {
4442 use super::*;
4443 use crate::{AttentionDelegatedAuthority, AttentionProjectionPolicy, WorkAttentionMode};
4444 use meerkat_core::SessionId;
4445
4446 fn test_attention(binding_id: &str) -> WorkAttentionBinding {
4447 WorkAttentionBinding {
4448 binding_id: WorkAttentionBindingId::new(binding_id).expect("binding id"),
4449 work_ref: crate::WorkItemRef {
4450 realm_id: "realm".to_string(),
4451 namespace: WorkNamespace::default(),
4452 item_id: WorkItemId::generated(),
4453 },
4454 target: crate::WorkAttentionTarget::Session {
4455 session_id: SessionId::new(),
4456 },
4457 mode: WorkAttentionMode::Pursue,
4458 status: WorkAttentionStatus::Active,
4459 machine_state: Default::default(),
4460 delegated_authority: AttentionDelegatedAuthority::AddEvidence,
4461 projection_policy: AttentionProjectionPolicy::default(),
4462 created_at: chrono::Utc::now(),
4463 updated_at: chrono::Utc::now(),
4464 }
4465 }
4466
4467 fn create_unledgered_v2_workgraph(path: &Path) -> Connection {
4468 let mut conn = Connection::open(path).expect("open raw");
4469 let tx = conn.transaction().expect("begin schema transaction");
4470 migration_0001_workgraph_schema(&tx).expect("create v1 workgraph schema");
4471 migration_0002_attention_query_columns(&tx).expect("create v2 workgraph schema");
4472 tx.commit().expect("commit v2 workgraph schema");
4473 conn
4474 }
4475
4476 #[test]
4477 fn explicit_bridge_authenticates_v2_and_repairs_null_attention_projections() {
4478 let dir = tempfile::tempdir().expect("tempdir");
4479 let path = dir.path().join("workgraph.sqlite3");
4480 let mut conn = create_unledgered_v2_workgraph(&path);
4481 let binding = test_attention("legacy-v2-binding");
4482 let expected_status = binding.status.status_key().to_string();
4483 let expected_target_key = binding.target.target_key();
4484 conn.execute(
4485 "INSERT INTO workgraph_attention
4486 (realm_id, namespace, binding_id, revision, updated_at_utc, attention_json)
4487 VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
4488 params![
4489 binding.work_ref.realm_id,
4490 binding.work_ref.namespace.as_str(),
4491 binding.binding_id.as_str(),
4492 binding.machine_state.revision,
4493 binding.updated_at.to_rfc3339(),
4494 serde_json::to_string(&binding).expect("serialize binding"),
4495 ],
4496 )
4497 .expect("insert mixed-version row");
4498
4499 let report = meerkat_sqlite::bridge_unledgered_domain(
4500 &mut conn,
4501 &WORKGRAPH_DOMAIN,
4502 WORKGRAPH_DOMAIN.supported_version(),
4503 &[1, 2],
4504 Some(prepare_pre_0_8_10_workgraph_attention),
4505 )
4506 .expect("bridge exact v2 catalog");
4507 assert_eq!(report.from_version, 2);
4508 assert_eq!(report.to_version, 3);
4509 assert_eq!(report.prepared, 1);
4510 let projections = conn
4511 .query_row(
4512 "SELECT status, target_key FROM workgraph_attention WHERE binding_id = ?1",
4513 [binding.binding_id.as_str()],
4514 |row| Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?)),
4515 )
4516 .expect("read repaired projections");
4517 assert_eq!(projections, (expected_status, expected_target_key));
4518 assert_eq!(
4519 meerkat_sqlite::domain_version(&conn, WORKGRAPH_DOMAIN.name).expect("ledger"),
4520 Some(3)
4521 );
4522
4523 let rerun = meerkat_sqlite::bridge_unledgered_domain(
4524 &mut conn,
4525 &WORKGRAPH_DOMAIN,
4526 WORKGRAPH_DOMAIN.supported_version(),
4527 &[1, 2],
4528 Some(prepare_pre_0_8_10_workgraph_attention),
4529 )
4530 .expect("idempotent target rerun");
4531 assert_eq!(rerun.from_version, 3);
4532 assert_eq!(rerun.to_version, 3);
4533 assert_eq!(rerun.prepared, 0);
4534 }
4535
4536 #[test]
4537 fn explicit_bridge_refuses_non_null_attention_projection_mismatch_without_mutation() {
4538 for (case, wrong_status, wrong_target) in [
4539 ("status", Some("stopped"), None),
4540 ("target_key", None, Some("session:wrong")),
4541 ] {
4542 let dir = tempfile::tempdir().expect("tempdir");
4543 let path = dir.path().join(format!("workgraph-{case}.sqlite3"));
4544 let mut conn = create_unledgered_v2_workgraph(&path);
4545 let binding = test_attention(&format!("legacy-v2-{case}"));
4546 let expected_status = binding.status.status_key().to_string();
4547 let expected_target_key = binding.target.target_key();
4548 let source_status = wrong_status.unwrap_or(&expected_status).to_string();
4549 let source_target_key = wrong_target.unwrap_or(&expected_target_key).to_string();
4550 let source_json = serde_json::to_string(&binding).expect("serialize binding");
4551 conn.execute(
4552 "INSERT INTO workgraph_attention
4553 (realm_id, namespace, binding_id, revision, updated_at_utc, attention_json,
4554 status, target_key)
4555 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
4556 params![
4557 binding.work_ref.realm_id,
4558 binding.work_ref.namespace.as_str(),
4559 binding.binding_id.as_str(),
4560 binding.machine_state.revision,
4561 binding.updated_at.to_rfc3339(),
4562 source_json,
4563 source_status,
4564 source_target_key,
4565 ],
4566 )
4567 .expect("insert mismatched projection");
4568
4569 let error = meerkat_sqlite::bridge_unledgered_domain(
4570 &mut conn,
4571 &WORKGRAPH_DOMAIN,
4572 WORKGRAPH_DOMAIN.supported_version(),
4573 &[1, 2],
4574 Some(prepare_pre_0_8_10_workgraph_attention),
4575 )
4576 .expect_err("non-null projection mismatch must be refused");
4577 assert!(
4578 error.to_string().contains("disagrees with typed authority"),
4579 "unexpected {case} refusal: {error}"
4580 );
4581 let unchanged = conn
4582 .query_row(
4583 "SELECT attention_json, status, target_key
4584 FROM workgraph_attention WHERE binding_id = ?1",
4585 [binding.binding_id.as_str()],
4586 |row| {
4587 Ok((
4588 row.get::<_, String>(0)?,
4589 row.get::<_, String>(1)?,
4590 row.get::<_, String>(2)?,
4591 ))
4592 },
4593 )
4594 .expect("read refused source");
4595 assert_eq!(unchanged, (source_json, source_status, source_target_key));
4596 assert_eq!(
4597 meerkat_sqlite::domain_version(&conn, WORKGRAPH_DOMAIN.name).expect("ledger"),
4598 None,
4599 "refused {case} row must not be stamped"
4600 );
4601 }
4602 }
4603
4604 #[test]
4605 fn explicit_bridge_refuses_near_miss_v2_catalog_before_preparation() {
4606 let dir = tempfile::tempdir().expect("tempdir");
4607 let path = dir.path().join("workgraph.sqlite3");
4608 let mut conn = create_unledgered_v2_workgraph(&path);
4609 conn.execute_batch(
4610 "DROP INDEX idx_workgraph_attention_scope_status;
4611 CREATE INDEX idx_workgraph_attention_scope_status
4612 ON workgraph_attention (realm_id, namespace, status);",
4613 )
4614 .expect("install near-miss index");
4615
4616 let error = meerkat_sqlite::bridge_unledgered_domain(
4617 &mut conn,
4618 &WORKGRAPH_DOMAIN,
4619 WORKGRAPH_DOMAIN.supported_version(),
4620 &[1, 2],
4621 Some(prepare_pre_0_8_10_workgraph_attention),
4622 )
4623 .expect_err("near-miss catalog must be refused");
4624 assert!(
4625 error
4626 .to_string()
4627 .contains("does not match any authorized source catalog"),
4628 "unexpected near-miss refusal: {error}"
4629 );
4630 assert_eq!(
4631 meerkat_sqlite::domain_version(&conn, WORKGRAPH_DOMAIN.name).expect("ledger"),
4632 None
4633 );
4634 let index_sql: String = conn
4635 .query_row(
4636 "SELECT sql FROM sqlite_schema
4637 WHERE type = 'index' AND name = 'idx_workgraph_attention_scope_status'",
4638 [],
4639 |row| row.get(0),
4640 )
4641 .expect("near-miss index remains");
4642 assert!(index_sql.ends_with("(realm_id, namespace, status)"));
4643 }
4644
4645 #[tokio::test]
4648 async fn unledgered_legacy_attention_rows_are_refused_unmutated() {
4649 let dir = tempfile::tempdir().expect("tempdir");
4650 let path = dir.path().join("workgraph.sqlite3");
4651 let session_id = SessionId::new();
4652
4653 {
4656 let conn = Connection::open(&path).expect("open raw");
4657 conn.execute_batch(
4658 r"
4659 CREATE TABLE workgraph_attention (
4660 realm_id TEXT NOT NULL,
4661 namespace TEXT NOT NULL,
4662 binding_id TEXT NOT NULL,
4663 revision INTEGER NOT NULL,
4664 updated_at_utc TEXT NOT NULL,
4665 attention_json TEXT NOT NULL,
4666 PRIMARY KEY (realm_id, namespace, binding_id)
4667 );
4668 ",
4669 )
4670 .expect("create legacy table");
4671 let legacy = WorkAttentionBinding {
4672 binding_id: WorkAttentionBindingId::new("legacy-binding").expect("binding id"),
4673 work_ref: crate::WorkItemRef {
4674 realm_id: "realm".to_string(),
4675 namespace: WorkNamespace::default(),
4676 item_id: WorkItemId::generated(),
4677 },
4678 target: crate::WorkAttentionTarget::Session { session_id },
4679 mode: WorkAttentionMode::Pursue,
4680 status: WorkAttentionStatus::Active,
4681 machine_state: Default::default(),
4682 delegated_authority: AttentionDelegatedAuthority::AddEvidence,
4683 projection_policy: AttentionProjectionPolicy::default(),
4684 created_at: chrono::Utc::now(),
4685 updated_at: chrono::Utc::now(),
4686 };
4687 conn.execute(
4688 "INSERT INTO workgraph_attention
4689 (realm_id, namespace, binding_id, revision, updated_at_utc, attention_json)
4690 VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
4691 params![
4692 legacy.work_ref.realm_id,
4693 legacy.work_ref.namespace.as_str(),
4694 legacy.binding_id.as_str(),
4695 legacy.machine_state.revision,
4696 legacy.updated_at.to_rfc3339(),
4697 serde_json::to_string(&legacy).expect("serialize legacy binding"),
4698 ],
4699 )
4700 .expect("insert legacy row");
4701 }
4702
4703 let error = crate::SqliteWorkGraphStore::open(&path)
4704 .err()
4705 .expect("unledgered owned workgraph schema must be refused");
4706 assert!(
4707 error.to_string().contains("no ledger row"),
4708 "unexpected refusal: {error}"
4709 );
4710 let conn = Connection::open(&path).expect("reopen raw");
4711 let row_count: i64 = conn
4712 .query_row("SELECT COUNT(*) FROM workgraph_attention", [], |row| {
4713 row.get(0)
4714 })
4715 .expect("legacy row remains");
4716 assert_eq!(row_count, 1);
4717 let projected_columns: i64 = conn
4718 .query_row(
4719 "SELECT COUNT(*) FROM pragma_table_info('workgraph_attention')
4720 WHERE name IN ('status', 'target_key')",
4721 [],
4722 |row| row.get(0),
4723 )
4724 .expect("legacy columns");
4725 assert_eq!(projected_columns, 0);
4726 assert_eq!(
4727 meerkat_sqlite::domain_version(&conn, WORKGRAPH_DOMAIN.name).expect("ledger"),
4728 None
4729 );
4730 }
4731}