diff --git a/apps/aether-gateway/src/task_runtime/mod.rs b/apps/aether-gateway/src/task_runtime/mod.rs index de9c32f24..33b0f7eca 100644 --- a/apps/aether-gateway/src/task_runtime/mod.rs +++ b/apps/aether-gateway/src/task_runtime/mod.rs @@ -8,7 +8,7 @@ use aether_data_contracts::repository::background_tasks::{ use aether_runtime::task::spawn_named; use aether_task_runtime::{RetryPolicy, TaskDefinition, TaskKind}; pub(crate) use aether_task_runtime::{TaskSupervisor, TaskSupervisorMetrics}; -use serde_json::Value; +use serde_json::{json, Value}; use sha2::{Digest, Sha256}; use tokio::task::JoinHandle; use tracing::warn; @@ -73,6 +73,41 @@ fn build_worker_boot_run_id(task_key: &str) -> String { format!("{}~{suffix}", &full_run_id[..prefix_bytes]) } +fn build_worker_boot_run( + task_key: &str, + kind: BackgroundTaskKind, + trigger: &str, + now: u64, +) -> UpsertBackgroundTaskRun { + UpsertBackgroundTaskRun { + id: build_worker_boot_run_id(task_key), + task_key: task_key.to_string(), + kind, + trigger: trigger.to_string(), + status: BackgroundTaskStatus::Running, + attempt: 1, + max_attempts: 1, + // This row represents the logical worker registration shared by every gateway. + // A supervisor starting does not mean that instance owns the singleton lease. + owner_instance: None, + progress_percent: 0, + progress_message: Some("worker registered".to_string()), + payload_json: None, + result_json: None, + error_message: None, + cancel_requested: false, + created_by: Some("system".to_string()), + created_at_unix_secs: now, + started_at_unix_secs: Some(now), + finished_at_unix_secs: None, + updated_at_unix_secs: now, + } +} + +fn worker_boot_event_payload(gateway_instance_id: &str) -> Value { + json!({ "gateway_instance_id": gateway_instance_id }) +} + pub(crate) fn spawn_singleton_worker( app: AppState, task_key: &'static str, @@ -510,35 +545,18 @@ pub(crate) fn spawn_record_worker_boot( ) -> JoinHandle<()> { spawn_named("task-runtime-record-worker-boot", async move { let now = now_unix_secs(); - let run_id = build_worker_boot_run_id(task_key); - let run = UpsertBackgroundTaskRun { - id: run_id.clone(), - task_key: task_key.to_string(), - kind, - trigger: trigger.to_string(), - status: BackgroundTaskStatus::Running, - attempt: 1, - max_attempts: 1, - owner_instance: Some(app.tunnel.local_instance_id().to_string()), - progress_percent: 0, - progress_message: Some("worker booted".to_string()), - payload_json: None, - result_json: None, - error_message: None, - cancel_requested: false, - created_by: Some("system".to_string()), - created_at_unix_secs: now, - started_at_unix_secs: Some(now), - finished_at_unix_secs: None, - updated_at_unix_secs: now, - }; - let _ = upsert_run_with_logging(&app, run).await; + let gateway_instance_id = app.tunnel.local_instance_id().to_string(); + let run = build_worker_boot_run(task_key, kind, trigger, now); + let run_id = run.id.clone(); + if upsert_run_with_logging(&app, run).await.is_none() { + return; + } append_event_with_logging( &app, &run_id, "worker_boot", - "background worker started", - None, + "background worker supervisor started", + Some(worker_boot_event_payload(&gateway_instance_id)), ) .await; }) @@ -784,7 +802,21 @@ pub(crate) async fn submit_provider_delete_task( #[cfg(test)] mod worker_boot_run_id_tests { - use super::{build_worker_boot_run_id, BACKGROUND_TASK_RUN_ID_MAX_BYTES}; + use std::collections::BTreeSet; + use std::sync::Arc; + + use super::{ + build_worker_boot_run, build_worker_boot_run_id, spawn_record_worker_boot, + worker_boot_event_payload, BACKGROUND_TASK_RUN_ID_MAX_BYTES, + }; + use aether_data::repository::background_tasks::InMemoryBackgroundTaskRepository; + use aether_data_contracts::repository::background_tasks::{ + BackgroundTaskKind, BackgroundTaskListQuery, BackgroundTaskReadRepository, + BackgroundTaskStatus, + }; + use serde_json::json; + + use crate::{data::GatewayDataState, AppState}; #[test] fn worker_boot_run_id_is_keyed_only_by_task() { @@ -832,4 +864,99 @@ mod worker_boot_run_id_tests { assert!(run_id.is_char_boundary(run_id.len())); assert!(run_id.starts_with("boot:后台任务")); } + + #[test] + fn worker_boot_run_is_a_gateway_neutral_logical_registration() { + let run = build_worker_boot_run( + "usage.queue.worker", + BackgroundTaskKind::Daemon, + "daemon", + 123, + ); + + assert_eq!(run.id, "boot:usage.queue.worker"); + assert_eq!(run.status, BackgroundTaskStatus::Running); + assert_eq!(run.owner_instance, None); + assert_eq!(run.progress_message.as_deref(), Some("worker registered")); + assert_eq!(run.created_at_unix_secs, 123); + assert_eq!(run.started_at_unix_secs, Some(123)); + assert_eq!(run.updated_at_unix_secs, 123); + } + + #[test] + fn worker_boot_event_keeps_the_observing_gateway_instance() { + assert_eq!( + worker_boot_event_payload("gateway-a"), + json!({ "gateway_instance_id": "gateway-a" }) + ); + } + + #[tokio::test] + async fn worker_boot_registration_is_shared_but_events_keep_each_gateway() { + let repository = Arc::new(InMemoryBackgroundTaskRepository::default()); + let state_for = |gateway_instance_id: &str| { + AppState::new() + .expect("gateway state should build") + .with_data_state_for_tests( + GatewayDataState::disabled() + .with_background_task_repository_for_tests(repository.clone()), + ) + .with_tunnel_identity_for_tests(gateway_instance_id, None) + }; + + spawn_record_worker_boot( + state_for("gateway-a"), + "usage.queue.worker", + BackgroundTaskKind::Daemon, + "daemon", + ) + .await + .expect("gateway-a worker boot recorder should finish"); + spawn_record_worker_boot( + state_for("gateway-b"), + "usage.queue.worker", + BackgroundTaskKind::Daemon, + "daemon", + ) + .await + .expect("gateway-b worker boot recorder should finish"); + + let page = repository + .list_runs(&BackgroundTaskListQuery { + task_key_substring: Some("usage.queue.worker".to_string()), + kind: None, + status: None, + trigger: None, + offset: 0, + limit: 10, + }) + .await + .expect("worker boot runs should load"); + assert_eq!(page.total, 1); + assert_eq!(page.items.len(), 1); + assert_eq!(page.items[0].id, "boot:usage.queue.worker"); + assert_eq!(page.items[0].owner_instance, None); + + let events = repository + .list_events("boot:usage.queue.worker", 0, 10) + .await + .expect("worker boot events should load"); + assert_eq!(events.len(), 2); + assert!(events.iter().all(|event| event.event_type == "worker_boot")); + let gateway_instances = events + .iter() + .filter_map(|event| { + event + .payload_json + .as_ref() + .and_then(|payload| payload.get("gateway_instance_id")) + .and_then(serde_json::Value::as_str) + .map(str::to_string) + }) + .collect::>(); + assert_eq!( + gateway_instances, + BTreeSet::from(["gateway-a".to_string(), "gateway-b".to_string()]) + ); + } } diff --git a/crates/aether-data/adapters/mysql/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql b/crates/aether-data/adapters/mysql/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql new file mode 100644 index 000000000..f7d07d142 --- /dev/null +++ b/crates/aether-data/adapters/mysql/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql @@ -0,0 +1,22 @@ +-- Worker supervisors are registered as one logical row per task. Older binaries +-- included the ephemeral gateway instance in the row id, leaving a permanently +-- running row after every restart. Remove only those system-generated boot rows; +-- current workers recreate the stable logical rows after migrations complete. +-- The metadata predicate also replaces task-only rows written by early builds of +-- this fix that still claimed an instance owner. Delete children explicitly so +-- cleanup remains complete after imports performed with FK checks disabled. +DELETE FROM background_task_events +WHERE run_id IN ( + SELECT id + FROM background_task_runs + WHERE id LIKE 'boot:%' + AND owner_instance IS NOT NULL + AND created_by = 'system' + AND progress_message = 'worker booted' +); + +DELETE FROM background_task_runs +WHERE id LIKE 'boot:%' + AND owner_instance IS NOT NULL + AND created_by = 'system' + AND progress_message = 'worker booted'; diff --git a/crates/aether-data/adapters/postgres/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql b/crates/aether-data/adapters/postgres/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql new file mode 100644 index 000000000..f7d07d142 --- /dev/null +++ b/crates/aether-data/adapters/postgres/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql @@ -0,0 +1,22 @@ +-- Worker supervisors are registered as one logical row per task. Older binaries +-- included the ephemeral gateway instance in the row id, leaving a permanently +-- running row after every restart. Remove only those system-generated boot rows; +-- current workers recreate the stable logical rows after migrations complete. +-- The metadata predicate also replaces task-only rows written by early builds of +-- this fix that still claimed an instance owner. Delete children explicitly so +-- cleanup remains complete after imports performed with FK checks disabled. +DELETE FROM background_task_events +WHERE run_id IN ( + SELECT id + FROM background_task_runs + WHERE id LIKE 'boot:%' + AND owner_instance IS NOT NULL + AND created_by = 'system' + AND progress_message = 'worker booted' +); + +DELETE FROM background_task_runs +WHERE id LIKE 'boot:%' + AND owner_instance IS NOT NULL + AND created_by = 'system' + AND progress_message = 'worker booted'; diff --git a/crates/aether-data/adapters/sqlite/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql b/crates/aether-data/adapters/sqlite/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql new file mode 100644 index 000000000..f7d07d142 --- /dev/null +++ b/crates/aether-data/adapters/sqlite/migrations/20260731000000_cleanup_duplicate_worker_boot_runs.sql @@ -0,0 +1,22 @@ +-- Worker supervisors are registered as one logical row per task. Older binaries +-- included the ephemeral gateway instance in the row id, leaving a permanently +-- running row after every restart. Remove only those system-generated boot rows; +-- current workers recreate the stable logical rows after migrations complete. +-- The metadata predicate also replaces task-only rows written by early builds of +-- this fix that still claimed an instance owner. Delete children explicitly so +-- cleanup remains complete after imports performed with FK checks disabled. +DELETE FROM background_task_events +WHERE run_id IN ( + SELECT id + FROM background_task_runs + WHERE id LIKE 'boot:%' + AND owner_instance IS NOT NULL + AND created_by = 'system' + AND progress_message = 'worker booted' +); + +DELETE FROM background_task_runs +WHERE id LIKE 'boot:%' + AND owner_instance IS NOT NULL + AND created_by = 'system' + AND progress_message = 'worker booted'; diff --git a/crates/aether-data/adapters/sqlite/src/migrations.rs b/crates/aether-data/adapters/sqlite/src/migrations.rs index 76acd9d16..aedab37e9 100644 --- a/crates/aether-data/adapters/sqlite/src/migrations.rs +++ b/crates/aether-data/adapters/sqlite/src/migrations.rs @@ -538,6 +538,146 @@ WHERE request_id = ? assert_eq!(remaining, 0); } + #[tokio::test] + async fn worker_boot_cleanup_removes_only_legacy_instance_rows_and_events() { + const CLEANUP_VERSION: i64 = 20260731000000; + const HASHED_LEGACY_ID: &str = + "boot:maintenance.request.candidate.cleanup:~0123456789abcdef0123"; + const OVERLONG_LEGACY_ID: &str = + "boot:maintenance.proxy.node.metrics.cleanup:gateway-instance-with-an-overlong-id"; + + assert_eq!(HASHED_LEGACY_ID.len(), 64); + assert!(OVERLONG_LEGACY_ID.len() > 64); + + let pool = sqlx::sqlite::SqlitePoolOptions::new() + .max_connections(1) + .connect("sqlite::memory:") + .await + .expect("in-memory sqlite pool"); + + for migration in MIGRATOR + .iter() + .filter(|migration| migration.version < CLEANUP_VERSION) + { + sqlx::raw_sql(migration.sql.as_ref()) + .execute(&pool) + .await + .unwrap_or_else(|err| panic!("migration {} should run: {err}", migration.version)); + } + + sqlx::query( + r#" +INSERT INTO background_task_runs ( + id, task_key, kind, "trigger", status, owner_instance, + progress_message, created_by, created_at_unix_secs, updated_at_unix_secs +) VALUES + ('boot:usage.queue.worker:gateway-a', 'usage.queue.worker', 'daemon', 'daemon', + 'running', 'gateway-a', 'worker booted', 'system', 1, 1), + ('boot:usage.queue.worker:gateway-b', 'usage.queue.worker', 'daemon', 'daemon', + 'running', 'gateway-b', 'worker booted', 'system', 2, 2), + ('boot:maintenance.request.candidate.cleanup:~0123456789abcdef0123', + 'maintenance.request.candidate.cleanup', 'scheduled', 'interval', + 'running', 'gateway-hash', 'worker booted', 'system', 3, 3), + ('boot:maintenance.proxy.node.metrics.cleanup:gateway-instance-with-an-overlong-id', + 'maintenance.proxy.node.metrics.cleanup', 'scheduled', 'interval', + 'running', 'gateway-overlong', 'worker booted', 'system', 4, 4), + ('boot:model.fetch.worker', 'model.fetch.worker', 'scheduled', 'interval', + 'running', 'gateway-early-fix', 'worker booted', 'system', 5, 5), + ('boot:usage.queue.worker', 'usage.queue.worker', 'daemon', 'daemon', + 'running', NULL, 'worker registered', 'system', 6, 6), + ('boot:ownerless-worker', 'ownerless.worker', 'daemon', 'daemon', + 'running', NULL, 'worker booted', 'system', 7, 7), + ('boot:custom-progress', 'custom.progress', 'daemon', 'daemon', + 'running', 'gateway-custom', 'worker healthy', 'system', 8, 8), + ('boot:user-request', 'user.request', 'on_demand', 'manual', + 'running', 'gateway-user', 'worker booted', 'admin', 9, 9); + +INSERT INTO background_task_events ( + id, run_id, event_type, message, created_at_unix_secs +) VALUES + ('legacy-event-a', 'boot:usage.queue.worker:gateway-a', 'worker_boot', 'legacy', 1), + ('legacy-event-b', 'boot:usage.queue.worker:gateway-b', 'worker_boot', 'legacy', 2), + ('legacy-event-hash', 'boot:maintenance.request.candidate.cleanup:~0123456789abcdef0123', + 'worker_boot', 'legacy hash', 3), + ('legacy-event-overlong', + 'boot:maintenance.proxy.node.metrics.cleanup:gateway-instance-with-an-overlong-id', + 'worker_boot', 'legacy overlong', 4), + ('early-fix-event', 'boot:model.fetch.worker', 'worker_boot', 'early fix', 5), + ('logical-event', 'boot:usage.queue.worker', 'worker_boot', 'logical', 6), + ('ownerless-event', 'boot:ownerless-worker', 'worker_boot', 'ownerless', 7), + ('custom-progress-event', 'boot:custom-progress', 'worker_boot', 'custom progress', 8), + ('manual-event', 'boot:user-request', 'manual', 'manual', 9); +"#, + ) + .execute(&pool) + .await + .expect("worker boot cleanup fixtures should insert"); + + let migration = MIGRATOR + .iter() + .find(|migration| migration.version == CLEANUP_VERSION) + .expect("worker boot cleanup migration should be embedded"); + sqlx::raw_sql(migration.sql.as_ref()) + .execute(&pool) + .await + .expect("worker boot cleanup migration should run"); + sqlx::raw_sql(migration.sql.as_ref()) + .execute(&pool) + .await + .expect("worker boot cleanup migration should be idempotent"); + + let remaining_runs = sqlx::query_as::<_, (String, Option, Option)>( + r#" +SELECT id, owner_instance, created_by +FROM background_task_runs +ORDER BY id +"#, + ) + .fetch_all(&pool) + .await + .expect("remaining worker task rows should load"); + assert_eq!( + remaining_runs, + vec![ + ( + "boot:custom-progress".to_string(), + Some("gateway-custom".to_string()), + Some("system".to_string()), + ), + ( + "boot:ownerless-worker".to_string(), + None, + Some("system".to_string()), + ), + ( + "boot:usage.queue.worker".to_string(), + None, + Some("system".to_string()), + ), + ( + "boot:user-request".to_string(), + Some("gateway-user".to_string()), + Some("admin".to_string()), + ), + ] + ); + + let remaining_events = + sqlx::query_scalar::<_, String>("SELECT id FROM background_task_events ORDER BY id") + .fetch_all(&pool) + .await + .expect("remaining worker task events should load"); + assert_eq!( + remaining_events, + vec![ + "custom-progress-event".to_string(), + "logical-event".to_string(), + "manual-event".to_string(), + "ownerless-event".to_string(), + ] + ); + } + #[tokio::test] async fn pending_and_startup_preparation_reject_dirty_migration_state() { let pool = sqlx::sqlite::SqlitePoolOptions::new() diff --git a/crates/aether-data/runtime/src/lifecycle/bootstrap/postgres.rs b/crates/aether-data/runtime/src/lifecycle/bootstrap/postgres.rs index 2fa7cfa33..b6aa14a0a 100644 --- a/crates/aether-data/runtime/src/lifecycle/bootstrap/postgres.rs +++ b/crates/aether-data/runtime/src/lifecycle/bootstrap/postgres.rs @@ -7,7 +7,7 @@ use tracing::info; // Generated by build.rs from schema/bootstrap/postgres. pub(crate) static EMPTY_DATABASE_SNAPSHOT_SQL: &str = include_str!(concat!(env!("OUT_DIR"), "/empty_database_snapshot.sql")); -pub(crate) const EMPTY_DATABASE_SNAPSHOT_CUTOFF_VERSION: i64 = 20260727000000; +pub(crate) const EMPTY_DATABASE_SNAPSHOT_CUTOFF_VERSION: i64 = 20260731000000; const PUBLIC_BASE_TABLE_COUNT_SQL: &str = r#" SELECT COUNT(*)::BIGINT diff --git a/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs b/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs index 83a98e289..9b753bb81 100644 --- a/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs +++ b/crates/aether-data/runtime/src/lifecycle/migrate/tests.rs @@ -410,6 +410,7 @@ fn empty_database_snapshot_covers_current_cutoff_versions() { 20260718010000, 20260720000000, 20260727000000, + 20260731000000, ] ); } @@ -983,6 +984,43 @@ fn mysql_and_sqlite_migrations_do_not_use_postgres_jsonb() { } } +#[test] +fn worker_boot_cleanup_migration_is_enabled_for_every_driver() { + const VERSION: i64 = 20260731000000; + + for (driver, migrator) in [ + ("postgres", &POSTGRES_MIGRATOR), + ("mysql", &super::mysql::MIGRATOR), + ("sqlite", &super::sqlite::MIGRATOR), + ] { + let migration = migrator + .iter() + .find(|migration| migration.version == VERSION) + .unwrap_or_else(|| panic!("{driver} worker boot cleanup migration should be embedded")); + let sql = migration.sql.as_ref(); + + for required in [ + "DELETE FROM background_task_events", + "DELETE FROM background_task_runs", + "id LIKE 'boot:%'", + "owner_instance IS NOT NULL", + "created_by = 'system'", + "progress_message = 'worker booted'", + ] { + assert!( + sql.contains(required), + "{driver} worker boot cleanup migration is missing {required}" + ); + } + + assert!( + sql.find("DELETE FROM background_task_events") + < sql.find("DELETE FROM background_task_runs"), + "{driver} must delete child events before worker boot runs" + ); + } +} + #[test] fn mysql_and_sqlite_migrations_include_enabled_incrementals() { let mysql_versions = super::mysql::MIGRATOR @@ -1025,6 +1063,7 @@ fn mysql_and_sqlite_migrations_include_enabled_incrementals() { 20260725020000, 20260725030000, 20260727000000, + 20260731000000, ] ); assert_eq!( @@ -1058,6 +1097,7 @@ fn mysql_and_sqlite_migrations_include_enabled_incrementals() { 20260725030000, 20260725040000, 20260727000000, + 20260731000000, ] ); } @@ -2162,6 +2202,8 @@ fn pending_migrations_from_applied_skips_versions_already_applied() { 20260718000000, 20260718010000, 20260720000000, + 20260727000000, + 20260731000000, ] ); }