Skip to main content

common_procedure/
event.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use 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/// `ProcedureEvent` represents an event emitted by a procedure during its execution lifecycle.
37#[derive(Debug)]
38pub struct ProcedureEvent {
39    /// Unique identifier associated with the originating procedure instance.
40    pub procedure_id: ProcedureId,
41    /// The timestamp of the event.
42    pub timestamp: Timestamp,
43    /// The state of the procedure.
44    pub state: ProcedureState,
45    /// The lifecycle trigger that caused the event to be emitted.
46    pub trigger: EventTrigger,
47    /// Context associated with the root submission.
48    pub context: ProcedureContext,
49    /// The event emitted by the procedure. It's generated by [Procedure::event].
50    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}