MCPcopy Create free account
hub / github.com/NodeDB-Lab/nodedb / package_and_enqueue

Method package_and_enqueue

nodedb/src/event/crdt_sync/packager.rs:83–146  ·  view source on GitHub ↗

Package a WriteEvent as an outbound delta and enqueue for delivery. Only packages events from `EventSource::User` (trigger side-effects are already captured by the triggering event's delta). Events from `CrdtSync` are skipped (prevent echo: Lite → Origin → Lite loop). Returns `true` if the delta was enqueued, `false` if skipped.

(&self, event: &WriteEvent, delivery: &CrdtSyncDelivery)

Source from the content-addressed store, hash-verified

81 ///
82 /// Returns `true` if the delta was enqueued, `false` if skipped.
83 pub fn package_and_enqueue(&self, event: &WriteEvent, delivery: &CrdtSyncDelivery) -> bool {
84 // Only package User-originated writes.
85 // CrdtSync events are inbound FROM Lite — don't echo back.
86 // Trigger/RaftFollower events are derivative — the original User
87 // event already covers the data change.
88 if event.source != EventSource::User {
89 return false;
90 }
91
92 // Check if any connected Lite session cares about this collection.
93 if !delivery.has_subscribers(event.tenant_id.as_u64(), &event.collection) {
94 self.deltas_skipped.fetch_add(1, Ordering::Relaxed);
95 return false;
96 }
97
98 let (op, payload) = match event.op {
99 WriteOp::Insert | WriteOp::Update => {
100 let payload = event
101 .new_value
102 .as_ref()
103 .map(|v| v.to_vec())
104 .unwrap_or_default();
105 (DeltaOp::Upsert, payload)
106 }
107 WriteOp::Delete => (DeltaOp::Delete, Vec::new()),
108 // Bulk operations: each row in the batch was already emitted
109 // as individual per-row events by the ring buffer path. The
110 // bulk event carries aggregate metadata, not per-row payloads.
111 WriteOp::BulkInsert { count } | WriteOp::BulkDelete { count } => {
112 trace!(
113 collection = %event.collection,
114 count,
115 "skipping bulk event for CRDT packaging (per-row events handle delivery)"
116 );
117 return false;
118 }
119 WriteOp::Heartbeat => return false,
120 };
121
122 let sequence = self.sequences.next(&event.collection);
123
124 let delta = OutboundDelta {
125 collection: event.collection.to_string(),
126 document_id: event.row_id.as_str().to_string(),
127 payload,
128 op,
129 lsn: event.lsn.as_u64(),
130 tenant_id: event.tenant_id.as_u64(),
131 peer_id: ORIGIN_PEER_ID.load(Ordering::Relaxed),
132 sequence,
133 };
134
135 delivery.enqueue(event.tenant_id.as_u64(), delta);
136 self.deltas_packaged.fetch_add(1, Ordering::Relaxed);
137
138 trace!(
139 collection = %event.collection,
140 doc_id = %event.row_id,

Callers 1

accumulate_data_eventFunction · 0.80

Calls 9

has_subscribersMethod · 0.80
to_stringMethod · 0.80
as_u64Method · 0.45
as_refMethod · 0.45
to_vecMethod · 0.45
nextMethod · 0.45
as_strMethod · 0.45
loadMethod · 0.45
enqueueMethod · 0.45

Tested by

no test coverage detected