Emit deferred trigger events for a completed transaction batch. Called after `execute_transaction_batch()` commits successfully. Each write in the transaction is emitted as a WriteEvent with `EventSource::Deferred`, which the Event Plane consumer routes to DEFERRED-mode triggers.
(
&mut self,
writes: Vec<DeferredWrite>,
tenant_id: crate::types::TenantId,
vshard_id: crate::types::VShardId,
)
| 28 | /// `EventSource::Deferred`, which the Event Plane consumer routes |
| 29 | /// to DEFERRED-mode triggers. |
| 30 | pub(in crate::data::executor) fn emit_deferred_events( |
| 31 | &mut self, |
| 32 | writes: Vec<DeferredWrite>, |
| 33 | tenant_id: crate::types::TenantId, |
| 34 | vshard_id: crate::types::VShardId, |
| 35 | ) { |
| 36 | let producer = match self.event_producer.as_mut() { |
| 37 | Some(p) => p, |
| 38 | None => return, |
| 39 | }; |
| 40 | |
| 41 | for write in writes { |
| 42 | self.event_sequence += 1; |
| 43 | |
| 44 | let (system_time_ms, valid_time_ms) = crate::event::bitemporal_extract::extract_stamps( |
| 45 | write.new_value.as_deref().or(write.old_value.as_deref()), |
| 46 | ); |
| 47 | let event = WriteEvent { |
| 48 | sequence: self.event_sequence, |
| 49 | collection: Arc::from(write.collection.as_str()), |
| 50 | op: write.op, |
| 51 | row_id: RowId::new(write.row_id.as_str()), |
| 52 | lsn: self.watermark, |
| 53 | tenant_id, |
| 54 | vshard_id, |
| 55 | source: EventSource::Deferred, |
| 56 | new_value: write.new_value.map(|v| Arc::from(v.as_slice())), |
| 57 | old_value: write.old_value.map(|v| Arc::from(v.as_slice())), |
| 58 | system_time_ms, |
| 59 | valid_time_ms, |
| 60 | user_id: None, |
| 61 | statement_digest: None, |
| 62 | }; |
| 63 | |
| 64 | producer.emit(event); |
| 65 | } |
| 66 | } |
| 67 | } |
no test coverage detected