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

Function dispatch_trigger_batch

nodedb/src/event/trigger/dispatcher/batch.rs:22–117  ·  view source on GitHub ↗
(
    batch: &crate::control::trigger::batch::collector::TriggerBatch,
    state: &Arc<SharedState>,
    retry_queue: &mut TriggerRetryQueue,
)

Source from the content-addressed store, hash-verified

20use super::identity::trigger_identity;
21
22pub async fn dispatch_trigger_batch(
23 batch: &crate::control::trigger::batch::collector::TriggerBatch,
24 state: &Arc<SharedState>,
25 retry_queue: &mut TriggerRetryQueue,
26) {
27 use crate::control::security::catalog::trigger_types::{TriggerGranularity, TriggerTiming};
28 use crate::control::trigger::batch::when_filter;
29 use crate::control::trigger::fire_common;
30 use crate::control::trigger::registry::DmlEvent;
31
32 let tenant_id = TenantId::new(batch.tenant_id);
33 let identity = trigger_identity(tenant_id);
34 let mode_filter = Some(TriggerExecutionMode::Async);
35
36 let dml_event = match batch.operation.as_str() {
37 "INSERT" => DmlEvent::Insert,
38 "UPDATE" => DmlEvent::Update,
39 "DELETE" => DmlEvent::Delete,
40 _ => return,
41 };
42
43 let triggers =
44 state
45 .trigger_registry
46 .get_matching(batch.tenant_id, &batch.collection, dml_event);
47
48 let after_row_triggers: Vec<_> = triggers
49 .iter()
50 .filter(|t| t.timing == TriggerTiming::After)
51 .filter(|t| t.granularity == TriggerGranularity::Row)
52 .filter(|t| mode_filter.is_none() || Some(t.execution_mode) == mode_filter)
53 .collect();
54
55 if after_row_triggers.is_empty() {
56 return;
57 }
58
59 for trigger in &after_row_triggers {
60 let mask = when_filter::filter_batch_by_when(
61 &batch.rows,
62 &batch.collection,
63 &batch.operation,
64 trigger.when_condition.as_deref(),
65 );
66
67 let passing = when_filter::count_passing(&mask);
68 if passing == 0 {
69 continue;
70 }
71
72 for (row, &passes) in batch.rows.iter().zip(mask.iter()) {
73 if !passes {
74 continue;
75 }
76
77 let bindings =
78 when_filter::build_row_bindings(row, &batch.collection, &batch.operation);
79

Callers 1

process_normal_batchFunction · 0.85

Calls 15

trigger_identityFunction · 0.85
filter_batch_by_whenFunction · 0.85
count_passingFunction · 0.85
build_row_bindingsFunction · 0.85
fire_triggersFunction · 0.85
nowFunction · 0.85
get_matchingMethod · 0.80
collectMethod · 0.80
new_fieldsMethod · 0.80
old_fieldsMethod · 0.80
to_stringMethod · 0.80
as_strMethod · 0.45

Tested by

no test coverage detected