refactor(event): separate procedure submission context (#8856)

* refactor(event): separate procedure submission context

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

* fix(event): map extensions and forward GC context

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

* fix(gc): initialize integration test context

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

* fix(test): pass procedure context to DDL helpers

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

* refactor: simplify procedure submission contexts

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

* refactor(event): separate procedure and query contexts

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

* fix(event): clarify procedure context propagation

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

* fix(event): preserve procedure submission context

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

* refactor(event): move DDL context by value

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

* fix(test): retain manual GC event context

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

* refactor(event): tighten procedure context API

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

* chore: update greptime-proto

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

---------

Signed-off-by: WenyXu <wenymedia@gmail.com>
This commit is contained in:
Weny Xu
2026-08-13 07:29:00 +00:00
committed by GitHub
parent d0fecdd6b0
commit 1af4c33524
41 changed files with 1487 additions and 770 deletions
+35 -17
View File
@@ -38,17 +38,20 @@ const DEFAULT_FULL_FILE_LISTING: bool = false;
)]
pub(crate) async fn gc_regions(
procedure_service_handler: &ProcedureServiceHandlerRef,
_ctx: &QueryContextRef,
query_ctx: &QueryContextRef,
params: &[ValueRef<'_>],
) -> Result<Value> {
let (region_ids, full_file_listing) = parse_gc_regions_params(params)?;
let resp = procedure_service_handler
.gc_regions(GcRegionsRequest {
region_ids,
full_file_listing,
timeout: None,
})
.gc_regions(
query_ctx.clone(),
GcRegionsRequest {
region_ids,
full_file_listing,
timeout: None,
},
)
.await?;
Ok(Value::from(resp.processed_regions))
@@ -69,13 +72,16 @@ pub(crate) async fn gc_table(
parse_gc_table_params(params, query_ctx)?;
let resp = procedure_service_handler
.gc_table(GcTableRequest {
catalog_name,
schema_name,
table_name,
full_file_listing,
timeout: None,
})
.gc_table(
query_ctx.clone(),
GcTableRequest {
catalog_name,
schema_name,
table_name,
full_file_listing,
timeout: None,
},
)
.await?;
Ok(Value::from(resp.processed_regions))
@@ -268,13 +274,17 @@ mod tests {
impl ProcedureServiceHandler for MockProcedureServiceHandler {
async fn purge_table(
&self,
_table_name: table::table_name::TableName,
_query_ctx: QueryContextRef,
_table_name: table::table_name::TableName,
) -> Result<()> {
unreachable!()
}
async fn migrate_region(&self, _request: MigrateRegionRequest) -> Result<Option<String>> {
async fn migrate_region(
&self,
_query_ctx: QueryContextRef,
_request: MigrateRegionRequest,
) -> Result<Option<String>> {
unreachable!()
}
@@ -297,12 +307,20 @@ mod tests {
unreachable!()
}
async fn gc_regions(&self, request: GcRegionsRequest) -> Result<GcResponse> {
async fn gc_regions(
&self,
_query_ctx: QueryContextRef,
request: GcRegionsRequest,
) -> Result<GcResponse> {
*self.gc_regions_request.lock().unwrap() = Some(request);
Ok(GcResponse::default())
}
async fn gc_table(&self, request: GcTableRequest) -> Result<GcResponse> {
async fn gc_table(
&self,
_query_ctx: QueryContextRef,
request: GcTableRequest,
) -> Result<GcResponse> {
*self.gc_table_request.lock().unwrap() = Some(request);
Ok(GcResponse::default())
}
@@ -47,7 +47,7 @@ const DEFAULT_TIMEOUT_SECS: u64 = 300;
)]
pub(crate) async fn migrate_region(
procedure_service_handler: &ProcedureServiceHandlerRef,
_ctx: &QueryContextRef,
query_ctx: &QueryContextRef,
params: &[ValueRef<'_>],
) -> Result<Value> {
let (region_id, from_peer, to_peer, timeout) = match params.len() {
@@ -82,12 +82,15 @@ pub(crate) async fn migrate_region(
match (region_id, from_peer, to_peer, timeout) {
(Some(region_id), Some(from_peer), Some(to_peer), Some(timeout)) => {
let pid = procedure_service_handler
.migrate_region(MigrateRegionRequest {
region_id,
from_peer,
to_peer,
timeout: Duration::from_secs(timeout),
})
.migrate_region(
query_ctx.clone(),
MigrateRegionRequest {
region_id,
from_peer,
to_peer,
timeout: Duration::from_secs(timeout),
},
)
.await?;
match pid {
+9 -5
View File
@@ -64,8 +64,8 @@ pub(crate) async fn purge_table(
procedure_service_handler
.purge_table(
TableName::new(catalog_name, schema_name, table_name),
query_ctx.clone(),
TableName::new(catalog_name, schema_name, table_name),
)
.await?;
Ok(Value::from(0_u64))
@@ -131,8 +131,8 @@ mod tests {
impl ProcedureServiceHandler for RecordingHandler {
async fn purge_table(
&self,
table_name: TableName,
query_ctx: QueryContextRef,
table_name: TableName,
) -> Result<()> {
if self.fail {
return InvalidFuncArgsSnafu {
@@ -144,7 +144,11 @@ mod tests {
Ok(())
}
async fn migrate_region(&self, _: MigrateRegionRequest) -> Result<Option<String>> {
async fn migrate_region(
&self,
_: QueryContextRef,
_: MigrateRegionRequest,
) -> Result<Option<String>> {
unreachable!()
}
async fn reconcile(&self, _: ReconcileRequest) -> Result<Option<String>> {
@@ -159,10 +163,10 @@ mod tests {
fn catalog_manager(&self) -> &CatalogManagerRef {
unreachable!()
}
async fn gc_regions(&self, _: GcRegionsRequest) -> Result<GcResponse> {
async fn gc_regions(&self, _: QueryContextRef, _: GcRegionsRequest) -> Result<GcResponse> {
unreachable!()
}
async fn gc_table(&self, _: GcTableRequest) -> Result<GcResponse> {
async fn gc_table(&self, _: QueryContextRef, _: GcTableRequest) -> Result<GcResponse> {
unreachable!()
}
}
+16 -4
View File
@@ -89,10 +89,14 @@ pub trait TableMutationHandler: Send + Sync {
#[async_trait]
pub trait ProcedureServiceHandler: Send + Sync {
/// Permanently purge a dropped table.
async fn purge_table(&self, table_name: TableName, query_ctx: QueryContextRef) -> Result<()>;
async fn purge_table(&self, query_ctx: QueryContextRef, table_name: TableName) -> Result<()>;
/// Migrate a region from source peer to target peer, returns the procedure id if success.
async fn migrate_region(&self, request: MigrateRegionRequest) -> Result<Option<String>>;
async fn migrate_region(
&self,
query_ctx: QueryContextRef,
request: MigrateRegionRequest,
) -> Result<Option<String>>;
/// Reconcile a table, database or catalog, returns the procedure id if success.
async fn reconcile(&self, request: ReconcileRequest) -> Result<Option<String>>;
@@ -107,10 +111,18 @@ pub trait ProcedureServiceHandler: Send + Sync {
fn catalog_manager(&self) -> &CatalogManagerRef;
/// Manually trigger GC for specific regions.
async fn gc_regions(&self, request: MetaGcRegionsRequest) -> Result<MetaGcResponse>;
async fn gc_regions(
&self,
query_ctx: QueryContextRef,
request: MetaGcRegionsRequest,
) -> Result<MetaGcResponse>;
/// Manually trigger GC for a table.
async fn gc_table(&self, request: MetaGcTableRequest) -> Result<MetaGcResponse>;
async fn gc_table(
&self,
query_ctx: QueryContextRef,
request: MetaGcTableRequest,
) -> Result<MetaGcResponse>;
}
/// This flow service handler is only use for flush flow for now.
+12 -3
View File
@@ -60,14 +60,15 @@ impl FunctionState {
impl ProcedureServiceHandler for MockProcedureServiceHandler {
async fn purge_table(
&self,
_table_name: table::table_name::TableName,
_query_ctx: QueryContextRef,
_table_name: table::table_name::TableName,
) -> Result<()> {
Ok(())
}
async fn migrate_region(
&self,
_ctx: QueryContextRef,
_request: MigrateRegionRequest,
) -> Result<Option<String>> {
Ok(Some("test_pid".to_string()))
@@ -92,7 +93,11 @@ impl FunctionState {
Ok(())
}
async fn gc_regions(&self, _request: GcRegionsRequest) -> Result<GcResponse> {
async fn gc_regions(
&self,
_context: QueryContextRef,
_request: GcRegionsRequest,
) -> Result<GcResponse> {
Ok(GcResponse {
processed_regions: 1,
need_retry_regions: vec![],
@@ -101,7 +106,11 @@ impl FunctionState {
})
}
async fn gc_table(&self, _request: GcTableRequest) -> Result<GcResponse> {
async fn gc_table(
&self,
_context: QueryContextRef,
_request: GcTableRequest,
) -> Result<GcResponse> {
Ok(GcResponse {
processed_regions: 1,
need_retry_regions: vec![],