1use std::any::Any;
16
17use api::v1::value::ValueData;
18use api::v1::{ColumnSchema, Row};
19use common_event_recorder::Event;
20use common_event_recorder::error::{Result, SerializeEventSnafu};
21use common_event_recorder::event_table::{
22 ACTOR_COLUMN, EVENT_CONTEXT_COLUMN, PROCEDURE_ERROR_COLUMN, PROCEDURE_ID_COLUMN,
23 PROCEDURE_STATE_COLUMN, PROCEDURE_TRIGGER_COLUMN, jsonb_value, nullable_json, nullable_string,
24 procedure_event_column_schemas,
25};
26use common_time::timestamp::{TimeUnit, Timestamp};
27use snafu::ResultExt;
28
29use crate::{EventTrigger, ProcedureContext, ProcedureId, ProcedureState};
30
31pub const EVENTS_TABLE_PROCEDURE_ID_COLUMN_NAME: &str = PROCEDURE_ID_COLUMN.name();
32pub const EVENTS_TABLE_PROCEDURE_STATE_COLUMN_NAME: &str = PROCEDURE_STATE_COLUMN.name();
33pub const EVENTS_TABLE_PROCEDURE_ERROR_COLUMN_NAME: &str = PROCEDURE_ERROR_COLUMN.name();
34pub const EVENTS_TABLE_PROCEDURE_TRIGGER_COLUMN_NAME: &str = PROCEDURE_TRIGGER_COLUMN.name();
35
36#[derive(Debug)]
38pub struct ProcedureEvent {
39 pub procedure_id: ProcedureId,
41 pub timestamp: Timestamp,
43 pub state: ProcedureState,
45 pub trigger: EventTrigger,
47 pub context: ProcedureContext,
49 pub internal_event: Box<dyn Event>,
51}
52
53impl ProcedureEvent {
54 pub fn new(
55 procedure_id: ProcedureId,
56 internal_event: Box<dyn Event>,
57 state: ProcedureState,
58 trigger: EventTrigger,
59 ) -> Self {
60 Self::new_with_context(
61 procedure_id,
62 internal_event,
63 state,
64 trigger,
65 ProcedureContext::default(),
66 )
67 }
68
69 pub fn new_with_context(
70 procedure_id: ProcedureId,
71 internal_event: Box<dyn Event>,
72 state: ProcedureState,
73 trigger: EventTrigger,
74 context: ProcedureContext,
75 ) -> Self {
76 Self {
77 procedure_id,
78 internal_event,
79 timestamp: Timestamp::current_time(TimeUnit::Nanosecond),
80 state,
81 trigger,
82 context,
83 }
84 }
85}
86
87impl Event for ProcedureEvent {
88 fn event_type(&self) -> &str {
89 self.internal_event.event_type()
90 }
91
92 fn timestamp(&self) -> Timestamp {
93 self.timestamp
94 }
95
96 fn json_payload(&self) -> Result<serde_json::Value> {
97 self.internal_event.json_payload()
98 }
99
100 fn extra_schema(&self) -> Vec<ColumnSchema> {
101 let mut schema = procedure_event_column_schemas();
102 let mut internal_schema = self.internal_event.extra_schema();
103 schema.append(&mut internal_schema);
104 schema.push(ACTOR_COLUMN.column_schema());
105 schema.push(EVENT_CONTEXT_COLUMN.column_schema());
106 schema
107 }
108
109 fn extra_rows(&self) -> Result<Vec<Row>> {
110 let mut internal_event_extra_rows = self.internal_event.extra_rows()?;
111 let mut rows = Vec::with_capacity(internal_event_extra_rows.len());
112 let procedure_id = self.procedure_id.to_string();
113 let state = self.state.as_str_name().to_string();
114 let error = match &self.state {
115 ProcedureState::Failed { error }
116 | ProcedureState::PrepareRollback { error }
117 | ProcedureState::RollingBack { error }
118 | ProcedureState::Retrying { error }
119 | ProcedureState::Poisoned { error, .. } => format!("{error:?}"),
120 _ => String::new(),
121 };
122 let trigger = serde_json::to_value(&self.trigger).context(SerializeEventSnafu)?;
123 let event_context = matches!(self.trigger, EventTrigger::Submitted)
124 .then_some(self.context.event_context.as_ref())
125 .flatten()
126 .map(serde_json::to_value)
127 .transpose()
128 .context(SerializeEventSnafu)?;
129 let actor = self.context.actor.as_deref();
130
131 for internal_event_extra_row in internal_event_extra_rows.iter_mut() {
132 let mut values = Vec::with_capacity(6 + internal_event_extra_row.values.len());
133 values.extend([
134 ValueData::StringValue(procedure_id.clone()).into(),
135 ValueData::StringValue(state.clone()).into(),
136 ValueData::StringValue(error.clone()).into(),
137 jsonb_value(&trigger),
138 ]);
139 values.append(&mut internal_event_extra_row.values);
140 values.push(nullable_string(actor));
141 values.push(nullable_json(event_context.as_ref()));
142 rows.push(Row { values });
143 }
144
145 Ok(rows)
146 }
147
148 fn as_any(&self) -> &dyn Any {
149 self
150 }
151}
152
153#[cfg(test)]
154mod tests {
155 use std::sync::Arc;
156
157 use api::v1::value::ValueData;
158 use api::v1::{ColumnDataType, ColumnSchema, Row, SemanticType, Value};
159 use common_error::mock::MockError;
160 use common_error::status_code::StatusCode;
161 use common_event_recorder::event_table::{
162 ACTOR_COLUMN, EVENT_CONTEXT_COLUMN, PROCEDURE_TRIGGER_COLUMN, jsonb_value,
163 };
164 use common_event_recorder::{Event, PersistentEventContext, TriggerReason};
165 use serde_json::json;
166
167 use crate::{
168 ChildSubmissionOutcome, Error, EventTrigger, ProcedureContext, ProcedureEvent, ProcedureId,
169 ProcedureState, RetryPhase,
170 };
171
172 #[derive(Debug)]
173 struct TestEvent;
174
175 impl Event for TestEvent {
176 fn event_type(&self) -> &str {
177 "test_event"
178 }
179
180 fn extra_schema(&self) -> Vec<ColumnSchema> {
181 vec![ColumnSchema {
182 column_name: "test_event_column".to_string(),
183 datatype: ColumnDataType::String.into(),
184 semantic_type: SemanticType::Field.into(),
185 ..Default::default()
186 }]
187 }
188
189 fn extra_rows(&self) -> common_event_recorder::error::Result<Vec<Row>> {
190 Ok(vec![
191 Row {
192 values: vec![ValueData::StringValue("test_event1".to_string()).into()],
193 },
194 Row {
195 values: vec![ValueData::StringValue("test_event2".to_string()).into()],
196 },
197 Row {
198 values: vec![Value { value_data: None }],
199 },
200 ])
201 }
202
203 fn as_any(&self) -> &dyn std::any::Any {
204 self
205 }
206 }
207
208 #[test]
209 fn procedure_event_extra_rows_preserve_envelope_values_and_internal_nulls() {
210 let procedure_id = ProcedureId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
211 let procedure_event = ProcedureEvent::new(
212 procedure_id,
213 Box::new(TestEvent {}),
214 ProcedureState::Running,
215 EventTrigger::Submitted,
216 );
217
218 let procedure_event_extra_rows = procedure_event.extra_rows().unwrap();
219 assert_eq!(
220 procedure_event_extra_rows,
221 vec![
222 Row {
223 values: vec![
224 ValueData::StringValue(procedure_id.to_string()).into(),
225 ValueData::StringValue("Running".to_string()).into(),
226 ValueData::StringValue(String::new()).into(),
227 jsonb_value(&json!({"type": "Submitted"})),
228 ValueData::StringValue("test_event1".to_string()).into(),
229 Value { value_data: None },
230 Value { value_data: None },
231 ],
232 },
233 Row {
234 values: vec![
235 ValueData::StringValue(procedure_id.to_string()).into(),
236 ValueData::StringValue("Running".to_string()).into(),
237 ValueData::StringValue(String::new()).into(),
238 jsonb_value(&json!({"type": "Submitted"})),
239 ValueData::StringValue("test_event2".to_string()).into(),
240 Value { value_data: None },
241 Value { value_data: None },
242 ],
243 },
244 Row {
245 values: vec![
246 ValueData::StringValue(procedure_id.to_string()).into(),
247 ValueData::StringValue("Running".to_string()).into(),
248 ValueData::StringValue(String::new()).into(),
249 jsonb_value(&json!({"type": "Submitted"})),
250 Value { value_data: None },
251 Value { value_data: None },
252 Value { value_data: None },
253 ],
254 },
255 ]
256 );
257 }
258
259 #[test]
260 fn procedure_event_extra_rows_include_error_for_failed_state() {
261 let error = Arc::new(Error::external(MockError::new(StatusCode::Unexpected)));
262 let procedure_event = ProcedureEvent::new(
263 ProcedureId::parse_str("00000000-0000-0000-0000-000000000001").unwrap(),
264 Box::new(TestEvent {}),
265 ProcedureState::failed(error.clone()),
266 EventTrigger::Failed,
267 );
268
269 let procedure_event_extra_rows = procedure_event.extra_rows().unwrap();
270
271 assert_eq!(procedure_event_extra_rows.len(), 3);
272 assert_eq!(
273 procedure_event_extra_rows[0].values[1],
274 ValueData::StringValue("Failed".to_string()).into()
275 );
276 assert_eq!(
277 procedure_event_extra_rows[0].values[2],
278 ValueData::StringValue(format!("{error:?}")).into()
279 );
280 assert_eq!(
281 procedure_event_extra_rows[0].values[3],
282 jsonb_value(&json!({"type": "Failed"}))
283 );
284 }
285
286 #[test]
287 fn test_procedure_event_extra_schema() {
288 let procedure_event = ProcedureEvent::new(
289 ProcedureId::random(),
290 Box::new(TestEvent {}),
291 ProcedureState::Running,
292 EventTrigger::Submitted,
293 );
294
295 assert_eq!(
296 procedure_event.extra_schema(),
297 vec![
298 ColumnSchema {
299 column_name: "procedure_id".to_string(),
300 datatype: ColumnDataType::String.into(),
301 semantic_type: SemanticType::Field.into(),
302 ..Default::default()
303 },
304 ColumnSchema {
305 column_name: "procedure_state".to_string(),
306 datatype: ColumnDataType::String.into(),
307 semantic_type: SemanticType::Field.into(),
308 ..Default::default()
309 },
310 ColumnSchema {
311 column_name: "procedure_error".to_string(),
312 datatype: ColumnDataType::String.into(),
313 semantic_type: SemanticType::Field.into(),
314 ..Default::default()
315 },
316 PROCEDURE_TRIGGER_COLUMN.column_schema(),
317 ColumnSchema {
318 column_name: "test_event_column".to_string(),
319 datatype: ColumnDataType::String.into(),
320 semantic_type: SemanticType::Field.into(),
321 ..Default::default()
322 },
323 ACTOR_COLUMN.column_schema(),
324 EVENT_CONTEXT_COLUMN.column_schema(),
325 ]
326 );
327 }
328
329 #[test]
330 fn procedure_event_records_actor_for_all_lifecycle_triggers() {
331 let mut context = ProcedureContext::from_event_context(PersistentEventContext::new(
332 TriggerReason::AutoRepartition,
333 ));
334 context.actor = Some("alice".to_string());
335 let submitted = ProcedureEvent::new_with_context(
336 ProcedureId::random(),
337 Box::new(TestEvent {}),
338 ProcedureState::Running,
339 EventTrigger::Submitted,
340 context.clone(),
341 );
342
343 assert_eq!(
344 submitted.extra_rows().unwrap()[0].values[5],
345 ValueData::StringValue("alice".to_string()).into()
346 );
347 assert_eq!(
348 submitted.extra_rows().unwrap()[0].values[6],
349 jsonb_value(&json!({"reason": "auto_repartition"}))
350 );
351
352 let child_id = ProcedureId::random();
353 for trigger in [
354 EventTrigger::Recovered,
355 EventTrigger::ChildSubmitted {
356 procedure_id: child_id,
357 outcome: ChildSubmissionOutcome::Accepted,
358 },
359 EventTrigger::Retrying {
360 phase: RetryPhase::Execute,
361 attempt: 1,
362 },
363 EventTrigger::RollingBack,
364 EventTrigger::Succeeded,
365 EventTrigger::Failed,
366 EventTrigger::Poisoned,
367 ] {
368 let event = ProcedureEvent::new_with_context(
369 ProcedureId::random(),
370 Box::new(TestEvent {}),
371 ProcedureState::Running,
372 trigger,
373 context.clone(),
374 );
375 let values = &event.extra_rows().unwrap()[0].values;
376 assert_eq!(
377 values[5],
378 ValueData::StringValue("alice".to_string()).into()
379 );
380 assert_eq!(values[6], Value { value_data: None });
381 }
382 }
383
384 #[test]
385 fn test_event_trigger_serialization() {
386 let procedure_id = ProcedureId::parse_str("00000000-0000-0000-0000-000000000001").unwrap();
387
388 for (trigger, name) in [
389 (EventTrigger::Submitted, "Submitted"),
390 (EventTrigger::Recovered, "Recovered"),
391 (EventTrigger::RollingBack, "RollingBack"),
392 (EventTrigger::Succeeded, "Succeeded"),
393 (EventTrigger::Failed, "Failed"),
394 (EventTrigger::Poisoned, "Poisoned"),
395 ] {
396 assert_eq!(
397 serde_json::to_value(trigger).unwrap(),
398 json!({"type": name})
399 );
400 }
401 assert_eq!(
402 serde_json::to_value(EventTrigger::ChildSubmitted {
403 procedure_id,
404 outcome: ChildSubmissionOutcome::Accepted,
405 })
406 .unwrap(),
407 json!({
408 "type": "ChildSubmitted",
409 "procedure_id": "00000000-0000-0000-0000-000000000001",
410 "outcome": "Accepted",
411 })
412 );
413 assert_eq!(
414 serde_json::to_value(EventTrigger::Retrying {
415 phase: RetryPhase::Execute,
416 attempt: 2,
417 })
418 .unwrap(),
419 json!({"type": "Retrying", "phase": "Execute", "attempt": 2})
420 );
421 }
422}