diff --git a/src/common/procedure/src/event.rs b/src/common/procedure/src/event.rs index 82e39f5900..9de3bc389c 100644 --- a/src/common/procedure/src/event.rs +++ b/src/common/procedure/src/event.rs @@ -20,11 +20,12 @@ use common_event_recorder::Event; use common_event_recorder::error::Result; use common_time::timestamp::{TimeUnit, Timestamp}; -use crate::{ProcedureId, ProcedureState}; +use crate::{EventTrigger, ProcedureId, ProcedureState}; pub const EVENTS_TABLE_PROCEDURE_ID_COLUMN_NAME: &str = "procedure_id"; pub const EVENTS_TABLE_PROCEDURE_STATE_COLUMN_NAME: &str = "procedure_state"; pub const EVENTS_TABLE_PROCEDURE_ERROR_COLUMN_NAME: &str = "procedure_error"; +pub const EVENTS_TABLE_PROCEDURE_TRIGGER_COLUMN_NAME: &str = "procedure_trigger"; /// `ProcedureEvent` represents an event emitted by a procedure during its execution lifecycle. #[derive(Debug)] @@ -35,6 +36,8 @@ pub struct ProcedureEvent { pub timestamp: Timestamp, /// The state of the procedure. pub state: ProcedureState, + /// The lifecycle trigger that caused the event to be emitted. + pub trigger: EventTrigger, /// The event emitted by the procedure. It's generated by [Procedure::event]. pub internal_event: Box, } @@ -44,12 +47,14 @@ impl ProcedureEvent { procedure_id: ProcedureId, internal_event: Box, state: ProcedureState, + trigger: EventTrigger, ) -> Self { Self { procedure_id, internal_event, timestamp: Timestamp::current_time(TimeUnit::Nanosecond), state, + trigger, } } } @@ -87,6 +92,12 @@ impl Event for ProcedureEvent { semantic_type: SemanticType::Field.into(), ..Default::default() }, + ColumnSchema { + column_name: EVENTS_TABLE_PROCEDURE_TRIGGER_COLUMN_NAME.to_string(), + datatype: ColumnDataType::String.into(), + semantic_type: SemanticType::Field.into(), + ..Default::default() + }, ]; schema.append(&mut self.internal_event.extra_schema()); schema @@ -95,20 +106,25 @@ impl Event for ProcedureEvent { fn extra_rows(&self) -> Result> { let mut internal_event_extra_rows = self.internal_event.extra_rows()?; let mut rows = Vec::with_capacity(internal_event_extra_rows.len()); + let procedure_id = self.procedure_id.to_string(); + let state = self.state.as_str_name().to_string(); + let error = match &self.state { + ProcedureState::Failed { error } + | ProcedureState::PrepareRollback { error } + | ProcedureState::RollingBack { error } + | ProcedureState::Retrying { error } + | ProcedureState::Poisoned { error, .. } => format!("{error:?}"), + _ => String::new(), + }; + let trigger = self.trigger.to_string(); + for internal_event_extra_row in internal_event_extra_rows.iter_mut() { - let error_str = match &self.state { - ProcedureState::Failed { error } => format!("{:?}", error), - ProcedureState::PrepareRollback { error } => format!("{:?}", error), - ProcedureState::RollingBack { error } => format!("{:?}", error), - ProcedureState::Retrying { error } => format!("{:?}", error), - ProcedureState::Poisoned { error, .. } => format!("{:?}", error), - _ => "".to_string(), - }; - let mut values = Vec::with_capacity(3 + internal_event_extra_row.values.len()); + let mut values = Vec::with_capacity(4 + internal_event_extra_row.values.len()); values.extend([ - ValueData::StringValue(self.procedure_id.to_string()).into(), - ValueData::StringValue(self.state.as_str_name().to_string()).into(), - ValueData::StringValue(error_str).into(), + ValueData::StringValue(procedure_id.clone()).into(), + ValueData::StringValue(state.clone()).into(), + ValueData::StringValue(error.clone()).into(), + ValueData::StringValue(trigger.clone()).into(), ]); values.append(&mut internal_event_extra_row.values); rows.push(Row { values }); @@ -128,7 +144,7 @@ mod tests { use api::v1::{ColumnDataType, ColumnSchema, Row, SemanticType}; use common_event_recorder::Event; - use crate::{ProcedureEvent, ProcedureId, ProcedureState}; + use crate::{EventTrigger, ProcedureEvent, ProcedureId, ProcedureState}; #[derive(Debug)] struct TestEvent; @@ -169,19 +185,49 @@ mod tests { ProcedureId::random(), Box::new(TestEvent {}), ProcedureState::Running, + EventTrigger::Submitted, ); let procedure_event_extra_rows = procedure_event.extra_rows().unwrap(); assert_eq!(procedure_event_extra_rows.len(), 2); - assert_eq!(procedure_event_extra_rows[0].values.len(), 4); + assert_eq!(procedure_event_extra_rows[0].values.len(), 5); assert_eq!( procedure_event_extra_rows[0].values[3], + ValueData::StringValue("Submitted".to_string()).into() + ); + assert_eq!( + procedure_event_extra_rows[0].values[4], ValueData::StringValue("test_event1".to_string()).into() ); - assert_eq!(procedure_event_extra_rows[1].values.len(), 4); + assert_eq!(procedure_event_extra_rows[1].values.len(), 5); assert_eq!( - procedure_event_extra_rows[1].values[3], + procedure_event_extra_rows[1].values[4], ValueData::StringValue("test_event2".to_string()).into() ); } + + #[test] + fn test_event_trigger_display() { + let procedure_id = ProcedureId::parse_str("00000000-0000-0000-0000-000000000001").unwrap(); + + assert_eq!(EventTrigger::Submitted.to_string(), "Submitted"); + assert_eq!(EventTrigger::Recovered.to_string(), "Recovered"); + assert_eq!( + EventTrigger::ChildSubmitted { + procedure_id, + outcome: crate::ChildSubmissionOutcome::Accepted, + } + .to_string(), + "ChildSubmitted(procedure_id=00000000-0000-0000-0000-000000000001, outcome=Accepted)" + ); + assert_eq!( + EventTrigger::Retrying { + phase: crate::RetryPhase::Execute, + attempt: 2, + } + .to_string(), + "Retrying(Execute, 2)" + ); + assert_eq!(EventTrigger::RollingBack.to_string(), "RollingBack"); + } } diff --git a/src/common/procedure/src/lib.rs b/src/common/procedure/src/lib.rs index 156a0ab78c..a1525ad242 100644 --- a/src/common/procedure/src/lib.rs +++ b/src/common/procedure/src/lib.rs @@ -29,9 +29,9 @@ pub mod test_util; pub use crate::error::{Error, Result}; pub use crate::event::ProcedureEvent; pub use crate::procedure::{ - BoxedProcedure, BoxedProcedureLoader, Context, ContextProvider, ContextProviderRef, LockKey, - Output, ParseIdError, PoisonKey, PoisonKeys, Procedure, ProcedureId, ProcedureInfo, - ProcedureManager, ProcedureManagerRef, ProcedureState, ProcedureWithId, Status, StringKey, - UserMetadata, + BoxedProcedure, BoxedProcedureLoader, ChildSubmissionOutcome, Context, ContextProvider, + ContextProviderRef, EventContext, EventTrigger, LockKey, Output, ParseIdError, PoisonKey, + PoisonKeys, Procedure, ProcedureId, ProcedureInfo, ProcedureManager, ProcedureManagerRef, + ProcedureState, ProcedureWithId, RetryPhase, Status, StringKey, }; pub use crate::watcher::Watcher; diff --git a/src/common/procedure/src/local.rs b/src/common/procedure/src/local.rs index 5e8717a53a..c704b6f8c0 100644 --- a/src/common/procedure/src/local.rs +++ b/src/common/procedure/src/local.rs @@ -37,20 +37,34 @@ use crate::error::{ Result, StartRemoveOutdatedMetaTaskSnafu, StopRemoveOutdatedMetaTaskSnafu, TooManyRunningProceduresSnafu, }; -use crate::event::ProcedureEvent; use crate::local::runner::Runner; use crate::procedure::{BoxedProcedureLoader, InitProcedureState, PoisonKeys, ProcedureInfo}; use crate::rwlock::{KeyRwLock, OwnedKeyRwLockGuard}; use crate::store::poison_store::PoisonStoreRef; use crate::store::{ProcedureMessage, ProcedureMessages, ProcedureStore, StateStoreRef}; use crate::{ - BoxedProcedure, ContextProvider, LockKey, PoisonKey, ProcedureId, ProcedureManager, - ProcedureState, ProcedureWithId, StringKey, UserMetadata, Watcher, + BoxedProcedure, ContextProvider, EventTrigger, LockKey, PoisonKey, ProcedureId, + ProcedureManager, ProcedureState, ProcedureWithId, StringKey, Watcher, }; /// The expired time of a procedure's metadata. const META_TTL: Duration = Duration::from_secs(60 * 10); +#[derive(Clone, Copy)] +enum RootSubmissionOrigin { + Fresh, + Recovery, +} + +impl RootSubmissionOrigin { + fn event_trigger(self) -> EventTrigger { + match self { + Self::Fresh => EventTrigger::Submitted, + Self::Recovery => EventTrigger::Recovered, + } + } +} + /// Shared metadata of a procedure. /// /// # Note @@ -83,10 +97,6 @@ pub(crate) struct ProcedureMeta { start_time_ms: AtomicI64, /// End execution time of this procedure. end_time_ms: AtomicI64, - /// Event recorder. - event_recorder: Option, - /// The user metadata of the procedure. It's generated by [Procedure::user_metadata]. - user_metadata: Option, } impl ProcedureMeta { @@ -98,8 +108,6 @@ impl ProcedureMeta { lock_key: LockKey, poison_keys: PoisonKeys, type_name: &str, - event_recorder: Option, - user_metadata: Option, ) -> ProcedureMeta { let (state_sender, state_receiver) = watch::channel(procedure_state); ProcedureMeta { @@ -114,8 +122,6 @@ impl ProcedureMeta { start_time_ms: AtomicI64::new(0), end_time_ms: AtomicI64::new(0), type_name: type_name.to_string(), - event_recorder, - user_metadata, } } @@ -126,14 +132,6 @@ impl ProcedureMeta { /// Update current [ProcedureState]. fn set_state(&self, state: ProcedureState) { - // Emit the event to the event recorder if the user metadata contains the eventable object. - if let (Some(event_recorder), Some(user_metadata)) = - (&self.event_recorder, &self.user_metadata) - && let Some(event) = user_metadata.to_event() - { - event_recorder.record(Box::new(ProcedureEvent::new(self.id, event, state.clone()))); - } - // Safety: ProcedureMeta also holds the receiver, so `send()` should never fail. self.state_sender.send(state).unwrap(); } @@ -693,10 +691,10 @@ impl LocalManager { procedure_state: ProcedureState, step: u32, procedure: BoxedProcedure, + origin: RootSubmissionOrigin, ) -> Result { ensure!(self.manager_ctx.running(), ManagerNotStartSnafu); - let user_metadata = procedure.user_metadata(); let meta = Arc::new(ProcedureMeta::new( procedure_id, procedure_state, @@ -704,8 +702,6 @@ impl LocalManager { procedure.lock_key(), procedure.poison_keys(), procedure.type_name(), - self.event_recorder.clone(), - user_metadata.clone(), )); let runner = Runner { meta: meta.clone(), @@ -718,19 +714,10 @@ impl LocalManager { store: self.procedure_store.clone(), rolling_back: false, event_recorder: self.event_recorder.clone(), + execute_retry_attempt: 0, + rollback_retry_attempt: 0, }; - if let (Some(event_recorder), Some(event)) = ( - self.event_recorder.as_ref(), - user_metadata.and_then(|m| m.to_event()), - ) { - event_recorder.record(Box::new(ProcedureEvent::new( - procedure_id, - event, - ProcedureState::Running, - ))); - } - let watcher = meta.state_receiver.clone(); ensure!( @@ -746,6 +733,8 @@ impl LocalManager { DuplicateProcedureSnafu { procedure_id }, ); + runner.record_event(origin.event_trigger()); + let tracing_context = TracingContext::from_current_span(); ensure!( @@ -816,6 +805,7 @@ impl LocalManager { procedure_state, loaded_procedure.step, loaded_procedure.procedure, + RootSubmissionOrigin::Recovery, ) { error!(e; "Failed to recover procedure {}", procedure_id); } @@ -942,6 +932,7 @@ impl ProcedureManager for LocalManager { ProcedureState::Running, 0, procedure.procedure, + RootSubmissionOrigin::Fresh, ) } @@ -992,8 +983,6 @@ pub(crate) mod test_util { LockKey::default(), PoisonKeys::default(), "ProcedureAdapter", - None, - None, ) } @@ -1007,10 +996,12 @@ pub(crate) mod test_util { #[cfg(test)] mod tests { use std::assert_matches; + use std::sync::Mutex; use std::sync::atomic::{AtomicBool, Ordering as AtomicOrdering}; use common_error::mock::MockError; use common_error::status_code::StatusCode; + use common_event_recorder::{Event, EventRecorder}; use common_test_util::temp_dir::create_temp_dir; use tokio::sync::oneshot; use tokio::time::timeout; @@ -1019,13 +1010,57 @@ mod tests { use crate::error::{self, Error}; use crate::store::state_store::ObjectStateStore; use crate::test_util::InMemoryPoisonStore; - use crate::{Context, Procedure, Status}; + use crate::{Context, EventContext, EventTrigger, Procedure, ProcedureEvent, Status}; fn new_test_manager_context() -> ManagerContext { let poison_manager = Arc::new(InMemoryPoisonStore::default()); ManagerContext::new(poison_manager) } + #[derive(Debug, Default)] + struct CapturingEventRecorder { + events: Mutex>>, + } + + impl CapturingEventRecorder { + fn triggers(&self) -> Vec { + self.events + .lock() + .unwrap() + .iter() + .map(|event| { + event + .as_any() + .downcast_ref::() + .unwrap() + .trigger + .clone() + }) + .collect() + } + } + + impl EventRecorder for CapturingEventRecorder { + fn record(&self, event: Box) { + self.events.lock().unwrap().push(event); + } + + fn close(&self) {} + } + + #[derive(Debug)] + struct TestProcedureEvent; + + impl Event for TestProcedureEvent { + fn event_type(&self) -> &str { + "test_procedure" + } + + fn as_any(&self) -> &dyn std::any::Any { + self + } + } + #[test] fn test_manager_context() { let ctx = new_test_manager_context(); @@ -1175,6 +1210,10 @@ mod tests { fn poison_keys(&self) -> PoisonKeys { self.poison_keys.clone() } + + fn event(&self, _ctx: &EventContext<'_>) -> Option> { + Some(Box::new(TestProcedureEvent)) + } } impl ProcedureToLoad { @@ -1195,6 +1234,80 @@ mod tests { } } + #[tokio::test] + async fn test_fresh_submission_emits_submitted_event() { + let dir = create_temp_dir("fresh_submission_event"); + let state_store = Arc::new(ObjectStateStore::new(test_util::new_object_store(&dir))); + let poison_manager = Arc::new(InMemoryPoisonStore::new()); + let event_recorder = Arc::new(CapturingEventRecorder::default()); + let manager = LocalManager::new( + ManagerConfig::default(), + state_store, + poison_manager, + None, + Some(event_recorder.clone()), + ); + manager.manager_ctx.start(); + + manager + .submit(ProcedureWithId { + id: ProcedureId::random(), + procedure: Box::new(ProcedureToLoad::new("fresh submission")), + }) + .await + .unwrap(); + + assert!(event_recorder.triggers().contains(&EventTrigger::Submitted)); + } + + #[tokio::test] + async fn test_recovery_emits_recovered_event() { + let dir = create_temp_dir("recovery_submission_event"); + let object_store = test_util::new_object_store(&dir); + let state_store = Arc::new(ObjectStateStore::new(object_store.clone())); + let poison_manager = Arc::new(InMemoryPoisonStore::new()); + let event_recorder = Arc::new(CapturingEventRecorder::default()); + let manager = LocalManager::new( + ManagerConfig { + parent_path: "data/".to_string(), + ..Default::default() + }, + state_store, + poison_manager, + None, + Some(event_recorder.clone()), + ); + manager.manager_ctx.start(); + manager + .register_loader("ProcedureToLoad", ProcedureToLoad::loader()) + .unwrap(); + + let procedure = ProcedureToLoad::new("recovered submission"); + let procedure_id = ProcedureId::random(); + ProcedureStore::from_object_store(object_store) + .store_procedure( + procedure_id, + 0, + procedure.type_name().to_string(), + procedure.dump().unwrap(), + None, + ) + .await + .unwrap(); + + manager.recover().await.unwrap(); + + assert!( + manager + .procedure_state(procedure_id) + .await + .unwrap() + .is_some() + ); + assert!(event_recorder.triggers().contains(&EventTrigger::Recovered)); + assert!(!event_recorder.triggers().contains(&EventTrigger::Submitted)); + } + #[derive(Debug)] struct BlockingProcedure { started_tx: Option>, diff --git a/src/common/procedure/src/local/runner.rs b/src/common/procedure/src/local/runner.rs index 509b3a7756..dff7990f82 100644 --- a/src/common/procedure/src/local/runner.rs +++ b/src/common/procedure/src/local/runner.rs @@ -33,7 +33,8 @@ use crate::procedure::{Output, StringKey}; use crate::rwlock::OwnedKeyRwLockGuard; use crate::store::{ProcedureMessage, ProcedureStore}; use crate::{ - BoxedProcedure, Context, Error, Procedure, ProcedureId, ProcedureState, ProcedureWithId, Status, + BoxedProcedure, ChildSubmissionOutcome, Context, Error, EventContext, EventTrigger, Procedure, + ProcedureId, ProcedureState, ProcedureWithId, RetryPhase, Status, }; /// A guard to cleanup procedure state. @@ -141,6 +142,8 @@ pub(crate) struct Runner { pub(crate) store: Arc, pub(crate) rolling_back: bool, pub(crate) event_recorder: Option, + pub(crate) execute_retry_attempt: u32, + pub(crate) rollback_retry_attempt: u32, } impl Runner { @@ -239,7 +242,7 @@ impl Runner { loop { // Don't store state if `ProcedureManager` is stopped. if !self.running() { - self.meta.set_state(ProcedureState::failed(Arc::new( + self.set_state_and_record(ProcedureState::failed(Arc::new( error::ManagerNotStartSnafu {}.build(), ))); return; @@ -270,14 +273,14 @@ impl Runner { | ProcedureState::RollingBack { error } => { rollback_times += 1; if let Some(d) = rollback.next() { - self.wait_on_err(d, rollback_times).await; + self.wait_on_err(d, rollback_times as u64).await; } else { let err = Err::<(), Arc>(error) .context(RollbackTimesExceededSnafu { procedure_id: self.meta.id, }) .unwrap_err(); - self.meta.set_state(ProcedureState::failed(Arc::new(err))); + self.set_state_and_record(ProcedureState::failed(Arc::new(err))); return; } } @@ -315,11 +318,10 @@ impl Runner { if self.procedure.rollback_supported() && let Err(e) = self.procedure.rollback(ctx).await { - self.meta - .set_state(ProcedureState::rolling_back(Arc::new(e))); + self.set_state_and_record(ProcedureState::rolling_back(Arc::new(e))); return; } - self.meta.set_state(ProcedureState::failed(err)); + self.set_state_and_record(ProcedureState::failed(err)); } async fn prepare_rollback(&mut self, err: Arc) { @@ -329,9 +331,9 @@ impl Runner { return; } if self.procedure.rollback_supported() { - self.meta.set_state(ProcedureState::rolling_back(err)); + self.set_state_and_record(ProcedureState::rolling_back(err)); } else { - self.meta.set_state(ProcedureState::failed(err)); + self.set_state_and_record(ProcedureState::failed(err)); } } @@ -350,7 +352,7 @@ impl Runner { // Don't store state if `ProcedureManager` is stopped. if !self.running() { - self.meta.set_state(ProcedureState::failed(Arc::new( + self.set_state_and_record(ProcedureState::failed(Arc::new( error::ManagerNotStartSnafu {}.build(), ))); return; @@ -361,7 +363,7 @@ impl Runner { && let Err(e) = self.clean_poisons().await { error!(e; "Failed to clean poison for procedure: {}", self.meta.id); - self.meta.set_state(ProcedureState::retrying(Arc::new(e))); + self.set_state_and_record(ProcedureState::retrying(Arc::new(e))); return; } @@ -369,7 +371,7 @@ impl Runner { && let Err(e) = self.persist_procedure().await { error!(e; "Failed to persist procedure: {}", self.meta.id); - self.meta.set_state(ProcedureState::retrying(Arc::new(e))); + self.set_state_and_record(ProcedureState::retrying(Arc::new(e))); return; } @@ -402,7 +404,9 @@ impl Runner { Status::Done { output } => { if let Err(e) = self.commit_procedure().await { error!(e; "Failed to commit procedure: {}", self.meta.id); - self.meta.set_state(ProcedureState::retrying(Arc::new(e))); + self.set_state_and_record(ProcedureState::retrying(Arc::new( + e, + ))); return; } @@ -416,8 +420,10 @@ impl Runner { self.meta.id, keys, ); - self.meta - .set_state(ProcedureState::poisoned(keys, Arc::new(error))); + self.set_state_and_record(ProcedureState::poisoned( + keys, + Arc::new(error), + )); } } } @@ -433,7 +439,7 @@ impl Runner { // Don't store state if `ProcedureManager` is stopped. if !self.running() { - self.meta.set_state(ProcedureState::failed(Arc::new( + self.set_state_and_record(ProcedureState::failed(Arc::new( error::ManagerNotStartSnafu {}.build(), ))); return; @@ -442,7 +448,7 @@ impl Runner { if e.need_clean_poisons() { if let Err(e) = self.clean_poisons().await { error!(e; "Failed to clean poison for procedure: {}", self.meta.id); - self.meta.set_state(ProcedureState::retrying(Arc::new(e))); + self.set_state_and_record(ProcedureState::retrying(Arc::new(e))); return; } debug!( @@ -453,7 +459,7 @@ impl Runner { } if e.is_retry_later() { - self.meta.set_state(ProcedureState::retrying(Arc::new(e))); + self.set_state_and_record(ProcedureState::retrying(Arc::new(e))); return; } @@ -461,7 +467,7 @@ impl Runner { self.meta .set_state(ProcedureState::prepare_rollback(Arc::new(e))); } else { - self.meta.set_state(ProcedureState::failed(Arc::new(e))); + self.set_state_and_record(ProcedureState::failed(Arc::new(e))); } } } @@ -480,19 +486,19 @@ impl Runner { procedure_id: ProcedureId, procedure_state: ProcedureState, procedure: BoxedProcedure, - ) { + ) -> ChildSubmissionOutcome { if !self.running() { warn!( "ProcedureManager is not running, skip submitting subprocedure {}-{}", procedure.type_name(), procedure_id ); - return; + return ChildSubmissionOutcome::ManagerStopped; } if self.manager_ctx.contains_procedure(procedure_id) { // If the parent has already submitted this procedure, don't submit it again. - return; + return ChildSubmissionOutcome::AlreadyAccepted; } let step = 0; @@ -504,8 +510,6 @@ impl Runner { procedure.lock_key(), procedure.poison_keys(), procedure.type_name(), - self.event_recorder.clone(), - procedure.user_metadata(), )); let runner = Runner { meta: meta.clone(), @@ -516,6 +520,8 @@ impl Runner { store: self.store.clone(), rolling_back: false, event_recorder: self.event_recorder.clone(), + execute_retry_attempt: 0, + rollback_retry_attempt: 0, }; // Insert the procedure. We already check the procedure existence before inserting @@ -548,11 +554,12 @@ impl Runner { }) }) { self.manager_ctx.remove_procedure(procedure_id); - return; + return ChildSubmissionOutcome::SpawnFailed; } // Add the id of the subprocedure to the metadata. self.meta.push_child(procedure_id); + ChildSubmissionOutcome::Accepted } /// Extend the retry time to wait for the next retry. @@ -598,7 +605,7 @@ impl Runner { if self.procedure.rollback_supported() { self.meta.set_state(ProcedureState::prepare_rollback(err)); } else { - self.meta.set_state(ProcedureState::failed(err)); + self.set_state_and_record(ProcedureState::failed(err)); } return; } @@ -613,11 +620,16 @@ impl Runner { subprocedure.id, ); - self.submit_subprocedure( + let child_id = subprocedure.id; + let outcome = self.submit_subprocedure( subprocedure.id, ProcedureState::Running, subprocedure.procedure, ); + self.record_event(EventTrigger::ChildSubmitted { + procedure_id: child_id, + outcome, + }); } info!( @@ -705,7 +717,7 @@ impl Runner { Ok(()) } - fn done(&self, output: Option) { + fn done(&mut self, output: Option) { // TODO(yingwen): Add files to remove list. info!( "Procedure {}-{} done", @@ -714,7 +726,65 @@ impl Runner { ); // Mark the state of this procedure to done. - self.meta.set_state(ProcedureState::Done { output }); + self.set_state_and_record(ProcedureState::Done { output }); + } + + /// Updates framework state and records the lifecycle event implied by that state. + fn set_state_and_record(&mut self, state: ProcedureState) { + let trigger = match &state { + ProcedureState::Retrying { .. } => { + self.execute_retry_attempt += 1; + Some(EventTrigger::Retrying { + phase: RetryPhase::Execute, + attempt: self.execute_retry_attempt, + }) + } + ProcedureState::RollingBack { .. } => { + if self.rolling_back { + self.rollback_retry_attempt += 1; + Some(EventTrigger::Retrying { + phase: RetryPhase::Rollback, + attempt: self.rollback_retry_attempt, + }) + } else { + self.rolling_back = true; + Some(EventTrigger::RollingBack) + } + } + ProcedureState::Done { .. } => Some(EventTrigger::Succeeded), + ProcedureState::Failed { .. } => Some(EventTrigger::Failed), + ProcedureState::Poisoned { .. } => Some(EventTrigger::Poisoned), + ProcedureState::Running => None, + // We don't record the prepare rollback state. + // The final result will be recorded as ProcedureState::Failed or ProcedureState::Succeeded. + ProcedureState::PrepareRollback { .. } => None, + }; + self.meta.set_state(state); + if let Some(trigger) = trigger { + self.record_event(trigger); + } + } + + /// Builds and dispatches an event from the live procedure. Delivery is best effort and is + /// intentionally not part of procedure execution or persistence. + pub(crate) fn record_event(&self, trigger: EventTrigger) { + let Some(recorder) = self.event_recorder.as_ref() else { + return; + }; + let state = self.meta.state(); + let context = EventContext { + procedure_id: self.meta.id, + lifecycle_state: &state, + trigger: trigger.clone(), + }; + if let Some(event) = self.procedure.event(&context) { + recorder.record(Box::new(crate::event::ProcedureEvent::new( + self.meta.id, + event, + state, + trigger, + ))); + } } } @@ -767,6 +837,8 @@ mod tests { store, rolling_back: false, event_recorder: None, + execute_retry_attempt: 0, + rollback_retry_attempt: 0, } } diff --git a/src/common/procedure/src/procedure.rs b/src/common/procedure/src/procedure.rs index 8e34f4bb30..6d07b160b6 100644 --- a/src/common/procedure/src/procedure.rs +++ b/src/common/procedure/src/procedure.rs @@ -19,7 +19,7 @@ use std::str::FromStr; use std::sync::Arc; use async_trait::async_trait; -use common_event_recorder::{Event, Eventable}; +use common_event_recorder::Event; use serde::{Deserialize, Serialize}; use smallvec::{SmallVec, smallvec}; use snafu::{ResultExt, Snafu}; @@ -228,27 +228,117 @@ pub trait Procedure: Send { PoisonKeys::default() } - /// Returns the user metadata of the procedure. If the metadata contains the eventable object, you can use [UserMetadata::to_event] to get the event and emit it to the event recorder. - fn user_metadata(&self) -> Option { + /// Builds an event for a framework lifecycle trigger. + /// + /// The hook is called with the current procedure instance, so an event can + /// include state that was produced after the procedure was submitted. A + /// return value of `None` means that this trigger should not be recorded. + /// Events that share an [`Event::event_type`] must return identical + /// [`Event::extra_schema`] values. The event recorder batches events by type + /// and rejects incompatible schemas; use a distinct event type for a + /// different schema. + fn event(&self, _ctx: &EventContext<'_>) -> Option> { None } } -/// The user metadata injected by the procedure caller. It can be used to emit events to the event recorder. -#[derive(Clone, Debug)] -pub struct UserMetadata { - event_object: Arc, +/// Framework-owned context supplied when a procedure builds a lifecycle event. +pub struct EventContext<'a> { + /// Id of the procedure associated with the event. + pub procedure_id: ProcedureId, + /// Current framework state of the procedure. + pub lifecycle_state: &'a ProcedureState, + /// Lifecycle action that caused the event hook to be called. + pub trigger: EventTrigger, } -impl UserMetadata { - /// Creates a new [UserMetadata] with the given event object. - pub fn new(event_object: Arc) -> Self { - Self { event_object } - } +/// Lifecycle action that causes the framework to invoke [`Procedure::event`]. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum EventTrigger { + /// The root procedure was submitted to the manager. + Submitted, + /// The root procedure was recovered from persisted state. + Recovered, + /// A child submission was attempted. + ChildSubmitted { + /// The submitted child procedure. + procedure_id: ProcedureId, + /// The result of the submission attempt. + outcome: ChildSubmissionOutcome, + }, + /// Procedure execution is being retried. + Retrying { + /// Phase in which the retry occurs. + phase: RetryPhase, + /// Retry attempt within the current runner lifecycle. + attempt: u32, + }, + /// Procedure rollback is starting. + RollingBack, + /// The procedure reached a successful terminal state. + Succeeded, + /// The procedure reached a failed terminal state. + Failed, + /// The procedure was poisoned and cannot proceed. + Poisoned, +} - /// Returns the event of the procedure. It can be None if the procedure does not emit any event. - pub fn to_event(&self) -> Option> { - self.event_object.to_event() +/// Phase of a procedure retry. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RetryPhase { + /// Retrying procedure execution. + Execute, + /// Retrying procedure rollback. + Rollback, +} + +/// Outcome of submitting a child procedure. +#[derive(Debug, Clone, PartialEq, Eq)] +pub enum ChildSubmissionOutcome { + Accepted, + AlreadyAccepted, + ManagerStopped, + SpawnFailed, +} + +impl Display for EventTrigger { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Submitted => write!(f, "Submitted"), + Self::Recovered => write!(f, "Recovered"), + Self::ChildSubmitted { + procedure_id, + outcome, + } => write!( + f, + "ChildSubmitted(procedure_id={procedure_id}, outcome={outcome})" + ), + Self::Retrying { phase, attempt } => write!(f, "Retrying({phase}, {attempt})"), + Self::RollingBack => write!(f, "RollingBack"), + Self::Succeeded => write!(f, "Succeeded"), + Self::Failed => write!(f, "Failed"), + Self::Poisoned => write!(f, "Poisoned"), + } + } +} + +impl Display for RetryPhase { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Execute => write!(f, "Execute"), + Self::Rollback => write!(f, "Rollback"), + } + } +} + +impl Display for ChildSubmissionOutcome { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::Accepted => write!(f, "Accepted"), + Self::AlreadyAccepted => write!(f, "AlreadyAccepted"), + Self::ManagerStopped => write!(f, "ManagerStopped"), + Self::SpawnFailed => write!(f, "SpawnFailed"), + } } } @@ -281,6 +371,10 @@ impl Procedure for Box { fn poison_keys(&self) -> PoisonKeys { (**self).poison_keys() } + + fn event(&self, ctx: &EventContext<'_>) -> Option> { + (**self).event(ctx) + } } #[derive(Clone, Debug, PartialEq, Eq, Hash, PartialOrd, Ord)] @@ -628,11 +722,46 @@ pub struct ProcedureInfo { #[cfg(test)] mod tests { + use async_trait::async_trait; use common_error::mock::MockError; use common_error::status_code::StatusCode; use super::*; + struct DefaultEventProcedure; + + #[async_trait] + impl Procedure for DefaultEventProcedure { + fn type_name(&self) -> &str { + "default_event" + } + + async fn execute(&mut self, _: &Context) -> Result { + Ok(Status::done()) + } + + fn dump(&self) -> Result { + Ok(String::new()) + } + + fn lock_key(&self) -> LockKey { + LockKey::default() + } + } + + #[test] + fn test_default_procedure_event_hook() { + let state = ProcedureState::Running; + let context = EventContext { + procedure_id: ProcedureId::random(), + lifecycle_state: &state, + trigger: EventTrigger::Succeeded, + }; + + assert!(DefaultEventProcedure.event(&context).is_none()); + assert!(Box::new(DefaultEventProcedure).event(&context).is_none()); + } + #[test] fn test_status() { let status = Status::executing(false); diff --git a/src/meta-srv/src/procedure/region_migration.rs b/src/meta-srv/src/procedure/region_migration.rs index f3c4860437..8dd6aac271 100644 --- a/src/meta-srv/src/procedure/region_migration.rs +++ b/src/meta-srv/src/procedure/region_migration.rs @@ -29,11 +29,10 @@ pub(crate) mod utils; use std::any::Any; use std::collections::{HashMap, HashSet}; use std::fmt::{Debug, Display}; -use std::sync::Arc; use std::time::Duration; use common_error::ext::BoxedError; -use common_event_recorder::{Event, Eventable}; +use common_event_recorder::Event; use common_meta::cache_invalidator::CacheInvalidatorRef; use common_meta::ddl::RegionFailureDetectorControllerRef; use common_meta::instruction::CacheIdent; @@ -50,7 +49,7 @@ use common_procedure::error::{ Error as ProcedureError, FromJsonSnafu, Result as ProcedureResult, ToJsonSnafu, }; use common_procedure::{ - Context as ProcedureContext, LockKey, Procedure, Status, StringKey, UserMetadata, + Context as ProcedureContext, EventContext, LockKey, Procedure, Status, StringKey, }; use common_telemetry::{debug, error, info}; use manager::RegionMigrationProcedureGuard; @@ -204,12 +203,6 @@ impl PersistentContext { } } -impl Eventable for PersistentContext { - fn to_event(&self) -> Option> { - Some(Box::new(RegionMigrationEvent::from_persistent_ctx(self))) - } -} - /// Metrics of region migration. #[derive(Debug, Clone, Default)] pub struct Metrics { @@ -971,8 +964,10 @@ impl Procedure for RegionMigrationProcedure { LockKey::new(self.context.persistent_ctx.lock_key()) } - fn user_metadata(&self) -> Option { - Some(UserMetadata::new(Arc::new(self.context.persistent_ctx()))) + fn event(&self, _ctx: &EventContext<'_>) -> Option> { + Some(Box::new(RegionMigrationEvent::from_persistent_ctx( + &self.context.persistent_ctx, + ))) } } @@ -1032,6 +1027,37 @@ mod tests { assert_eq!(expected, serialized); } + #[test] + fn test_event_hook_uses_current_persistent_context() { + let env = TestingEnv::new(); + let procedure = + RegionMigrationProcedure::new(new_persistent_context(), env.context_factory(), vec![]); + let state = common_procedure::ProcedureState::Running; + let triggers = [ + common_procedure::EventTrigger::Recovered, + common_procedure::EventTrigger::Retrying { + phase: common_procedure::RetryPhase::Execute, + attempt: 1, + }, + common_procedure::EventTrigger::RollingBack, + common_procedure::EventTrigger::Succeeded, + common_procedure::EventTrigger::Failed, + common_procedure::EventTrigger::Poisoned, + ]; + + for trigger in triggers { + let event = procedure + .event(&EventContext { + procedure_id: common_procedure::ProcedureId::random(), + lifecycle_state: &state, + trigger, + }) + .unwrap(); + assert_eq!(event.event_type(), "region_migration"); + assert_eq!(event.extra_rows().unwrap().len(), 1); + } + } + #[test] fn test_backward_compatibility() { let persistent_ctx = PersistentContext { diff --git a/src/meta-srv/src/procedure/repartition.rs b/src/meta-srv/src/procedure/repartition.rs index 2aa8f18799..bc3757319b 100644 --- a/src/meta-srv/src/procedure/repartition.rs +++ b/src/meta-srv/src/procedure/repartition.rs @@ -48,7 +48,7 @@ use common_meta::wal_provider::RegionWalOptions; use common_procedure::error::{FromJsonSnafu, ToJsonSnafu}; use common_procedure::{ BoxedProcedure, Context as ProcedureContext, Error as ProcedureError, LockKey, Procedure, - ProcedureManagerRef, Result as ProcedureResult, Status, StringKey, UserMetadata, + ProcedureManagerRef, Result as ProcedureResult, Status, StringKey, }; use common_telemetry::{error, info, warn}; use partition::expr::PartitionExpr; @@ -787,11 +787,6 @@ impl Procedure for RepartitionProcedure { fn lock_key(&self) -> LockKey { LockKey::new(self.context.persistent_ctx.lock_key()) } - - fn user_metadata(&self) -> Option { - // TODO(weny): support user metadata. - None - } } pub struct DefaultRepartitionProcedureFactory { diff --git a/src/meta-srv/src/procedure/repartition/group.rs b/src/meta-srv/src/procedure/repartition/group.rs index 2dc1117467..5a1d880d6b 100644 --- a/src/meta-srv/src/procedure/repartition/group.rs +++ b/src/meta-srv/src/procedure/repartition/group.rs @@ -39,7 +39,7 @@ use common_meta::rpc::router::RegionRoute; use common_procedure::error::{FromJsonSnafu, ToJsonSnafu}; use common_procedure::{ Context as ProcedureContext, Error as ProcedureError, LockKey, Procedure, - Result as ProcedureResult, Status, StringKey, UserMetadata, + Result as ProcedureResult, Status, StringKey, }; use common_telemetry::{error, info}; use serde::{Deserialize, Serialize}; @@ -260,11 +260,6 @@ impl Procedure for RepartitionGroupProcedure { fn lock_key(&self) -> LockKey { LockKey::new(self.context.persistent_ctx.lock_key()) } - - fn user_metadata(&self) -> Option { - // TODO(weny): support user metadata. - None - } } pub struct Context { diff --git a/tests-integration/tests/region_migration.rs b/tests-integration/tests/region_migration.rs index 4e26fab3bb..f43de91c24 100644 --- a/tests-integration/tests/region_migration.rs +++ b/tests-integration/tests/region_migration.rs @@ -25,6 +25,7 @@ use common_meta::key::{RegionDistribution, RegionRoleSet, TableMetadataManagerRe use common_meta::peer::Peer; use common_procedure::event::{ EVENTS_TABLE_PROCEDURE_ID_COLUMN_NAME, EVENTS_TABLE_PROCEDURE_STATE_COLUMN_NAME, + EVENTS_TABLE_PROCEDURE_TRIGGER_COLUMN_NAME, }; use common_query::Output; use common_recordbatch::RecordBatches; @@ -1288,6 +1289,7 @@ enum RegionMigrationEvents { ProcedureId, Timestamp, ProcedureState, + ProcedureTrigger, Schema, Table, EventType, @@ -1306,6 +1308,7 @@ impl Iden for RegionMigrationEvents { Self::ProcedureId => EVENTS_TABLE_PROCEDURE_ID_COLUMN_NAME, Self::Timestamp => EVENTS_TABLE_TIMESTAMP_COLUMN_NAME, Self::ProcedureState => EVENTS_TABLE_PROCEDURE_STATE_COLUMN_NAME, + Self::ProcedureTrigger => EVENTS_TABLE_PROCEDURE_TRIGGER_COLUMN_NAME, Self::Schema => DEFAULT_PRIVATE_SCHEMA_NAME, Self::Table => DEFAULT_EVENTS_TABLE_NAME, Self::EventType => EVENTS_TABLE_TYPE_COLUMN_NAME, @@ -1342,6 +1345,7 @@ async fn check_region_migration_events_system_table( let query = Query::select() .column(RegionMigrationEvents::RegionMigrationTriggerReason) .column(RegionMigrationEvents::ProcedureState) + .column(RegionMigrationEvents::ProcedureTrigger) .from((RegionMigrationEvents::Schema, RegionMigrationEvents::Table)) .and_where(Expr::col(RegionMigrationEvents::EventType).eq(REGION_MIGRATION_EVENT_TYPE)) .and_where(Expr::col(RegionMigrationEvents::ProcedureId).eq(procedure_id)) @@ -1357,11 +1361,11 @@ async fn check_region_migration_events_system_table( .remove(0); let expected = "\ -+---------------------------------+-----------------+ -| region_migration_trigger_reason | procedure_state | -+---------------------------------+-----------------+ -| Manual | Running | -| Manual | Done | -+---------------------------------+-----------------+"; ++---------------------------------+-----------------+-------------------+ +| region_migration_trigger_reason | procedure_state | procedure_trigger | ++---------------------------------+-----------------+-------------------+ +| Manual | Running | Submitted | +| Manual | Done | Succeeded | ++---------------------------------+-----------------+-------------------+"; check_output_stream(result.unwrap().data, expected).await; }