(seq: u64, op: WriteOp, payload_bytes: Vec<u8>, is_delete: bool)
| 34 | } |
| 35 | |
| 36 | fn write_event(seq: u64, op: WriteOp, payload_bytes: Vec<u8>, is_delete: bool) -> WriteEvent { |
| 37 | let arc: Arc<[u8]> = Arc::from(payload_bytes.as_slice()); |
| 38 | let (system_time_ms, valid_time_ms) = extract_stamps(Some(&arc)); |
| 39 | WriteEvent { |
| 40 | sequence: seq, |
| 41 | collection: Arc::from("users"), |
| 42 | op, |
| 43 | row_id: RowId::new("u-1"), |
| 44 | lsn: Lsn::new(seq * 10), |
| 45 | tenant_id: TenantId::new(1), |
| 46 | vshard_id: VShardId::new(0), |
| 47 | source: EventSource::User, |
| 48 | new_value: if is_delete { None } else { Some(arc.clone()) }, |
| 49 | old_value: if is_delete { Some(arc) } else { None }, |
| 50 | system_time_ms, |
| 51 | valid_time_ms, |
| 52 | user_id: None, |
| 53 | statement_digest: None, |
| 54 | } |
| 55 | } |
| 56 | |
| 57 | fn stream_def() -> ChangeStreamDef { |
| 58 | ChangeStreamDef { |
no test coverage detected