Skip to main content

meerkat_workgraph/
store.rs

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    /// Insert one immutable execution binding after proving the referenced
147    /// WorkGraph item revision and retry-chain predecessor in the same store
148    /// transaction.
149    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    /// Resolve the unique execution binding for one target run in a realm.
168    /// This powers reverse linkage from Flow status without scanning a
169    /// bounded public binding list.
170    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    /// Enumerate only nonterminal execution obligations for host recovery.
195    /// Shipping stores override this with an active-queue projection so the
196    /// hot recovery path never scans historical terminal bindings.
197    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    /// Return at most `limit` attention rows. Backends should push this bound
220    /// into iteration/query ownership; the default is compatibility-only for
221    /// custom stores.
222    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    /// Delete TERMINAL (superseded/stopped) attention binding rows in scope.
233    /// The event stream keeps the audit history; binding rows otherwise grow
234    /// monotonically with reassignment churn. Returns the pruned row count.
235    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    /// Return at most `limit` edges in one namespace.
263    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    /// Return a bounded public event page while omitting internal execution
280    /// lifecycle events before applying the caller's visible limit.
281    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        // Custom stores get a single bounded read by default. Built-in stores
290        // override this method so visibility is filtered before applying the
291        // caller's limit.
292        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    /// Highest sequence matching a scope without retaining the event history.
302    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        // Occupancy over the post-reassign state: `previous` is being
860        // superseded in this same mutation, so it is excluded from the probe.
861        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        // Occupancy over the post-update state: exclude every binding this
923        // batch rewrites, then judge each Active-status update against the
924        // survivors plus its already-applied batch predecessors.
925        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    // The terminality verdict (which lifecycle phases are terminal) is a machine
1341    // fact owned by WorkGraphLifecycleMachine, not this filter. We drive the
1342    // machine's ClassifyTerminality over the item's recovered state and mirror the
1343    // verdict, failing closed: an item the machine cannot classify is treated as
1344    // terminal so it is never surfaced as live work when terminals are excluded.
1345    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        // Probe open: `with_connection` brings the schema domain up to date.
1438        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            // Rebuild is a whole-projection writer: it must acquire the write
1449            // lock before deleting projected rows so concurrent writers either
1450            // wait on busy_timeout or proceed after the rebuild commits.
1451            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        // Per-operation fence guard: lives exactly as long as the connection
1491        // it admits.
1492        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            // Occupancy after the row rewrite (the probe excludes the
1863            // candidate itself); a conflict drops the transaction, rolling
1864            // the rewrite back.
1865            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            // Occupancy over the post-reassign state: `previous` was just
1909            // rewritten to Superseded inside this transaction, so the probe
1910            // no longer sees it as active.
1911            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                // Occupancy after the row rewrite (the probe excludes the
1975                // candidate itself); a conflict drops the transaction.
1976                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            // Candidate scan is NULL-tolerant (rows written by older binaries
2020            // carry NULL status); each candidate is decoded and judged in
2021            // Rust before deletion, so only provably terminal rows go.
2022            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/// The workgraph store's schema domain in the per-file migration ledger.
2216///
2217/// Migration 0001 is the base DDL; 0002 lifts the historical attention
2218/// query-column upgrade (previously re-run on every open, idempotent only
2219/// via "duplicate column name" error matching) into a once-per-file,
2220/// transaction-wrapped migration with a `table_info` guard.
2221#[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/// One-time migration adding the indexed `status` / `target_key` query
2576/// columns to `workgraph_attention` (SQL filter pushdown + the
2577/// active-binding-per-target occupancy guard) and backfilling existing rows.
2578/// The current direct initializer composes this with the historical base DDL.
2579/// Ledger v1 is below the supported floor and is refused rather than inferred
2580/// or upgraded.
2581#[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/// Reconcile the exact v2 attention projection published before the schema
2646/// ledger existed.
2647///
2648/// The explicit maintenance bridge authenticates the catalog before calling
2649/// this function. A physical v1 source has no projection columns and is left
2650/// for migration 0002. A physical v2 source is data-bearing: every non-NULL
2651/// projection must already agree with its typed `attention_json` authority,
2652/// while NULL projections written by an older mixed-version process are
2653/// backfilled. Any disagreement is refused inside the bridge transaction.
2654#[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/// Occupancy probe for the active-binding-per-target invariant, run INSIDE
2778/// the same immediate write transaction as the mutation it guards so the
2779/// check is race-free next to the data. NULL-column rows (written by older
2780/// binaries) are decoded from JSON before judging, so mixed-version stores
2781/// cannot dodge the guard.
2782#[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
2826/// Typed conflict naming the occupant, so hosts get the invariant they were
2827/// building by hand (mobkit admission guards demote to defense-in-depth).
2828fn 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
2840/// Memory-store twin of [`active_target_occupant_tx`], run under the store's
2841/// write lock.
2842fn 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    // SQL filter pushdown over the indexed query columns. Every predicate is
3186    // NULL-tolerant: rows written by older binaries carry NULL status /
3187    // target_key and must still reach the Rust-side filter, which remains the
3188    // final authority over every returned row.
3189    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        // Terminality is a WorkGraph machine fact: the shell mirrors the
3536        // canonical classify verdict rather than re-deciding `is_terminal()`.
3537        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    /// Pins the SQLite UNIQUE-violation mapping for duplicate item inserts:
4039    /// a second insert of an existing item id must surface as the typed
4040    /// `Conflict`, not a generic `Store` error.
4041    #[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    /// Pins the SQLite UNIQUE-violation mapping for duplicate attention
4086    /// binding inserts (via the compound goal insert): the typed `Conflict`,
4087    /// not a generic `Store` error.
4088    #[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        // machine_state is the sole machine-owned lifecycle/revision authority.
4236        // A persisted item missing it can no longer be backfilled from projected
4237        // fields (that fabrication path was deleted); reading it must FAIL CLOSED
4238        // with a typed error rather than reconstructing machine truth.
4239        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    /// The released v2 floor is exact: an unledgered v1 attention table is
4646    /// refused without schema/data mutation or a ledger stamp.
4647    #[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        // Simulate the old binary: old-schema table + one active binding row
4654        // written without the query columns.
4655        {
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}