From 46edabe916a10c2cafe816b8fe01f0d5da8fffc8 Mon Sep 17 00:00:00 2001 From: discord9 Date: Thu, 20 Aug 2026 13:19:13 +0800 Subject: [PATCH] fix(event): preserve legacy DDL event context compatibility Signed-off-by: discord9 --- src/common/meta/src/rpc/ddl.rs | 47 ++++++++++++++++++++++++--- src/meta-srv/src/service/procedure.rs | 5 +++ 2 files changed, 48 insertions(+), 4 deletions(-) diff --git a/src/common/meta/src/rpc/ddl.rs b/src/common/meta/src/rpc/ddl.rs index 47bd537b9b..b1c5ca8f21 100644 --- a/src/common/meta/src/rpc/ddl.rs +++ b/src/common/meta/src/rpc/ddl.rs @@ -42,10 +42,10 @@ use api::v1::{ }; use base64::Engine as _; use base64::engine::general_purpose; -use common_base::protocol::Channel; use common_catalog::{format_full_flow_name, format_full_table_name}; use common_error::ext::BoxedError; -pub use common_event_recorder::TriggerReason; +pub use common_event_recorder::{PersistentEventContext, TriggerReason}; +use common_session::channel_protocol; use common_time::{DatabaseTimeToLive, Timestamp}; use prost::Message; use serde::{Deserialize, Serialize}; @@ -68,6 +68,8 @@ use crate::key::table_name::{TableNameKey, TableNameManager}; /// Reserved query-context extension key for the frontend peer address that submitted a DDL request. pub const ORIGIN_FRONTEND_ADDR_EXTENSION_KEY: &str = "__greptime_origin_frontend.addr"; +/// Reserved query-context extension key for the trigger reason supplied by a legacy frontend. +pub const TRIGGER_REASON_EXTENSION_KEY: &str = "__greptime_event.trigger_reason"; /// Reserved query-context extension key for the authenticated database creator. pub const CREATE_DATABASE_CREATOR_EXTENSION_KEY: &str = "__greptime_create_database.creator"; /// Internal gRPC metadata key for the authenticated database creator. @@ -1651,6 +1653,21 @@ pub struct QueryContext { pub sst_min_sequences: HashMap, } +/// Derives the persistent event context used by legacy DDL submitters that do +/// not populate the newer protobuf `event_context` field. +pub fn event_context_from_query_context(query_context: &QueryContext) -> PersistentEventContext { + let reason = query_context + .extensions + .get(TRIGGER_REASON_EXTENSION_KEY) + .map(|value| TriggerReason::from_extension(value)) + .unwrap_or_default(); + let mut context = PersistentEventContext::new(reason); + if let Some(protocol) = channel_protocol(query_context.channel) { + context = context.with_protocol(protocol); + } + context +} + impl QueryContext { /// Get the current catalog pub fn current_catalog(&self) -> &str { @@ -1679,8 +1696,7 @@ impl QueryContext { /// Returns the protocol derived from the typed query channel. pub fn protocol(&self) -> Option { - let channel = Channel::from(u32::from(self.channel)); - (channel != Channel::Unknown).then(|| channel.as_ref().to_string()) + channel_protocol(self.channel).map(str::to_string) } pub fn snapshot_seqs(&self) -> &HashMap { @@ -1856,6 +1872,29 @@ mod tests { use super::{AlterTableTask, CreateTableTask, *}; + #[test] + fn test_legacy_event_context_fallback_uses_reason_and_typed_channel() { + let mut query_context = QueryContext::default(); + query_context.channel = 4; + query_context.extensions.insert( + TRIGGER_REASON_EXTENSION_KEY.to_string(), + "auto_create".to_string(), + ); + let request = api::v1::meta::DdlTaskRequest { + query_context: Some(query_context.clone().into()), + event_context: None, + ..Default::default() + }; + + // This is the wire shape sent by an older frontend: the new field is + // absent, while the legacy query context still carries both values. + assert!(request.event_context.is_none()); + let query_context = QueryContext::from(request.query_context.unwrap()); + let context = event_context_from_query_context(&query_context); + assert_eq!(context.reason, TriggerReason::AutoCreate); + assert_eq!(context.protocol.as_deref(), Some("prometheus")); + } + #[test] fn test_ddl_timeout_secs() { assert_eq!(ddl_timeout_secs(Duration::ZERO), 0); diff --git a/src/meta-srv/src/service/procedure.rs b/src/meta-srv/src/service/procedure.rs index df27e2cb32..759a546c53 100644 --- a/src/meta-srv/src/service/procedure.rs +++ b/src/meta-srv/src/service/procedure.rs @@ -30,6 +30,7 @@ use common_meta::procedure_executor::ExecutorContext; use common_meta::rpc::ddl::{ CREATE_DATABASE_CREATOR_EXTENSION_KEY, CREATE_DATABASE_CREATOR_METADATA_KEY, CreatorGrantIntent, DdlTask, QueryContext, SubmitDdlTaskRequest, + event_context_from_query_context, }; use common_meta::rpc::procedure::{ self, GcRegionsRequest as MetaGcRegionsRequest, GcResponse, @@ -131,6 +132,10 @@ impl procedure_service_server::ProcedureService for Metasrv { param: "query_context", })? .into(); + // Older frontends omit `event_context` and carry its reason in the + // query-context extensions. Preserve that context at the DDL boundary. + let event_context = + event_context.or_else(|| Some(event_context_from_query_context(&query_context))); let mut task: DdlTask = task .context(error::MissingRequiredParameterSnafu { param: "task" })? .try_into()