feat(procedure): support trigger-aware procedure events (#8549)

* feat(procedure): add live event hook

Signed-off-by: WenyXu <wenymedia@gmail.com>

* feat(procedure): add event trigger metadata

Signed-off-by: WenyXu <wenymedia@gmail.com>

* fix(procedure): distinguish recovery events

Signed-off-by: WenyXu <wenymedia@gmail.com>

* chore: apply suggestions

Signed-off-by: WenyXu <wenymedia@gmail.com>

* chore: apply suggestions

Signed-off-by: WenyXu <wenymedia@gmail.com>

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-07-24 17:44:33 +08:00
committed by GitHub
parent 8afb960dda
commit d39f2292fa
9 changed files with 510 additions and 130 deletions
+63 -17
View File
@@ -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<dyn Event>,
}
@@ -44,12 +47,14 @@ impl ProcedureEvent {
procedure_id: ProcedureId,
internal_event: Box<dyn Event>,
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<Vec<Row>> {
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");
}
}
+4 -4
View File
@@ -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;
+149 -36
View File
@@ -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<EventRecorderRef>,
/// The user metadata of the procedure. It's generated by [Procedure::user_metadata].
user_metadata: Option<UserMetadata>,
}
impl ProcedureMeta {
@@ -98,8 +108,6 @@ impl ProcedureMeta {
lock_key: LockKey,
poison_keys: PoisonKeys,
type_name: &str,
event_recorder: Option<EventRecorderRef>,
user_metadata: Option<UserMetadata>,
) -> 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<Watcher> {
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<Vec<Box<dyn Event>>>,
}
impl CapturingEventRecorder {
fn triggers(&self) -> Vec<EventTrigger> {
self.events
.lock()
.unwrap()
.iter()
.map(|event| {
event
.as_any()
.downcast_ref::<ProcedureEvent>()
.unwrap()
.trigger
.clone()
})
.collect()
}
}
impl EventRecorder for CapturingEventRecorder {
fn record(&self, event: Box<dyn Event>) {
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<Box<dyn Event>> {
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<oneshot::Sender<()>>,
+101 -29
View File
@@ -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<ProcedureStore>,
pub(crate) rolling_back: bool,
pub(crate) event_recorder: Option<EventRecorderRef>,
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>>(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<Error>) {
@@ -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<Output>) {
fn done(&mut self, output: Option<Output>) {
// 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,
}
}
+144 -15
View File
@@ -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<UserMetadata> {
/// 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<Box<dyn Event>> {
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<dyn Eventable>,
/// 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<dyn Eventable>) -> 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<Box<dyn Event>> {
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<T: Procedure + ?Sized> Procedure for Box<T> {
fn poison_keys(&self) -> PoisonKeys {
(**self).poison_keys()
}
fn event(&self, ctx: &EventContext<'_>) -> Option<Box<dyn Event>> {
(**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<Status> {
Ok(Status::done())
}
fn dump(&self) -> Result<String> {
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);
+37 -11
View File
@@ -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<Box<dyn Event>> {
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<UserMetadata> {
Some(UserMetadata::new(Arc::new(self.context.persistent_ctx())))
fn event(&self, _ctx: &EventContext<'_>) -> Option<Box<dyn Event>> {
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 {
+1 -6
View File
@@ -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<UserMetadata> {
// TODO(weny): support user metadata.
None
}
}
pub struct DefaultRepartitionProcedureFactory {
@@ -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<UserMetadata> {
// TODO(weny): support user metadata.
None
}
}
pub struct Context {
+10 -6
View File
@@ -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;
}