use std::sync::Arc;
use std::time::Duration;
use meerkat::{
WorkGraphError, WorkGraphEvent, WorkGraphEventFilter, WorkGraphFact, WorkGraphService,
WorkGraphStoreKind, WorkItemId, WorkNamespace,
};
use serde::{Deserialize, Serialize};
use serde_json::json;
use tokio::sync::{Notify, broadcast};
use crate::types::{ModuleEvent, UnifiedEvent};
pub const DEFAULT_FACT_POLL_LIMIT: usize = 256;
pub const MAX_FACT_POLL_LIMIT: usize = 1000;
pub const DEFAULT_FACT_POLL_INTERVAL: Duration = Duration::from_secs(2);
pub const WORKGRAPH_FACTS_CHANNEL_CAP: usize = 512;
pub const WORKGRAPH_EVENT_MODULE: &str = "mobkit.workgraph";
pub const WORKGRAPH_FACT_EVENT_TYPE: &str = "workgraph.fact";
pub const WORKGRAPH_RESYNC_REQUIRED_EVENT_TYPE: &str = "workgraph.resync_required";
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkGraphFactEnvelope {
pub seq: i64,
pub realm_id: String,
pub namespace: WorkNamespace,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub item_id: Option<WorkItemId>,
pub fact: WorkGraphFact,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct WorkGraphFactPage {
pub facts: Vec<WorkGraphFactEnvelope>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub next_after_seq: Option<i64>,
pub more_may_be_pending: bool,
pub events_without_seq: usize,
}
impl WorkGraphFactPage {
pub fn is_empty(&self) -> bool {
self.facts.is_empty()
}
}
pub fn project_workgraph_fact_page(
events: &[WorkGraphEvent],
after_seq: Option<i64>,
limit: usize,
) -> WorkGraphFactPage {
let limit = limit.max(1);
let mut facts = Vec::new();
let mut next_after_seq = after_seq;
let mut events_without_seq = 0usize;
for event in events {
let Some(seq) = event.seq else {
events_without_seq = events_without_seq.saturating_add(1);
continue;
};
next_after_seq = Some(match next_after_seq {
Some(current) => current.max(seq),
None => seq,
});
for fact in &event.facts {
facts.push(WorkGraphFactEnvelope {
seq,
realm_id: event.realm_id.clone(),
namespace: event.namespace.clone(),
item_id: event.item_id.clone(),
fact: fact.clone(),
});
}
}
WorkGraphFactPage {
facts,
next_after_seq,
more_may_be_pending: events.len() >= limit,
events_without_seq,
}
}
pub async fn poll_workgraph_facts(
service: &WorkGraphService,
after_seq: Option<i64>,
limit: usize,
) -> Result<WorkGraphFactPage, WorkGraphError> {
let limit = limit.clamp(1, MAX_FACT_POLL_LIMIT);
let events = service
.events(WorkGraphEventFilter {
realm_id: None,
namespace: None,
all_namespaces: false,
after_seq,
limit: Some(limit),
})
.await?;
Ok(project_workgraph_fact_page(&events, after_seq, limit))
}
pub async fn latest_workgraph_fact_seq(
service: &WorkGraphService,
) -> Result<Option<i64>, WorkGraphError> {
service
.store()
.latest_event_seq(WorkGraphEventFilter {
realm_id: Some(service.default_realm_id().to_string()),
namespace: Some(service.default_namespace().clone()),
all_namespaces: false,
after_seq: None,
limit: None,
})
.await
}
pub fn workgraph_fact_event(fact: &WorkGraphFactEnvelope) -> UnifiedEvent {
UnifiedEvent::Module(ModuleEvent {
module: WORKGRAPH_EVENT_MODULE.to_string(),
event_type: WORKGRAPH_FACT_EVENT_TYPE.to_string(),
payload: serde_json::to_value(fact).unwrap_or_else(|_| {
json!({
"seq": fact.seq,
"projection_error": true,
})
}),
})
}
pub fn workgraph_resync_required_event(reason: &'static str, skipped: Option<u64>) -> UnifiedEvent {
UnifiedEvent::Module(ModuleEvent {
module: WORKGRAPH_EVENT_MODULE.to_string(),
event_type: WORKGRAPH_RESYNC_REQUIRED_EVENT_TYPE.to_string(),
payload: json!({
"reason": reason,
"skipped": skipped,
"authority": "durable_workgraph_pull",
}),
})
}
#[derive(Clone)]
pub struct WorkGraphFactHub {
tx: broadcast::Sender<UnifiedEvent>,
subscriber_arrived: Arc<Notify>,
tail_ready: Arc<Notify>,
}
impl Default for WorkGraphFactHub {
fn default() -> Self {
Self::new()
}
}
impl std::fmt::Debug for WorkGraphFactHub {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("WorkGraphFactHub")
.field("receiver_count", &self.tx.receiver_count())
.finish()
}
}
impl WorkGraphFactHub {
pub fn new() -> Self {
Self::with_capacity(WORKGRAPH_FACTS_CHANNEL_CAP)
}
pub fn with_capacity(capacity: usize) -> Self {
let (tx, _) = broadcast::channel(capacity.max(1));
Self {
tx,
subscriber_arrived: Arc::new(Notify::new()),
tail_ready: Arc::new(Notify::new()),
}
}
pub fn subscribe(&self) -> broadcast::Receiver<UnifiedEvent> {
let rx = self.tx.subscribe();
self.subscriber_arrived.notify_one();
rx
}
pub fn receiver_count(&self) -> usize {
self.tx.receiver_count()
}
pub async fn wait_for_subscriber(&self) {
while self.receiver_count() == 0 {
self.subscriber_arrived.notified().await;
}
}
#[doc(hidden)]
pub async fn wait_for_tail_ready(&self) {
self.tail_ready.notified().await;
}
fn mark_tail_ready(&self) {
self.tail_ready.notify_one();
}
fn publish(&self, event: UnifiedEvent) -> bool {
self.receiver_count() > 0 && self.tx.send(event).is_ok()
}
pub fn publish_fact(&self, fact: &WorkGraphFactEnvelope) -> bool {
self.publish(workgraph_fact_event(fact))
}
pub fn publish_page(&self, page: &WorkGraphFactPage) -> usize {
if self.receiver_count() == 0 {
return 0;
}
page.facts
.iter()
.filter(|fact| self.publish_fact(fact))
.count()
}
}
pub type WorkGraphFactStream = WorkGraphFactHub;
#[derive(Debug, Clone, Copy)]
pub struct WorkGraphFactTailOptions {
pub poll_interval: Duration,
pub page_limit: usize,
}
impl Default for WorkGraphFactTailOptions {
fn default() -> Self {
Self {
poll_interval: DEFAULT_FACT_POLL_INTERVAL,
page_limit: DEFAULT_FACT_POLL_LIMIT,
}
}
}
pub fn spawn_workgraph_fact_tail(
service: WorkGraphService,
hub: WorkGraphFactHub,
options: WorkGraphFactTailOptions,
) -> tokio::task::JoinHandle<()> {
tokio::spawn(async move {
loop {
hub.wait_for_subscriber().await;
let mut after_seq = match service.store().kind() {
WorkGraphStoreKind::Memory | WorkGraphStoreKind::Sqlite => {
match latest_workgraph_fact_seq(&service).await {
Ok(seq) => seq,
Err(error) => {
tracing::warn!(
%error,
"workgraph fact tail could not discover the ledger frontier",
);
None
}
}
}
WorkGraphStoreKind::Disabled => None,
WorkGraphStoreKind::Custom => {
let mut cursor = None;
while hub.receiver_count() > 0 {
match poll_workgraph_facts(&service, cursor, options.page_limit).await {
Ok(page) => {
cursor = page.next_after_seq;
if !page.more_may_be_pending {
break;
}
}
Err(error) => {
tracing::warn!(
%error,
"workgraph fact tail custom-store cursor discovery failed",
);
break;
}
}
tokio::time::sleep(options.poll_interval).await;
}
cursor
}
};
if hub.receiver_count() == 0 {
continue;
}
let mut first_live_poll_pending = true;
loop {
if hub.receiver_count() == 0 {
break;
}
let first_live_poll = first_live_poll_pending;
let polled_after_seq = after_seq;
match poll_workgraph_facts(&service, after_seq, options.page_limit).await {
Ok(page) => {
if page.events_without_seq > 0 {
tracing::warn!(
events_without_seq = page.events_without_seq,
"workgraph store returned events with no ledger sequence; \
their facts were not projected",
);
}
hub.publish_page(&page);
let cursor_advanced = page.next_after_seq != polled_after_seq;
after_seq = page.next_after_seq;
if page.more_may_be_pending && cursor_advanced {
tokio::task::yield_now().await;
continue;
}
}
Err(error) => {
tracing::warn!(
%error,
"workgraph fact poll failed; retrying after the poll interval",
);
}
}
if first_live_poll {
hub.mark_tail_ready();
first_live_poll_pending = false;
}
tokio::time::sleep(options.poll_interval).await;
}
}
})
}
#[cfg(test)]
#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
mod tests {
use super::*;
use futures::FutureExt;
use meerkat::WorkGraphEventKind;
fn item_id(value: &str) -> WorkItemId {
WorkItemId::new(value).expect("valid work item id")
}
fn ready_fact(id: &str, revision: u64) -> WorkGraphFact {
WorkGraphFact::ItemReady {
item_id: item_id(id),
item_revision: revision,
}
}
fn event(seq: Option<i64>, id: &str, facts: Vec<WorkGraphFact>) -> WorkGraphEvent {
let mut event = WorkGraphEvent::item(
"realm".to_string(),
WorkNamespace::default(),
item_id(id),
WorkGraphEventKind::Updated,
chrono::Utc::now(),
serde_json::Value::Null,
);
event.seq = seq;
event.facts = facts;
event
}
#[test]
fn fact_envelope_carries_identifiers_only() {
let envelope = WorkGraphFactEnvelope {
seq: 7,
realm_id: "realm".to_string(),
namespace: WorkNamespace::default(),
item_id: Some(item_id("work_a")),
fact: ready_fact("work_a", 3),
};
let value = serde_json::to_value(&envelope).expect("envelope serializes");
let object = value.as_object().expect("envelope is a JSON object");
assert!(object.contains_key("fact"), "fact must be projected");
let mut keys = object.keys().map(String::as_str).collect::<Vec<_>>();
keys.sort_unstable();
assert_eq!(keys, ["fact", "item_id", "namespace", "realm_id", "seq"]);
for forbidden in [
"item", "status", "owner", "claim", "evidence", "payload", "title", "revision",
] {
assert!(
!object.contains_key(forbidden),
"'{forbidden}' would let a consumer treat the accelerant as authority",
);
}
}
#[test]
fn projection_emits_one_envelope_per_fact_and_advances_the_cursor() {
let events = vec![
event(Some(4), "work_a", vec![ready_fact("work_a", 1)]),
event(
Some(9),
"work_b",
vec![ready_fact("work_b", 2), ready_fact("work_parent", 5)],
),
];
let page = project_workgraph_fact_page(&events, Some(1), 64);
assert_eq!(page.facts.len(), 3);
assert_eq!(page.next_after_seq, Some(9));
assert!(!page.more_may_be_pending);
assert_eq!(page.events_without_seq, 0);
assert_eq!(page.facts[0].seq, 4);
assert_eq!(page.facts[2].seq, 9);
assert_eq!(page.facts[2].item_id, Some(item_id("work_b")));
}
#[test]
fn events_without_facts_still_advance_the_cursor() {
let events = vec![event(Some(11), "work_a", Vec::new())];
let page = project_workgraph_fact_page(&events, Some(2), 64);
assert!(page.is_empty());
assert_eq!(page.next_after_seq, Some(11));
}
#[test]
fn seq_less_events_never_emit_and_never_stall_a_live_cursor() {
let events = vec![
event(None, "work_a", vec![ready_fact("work_a", 1)]),
event(Some(6), "work_b", vec![ready_fact("work_b", 1)]),
];
let page = project_workgraph_fact_page(&events, Some(3), 64);
assert_eq!(page.events_without_seq, 1);
assert_eq!(page.facts.len(), 1, "the seq-less fact must not be emitted");
assert_eq!(page.facts[0].seq, 6);
assert_eq!(page.next_after_seq, Some(6));
}
#[test]
fn a_page_of_only_seq_less_events_preserves_the_incoming_cursor() {
let events = vec![event(None, "work_a", vec![ready_fact("work_a", 1)])];
let page = project_workgraph_fact_page(&events, Some(3), 64);
assert!(page.is_empty());
assert_eq!(
page.next_after_seq,
Some(3),
"an unsequenced page must never rewind the cursor",
);
}
#[test]
fn the_cursor_never_moves_backwards() {
let events = vec![event(Some(2), "work_a", vec![ready_fact("work_a", 1)])];
let page = project_workgraph_fact_page(&events, Some(40), 64);
assert_eq!(page.next_after_seq, Some(40));
}
#[test]
fn a_full_page_reports_that_more_may_be_pending() {
let events = vec![
event(Some(1), "work_a", Vec::new()),
event(Some(2), "work_b", Vec::new()),
];
assert!(project_workgraph_fact_page(&events, None, 2).more_may_be_pending);
assert!(!project_workgraph_fact_page(&events, None, 3).more_may_be_pending);
assert!(!project_workgraph_fact_page(&[], None, 0).more_may_be_pending);
}
#[test]
fn a_full_page_of_seq_less_rows_reports_fullness_without_progress() {
let events = vec![
event(None, "work_a", vec![ready_fact("work_a", 1)]),
event(None, "work_b", vec![ready_fact("work_b", 1)]),
];
let page = project_workgraph_fact_page(&events, Some(12), 2);
assert!(page.more_may_be_pending, "a filled page reports fullness");
assert_eq!(
page.next_after_seq,
Some(12),
"no sequenced row was observed, so the cursor cannot have moved",
);
assert_eq!(page.events_without_seq, 2);
assert!(page.is_empty());
let sequenced = vec![
event(Some(13), "work_a", vec![ready_fact("work_a", 1)]),
event(Some(14), "work_b", vec![ready_fact("work_b", 1)]),
];
let progressed = project_workgraph_fact_page(&sequenced, Some(12), 2);
assert!(progressed.more_may_be_pending);
assert_eq!(progressed.next_after_seq, Some(14));
}
#[tokio::test]
async fn publishing_without_subscribers_is_not_an_error() {
let stream = WorkGraphFactStream::new();
let page = project_workgraph_fact_page(
&[event(Some(1), "work_a", vec![ready_fact("work_a", 1)])],
None,
64,
);
assert_eq!(stream.receiver_count(), 0);
assert_eq!(stream.publish_page(&page), 0);
}
#[test]
fn idle_hub_waits_for_notification_instead_of_scheduling_a_poll() {
let hub = WorkGraphFactHub::new();
assert_eq!(hub.receiver_count(), 0);
assert!(
hub.wait_for_subscriber().now_or_never().is_none(),
"zero-subscriber state must remain pending until subscribe notifies it",
);
}
#[tokio::test]
async fn subscribers_receive_projected_facts() {
let stream = WorkGraphFactStream::new();
let mut rx = stream.subscribe();
let page = project_workgraph_fact_page(
&[event(Some(5), "work_a", vec![ready_fact("work_a", 2)])],
None,
64,
);
assert_eq!(stream.publish_page(&page), 1);
let event = rx.try_recv().expect("published fact is delivered");
let UnifiedEvent::Module(module) = event else {
panic!("fact hub emitted a non-module event");
};
assert_eq!(module.module, WORKGRAPH_EVENT_MODULE);
assert_eq!(module.event_type, WORKGRAPH_FACT_EVENT_TYPE);
assert_eq!(module.payload["seq"], 5);
assert_eq!(
module.payload["fact"],
serde_json::to_value(ready_fact("work_a", 2)).expect("fact JSON"),
);
}
#[test]
fn tail_defaults_stay_inside_the_upstream_collection_ceiling() {
let options = WorkGraphFactTailOptions::default();
assert!(options.page_limit <= MAX_FACT_POLL_LIMIT);
const { assert!(DEFAULT_FACT_POLL_LIMIT <= MAX_FACT_POLL_LIMIT) };
assert!(!options.poll_interval.is_zero());
}
}