diff --git a/crates/erp-core/src/entity/mod.rs b/crates/erp-core/src/entity/mod.rs index f9b0997..00a81aa 100644 --- a/crates/erp-core/src/entity/mod.rs +++ b/crates/erp-core/src/entity/mod.rs @@ -1,2 +1,3 @@ pub mod audit_log; pub mod domain_event; +pub mod processed_event; diff --git a/crates/erp-core/src/entity/processed_event.rs b/crates/erp-core/src/entity/processed_event.rs new file mode 100644 index 0000000..b03d124 --- /dev/null +++ b/crates/erp-core/src/entity/processed_event.rs @@ -0,0 +1,18 @@ +use sea_orm::entity::prelude::*; +use serde::{Deserialize, Serialize}; + +/// 已处理事件记录 — 幂等性去重表。 +#[derive(Clone, Debug, PartialEq, DeriveEntityModel, Serialize, Deserialize)] +#[sea_orm(table_name = "processed_events")] +pub struct Model { + #[sea_orm(primary_key, auto_increment = false)] + pub event_id: Uuid, + #[sea_orm(primary_key, auto_increment = false)] + pub consumer_id: String, + pub processed_at: DateTimeUtc, +} + +#[derive(Copy, Clone, Debug, EnumIter, DeriveRelation)] +pub enum Relation {} + +impl ActiveModelBehavior for ActiveModel {} diff --git a/crates/erp-core/src/events.rs b/crates/erp-core/src/events.rs index a01cc4a..2d78d7e 100644 --- a/crates/erp-core/src/events.rs +++ b/crates/erp-core/src/events.rs @@ -1,5 +1,5 @@ use chrono::Utc; -use sea_orm::{ActiveModelTrait, ConnectionTrait, Set}; +use sea_orm::{ActiveModelTrait, ConnectionTrait, PaginatorTrait, Set}; use serde::{Deserialize, Serialize}; use tokio::sync::{broadcast, mpsc}; use tracing::{error, info}; @@ -53,6 +53,54 @@ pub fn build_event_payload(data: serde_json::Value) -> serde_json::Value { envelope } +/// 检查事件是否已被指定消费者处理。 +/// +/// 查询 `processed_events` 表判断 event_id + consumer_id 是否已存在。 +pub async fn is_event_processed( + db: &sea_orm::DatabaseConnection, + event_id: Uuid, + consumer_id: &str, +) -> Result { + use sea_orm::{ColumnTrait, EntityTrait, QueryFilter}; + + let count = crate::entity::processed_event::Entity::find() + .filter(crate::entity::processed_event::Column::EventId.eq(event_id)) + .filter(crate::entity::processed_event::Column::ConsumerId.eq(consumer_id)) + .count(db) + .await?; + Ok(count > 0) +} + +/// 标记事件已被指定消费者处理。 +/// +/// 插入 `processed_events` 记录,重复插入会因主键冲突被安全忽略。 +pub async fn mark_event_processed( + db: &sea_orm::DatabaseConnection, + event_id: Uuid, + consumer_id: &str, +) -> Result<(), sea_orm::DbErr> { + use sea_orm::ActiveModelTrait; + use sea_orm::Set; + + let model = crate::entity::processed_event::ActiveModel { + event_id: Set(event_id), + consumer_id: Set(consumer_id.to_string()), + processed_at: Set(Utc::now()), + }; + // INSERT ... ON CONFLICT DO NOTHING(主键冲突时安全忽略) + match model.insert(db).await { + Ok(_) => Ok(()), + Err(e) => { + // 唯一约束冲突 = 已处理,不是错误 + if e.to_string().contains("duplicate") || e.to_string().contains("violates unique") { + Ok(()) + } else { + Err(e) + } + } + } +} + /// 过滤事件接收器 — 只接收匹配 `event_type_prefix` 的事件 pub struct FilteredEventReceiver { receiver: mpsc::Receiver, diff --git a/crates/erp-server/migration/src/lib.rs b/crates/erp-server/migration/src/lib.rs index 82537f4..f68b5cc 100644 --- a/crates/erp-server/migration/src/lib.rs +++ b/crates/erp-server/migration/src/lib.rs @@ -83,6 +83,8 @@ mod m20260427_000080_create_medication_record; mod m20260427_000081_create_dialysis_prescription; mod m20260427_000082_seed_ai_prompts; mod m20260427_000083_create_follow_up_template; +mod m20260427_000084_domain_events_cleanup; +mod m20260427_000085_processed_events; pub struct Migrator; @@ -173,6 +175,8 @@ impl MigratorTrait for Migrator { Box::new(m20260427_000081_create_dialysis_prescription::Migration), Box::new(m20260427_000082_seed_ai_prompts::Migration), Box::new(m20260427_000083_create_follow_up_template::Migration), + Box::new(m20260427_000084_domain_events_cleanup::Migration), + Box::new(m20260427_000085_processed_events::Migration), ] } } diff --git a/crates/erp-server/migration/src/m20260427_000084_domain_events_cleanup.rs b/crates/erp-server/migration/src/m20260427_000084_domain_events_cleanup.rs new file mode 100644 index 0000000..faa0ea1 --- /dev/null +++ b/crates/erp-server/migration/src/m20260427_000084_domain_events_cleanup.rs @@ -0,0 +1,89 @@ +use sea_orm_migration::prelude::*; + +#[derive(DeriveMigrationName)] +pub struct Migration; + +#[async_trait::async_trait] +impl MigrationTrait for Migration { + async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> { + // 归档表 — 与 domain_events 结构相同,用于存放 >90 天的已发布事件 + manager + .create_table( + Table::create() + .table(Alias::new("domain_events_archive")) + .if_not_exists() + .col(ColumnDef::new(Alias::new("id")).uuid().not_null().primary_key()) + .col(ColumnDef::new(Alias::new("tenant_id")).uuid().not_null()) + .col(ColumnDef::new(Alias::new("event_type")).string_len(200).not_null()) + .col(ColumnDef::new(Alias::new("payload")).json().null()) + .col(ColumnDef::new(Alias::new("correlation_id")).uuid().null()) + .col(ColumnDef::new(Alias::new("status")).string_len(20).not_null()) + .col(ColumnDef::new(Alias::new("attempts")).integer().not_null().default(0)) + .col(ColumnDef::new(Alias::new("last_error")).text().null()) + .col(ColumnDef::new(Alias::new("created_at")).timestamp_with_time_zone().not_null()) + .col(ColumnDef::new(Alias::new("published_at")).timestamp_with_time_zone().null()) + .col(ColumnDef::new(Alias::new("archived_at")).timestamp_with_time_zone().not_null().default(Expr::current_timestamp())) + .to_owned(), + ) + .await?; + + manager + .create_index( + Index::create() + .name("idx_domain_events_archive_created") + .table(Alias::new("domain_events_archive")) + .col(Alias::new("created_at")) + .to_owned(), + ) + .await?; + + // 清理函数:将 >90 天的已发布事件迁移到归档表 + manager + .get_connection() + .execute_unprepared( + r#" + CREATE OR REPLACE FUNCTION cleanup_old_published_events( + retention_days INT DEFAULT 90, + batch_size INT DEFAULT 1000 + ) RETURNS INT AS $$ + DECLARE + moved_count INT; + BEGIN + INSERT INTO domain_events_archive (id, tenant_id, event_type, payload, correlation_id, status, attempts, last_error, created_at, published_at) + SELECT id, tenant_id, event_type, payload, correlation_id, status, attempts, last_error, created_at, published_at + FROM domain_events + WHERE status = 'published' + AND published_at < NOW() - (retention_days || ' days')::INTERVAL + ORDER BY created_at ASC + LIMIT batch_size; + + GET DIAGNOSTICS moved_count = ROW_COUNT; + + DELETE FROM domain_events + WHERE status = 'published' + AND published_at < NOW() - (retention_days || ' days')::INTERVAL + LIMIT batch_size; + + RETURN moved_count; + END; + $$ LANGUAGE plpgsql; + "#, + ) + .await?; + + Ok(()) + } + + async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> { + manager + .get_connection() + .execute_unprepared("DROP FUNCTION IF EXISTS cleanup_old_published_events(INT, INT);") + .await?; + + manager + .drop_table(Table::drop().table(Alias::new("domain_events_archive")).to_owned()) + .await?; + + Ok(()) + } +} diff --git a/crates/erp-server/migration/src/m20260427_000085_processed_events.rs b/crates/erp-server/migration/src/m20260427_000085_processed_events.rs new file mode 100644 index 0000000..b619713 --- /dev/null +++ b/crates/erp-server/migration/src/m20260427_000085_processed_events.rs @@ -0,0 +1,61 @@ +use sea_orm_migration::prelude::*; + +#[derive(DeriveMigrationName)] +pub struct Migration; + +#[async_trait::async_trait] +impl MigrationTrait for Migration { + async fn up(&self, manager: &SchemaManager) -> Result<(), DbErr> { + manager + .create_table( + Table::create() + .table(Alias::new("processed_events")) + .if_not_exists() + .col(ColumnDef::new(Alias::new("event_id")).uuid().not_null()) + .col(ColumnDef::new(Alias::new("consumer_id")).string_len(200).not_null()) + .col(ColumnDef::new(Alias::new("processed_at")).timestamp_with_time_zone().not_null().default(Expr::current_timestamp())) + .primary_key(Index::create().col(Alias::new("event_id")).col(Alias::new("consumer_id"))) + .to_owned(), + ) + .await?; + + // 7 天 TTL 清理函数 + manager + .get_connection() + .execute_unprepared( + r#" + CREATE OR REPLACE FUNCTION cleanup_old_processed_events( + retention_days INT DEFAULT 7, + batch_size INT DEFAULT 1000 + ) RETURNS INT AS $$ + DECLARE + deleted_count INT; + BEGIN + DELETE FROM processed_events + WHERE processed_at < NOW() - (retention_days || ' days')::INTERVAL + LIMIT batch_size; + + GET DIAGNOSTICS deleted_count = ROW_COUNT; + RETURN deleted_count; + END; + $$ LANGUAGE plpgsql; + "#, + ) + .await?; + + Ok(()) + } + + async fn down(&self, manager: &SchemaManager) -> Result<(), DbErr> { + manager + .get_connection() + .execute_unprepared("DROP FUNCTION IF EXISTS cleanup_old_processed_events(INT, INT);") + .await?; + + manager + .drop_table(Table::drop().table(Alias::new("processed_events")).to_owned()) + .await?; + + Ok(()) + } +} diff --git a/crates/erp-server/src/main.rs b/crates/erp-server/src/main.rs index 93787fd..a10d408 100644 --- a/crates/erp-server/src/main.rs +++ b/crates/erp-server/src/main.rs @@ -4,6 +4,7 @@ mod handlers; mod middleware; mod outbox; mod state; +mod tasks; /// OpenAPI 规范定义 — 通过 utoipa derive 合并各模块 schema。 #[derive(OpenApi)] @@ -410,6 +411,9 @@ async fn main() -> anyhow::Result<()> { outbox::start_outbox_relay(db.clone(), event_bus.clone(), config.database.url.clone()); tracing::info!("Outbox relay started"); + // Start event cleanup (archive old published events + purge processed_events) + tasks::start_event_cleanup(db.clone()); + // Start timeout checker (scan overdue tasks every 60s) erp_workflow::WorkflowModule::start_timeout_checker(db.clone(), event_bus.clone()); tracing::info!("Timeout checker started"); diff --git a/crates/erp-server/src/tasks.rs b/crates/erp-server/src/tasks.rs new file mode 100644 index 0000000..b3105d2 --- /dev/null +++ b/crates/erp-server/src/tasks.rs @@ -0,0 +1,53 @@ +use std::time::Duration; + +/// 启动事件清理后台任务。 +/// +/// 每日执行一次: +/// - 调用 `cleanup_old_published_events()` 归档 >90 天的已发布事件 +/// - 调用 `cleanup_old_processed_events()` 清理 >7 天的去重记录 +pub fn start_event_cleanup(db: sea_orm::DatabaseConnection) { + tokio::spawn(async move { + let mut interval = tokio::time::interval(Duration::from_secs(86400)); + loop { + interval.tick().await; + if let Err(e) = run_cleanup(&db).await { + tracing::warn!(error = %e, "事件清理任务执行失败"); + } + } + }); + tracing::info!("事件清理任务已启动(每 24 小时执行一次)"); +} + +async fn run_cleanup(db: &sea_orm::DatabaseConnection) -> Result<(), sea_orm::DbErr> { + use sea_orm::ConnectionTrait; + + // 归档 >90 天的已发布事件 + match db + .execute_unprepared("SELECT cleanup_old_published_events(90, 1000)") + .await + { + Ok(result) => { + tracing::info!( + rows_affected = result.rows_affected(), + "已发布事件归档完成" + ); + } + Err(e) => tracing::warn!(error = %e, "已发布事件归档失败"), + } + + // 清理 >7 天的去重记录 + match db + .execute_unprepared("SELECT cleanup_old_processed_events(7, 1000)") + .await + { + Ok(result) => { + tracing::info!( + rows_affected = result.rows_affected(), + "去重记录清理完成" + ); + } + Err(e) => tracing::warn!(error = %e, "去重记录清理失败"), + } + + Ok(()) +}