feat(core): 事件归档 + 消费者幂等性 — 迁移 084/085 + 清理任务
- 迁移 084: domain_events_archive 归档表 + cleanup_old_published_events() - 迁移 085: processed_events 去重表 + cleanup_old_processed_events() - erp-core: is_event_processed() / mark_event_processed() 幂等性辅助 - erp-server: tasks::start_event_cleanup() 每 24h 归档 >90 天事件
This commit is contained in:
@@ -1,2 +1,3 @@
|
||||
pub mod audit_log;
|
||||
pub mod domain_event;
|
||||
pub mod processed_event;
|
||||
|
||||
18
crates/erp-core/src/entity/processed_event.rs
Normal file
18
crates/erp-core/src/entity/processed_event.rs
Normal file
@@ -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 {}
|
||||
@@ -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<bool, sea_orm::DbErr> {
|
||||
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<DomainEvent>,
|
||||
|
||||
@@ -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),
|
||||
]
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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(())
|
||||
}
|
||||
}
|
||||
@@ -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(())
|
||||
}
|
||||
}
|
||||
@@ -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");
|
||||
|
||||
53
crates/erp-server/src/tasks.rs
Normal file
53
crates/erp-server/src/tasks.rs
Normal file
@@ -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(())
|
||||
}
|
||||
Reference in New Issue
Block a user