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)
| 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, |
no test coverage detected