Deliver one op to the session's delivery channel.
(&self, op: &ArrayOp, delivery: &ArrayDeliveryRegistry)
| 191 | |
| 192 | /// Deliver one op to the session's delivery channel. |
| 193 | fn deliver_op(&self, op: &ArrayOp, delivery: &ArrayDeliveryRegistry) { |
| 194 | let op_payload = match op_codec::encode_op(op) { |
| 195 | Ok(b) => b, |
| 196 | Err(e) => { |
| 197 | warn!( |
| 198 | session = %self.session_id, |
| 199 | array = %self.array, |
| 200 | error = %e, |
| 201 | "multi_shard_merger: encode_op failed — skipping op" |
| 202 | ); |
| 203 | return; |
| 204 | } |
| 205 | }; |
| 206 | |
| 207 | let msg = ArrayDeltaMsg { |
| 208 | array: op.header.array.clone(), |
| 209 | op_payload, |
| 210 | }; |
| 211 | |
| 212 | let frame = match nodedb_types::sync::wire::SyncFrame::try_encode( |
| 213 | SyncMessageType::ArrayDelta, |
| 214 | &msg, |
| 215 | ) { |
| 216 | Some(f) => f.to_bytes(), |
| 217 | None => { |
| 218 | warn!( |
| 219 | session = %self.session_id, |
| 220 | array = %self.array, |
| 221 | "multi_shard_merger: SyncFrame encode failed — skipping op" |
| 222 | ); |
| 223 | return; |
| 224 | } |
| 225 | }; |
| 226 | |
| 227 | delivery.enqueue(&self.session_id, frame); |
| 228 | } |
| 229 | } |
| 230 | |
| 231 | // ─── Registry ──────────────────────────────────────────────────────────────── |