use std::sync::{Arc, OnceLock};
use std::time::Duration;
use async_trait::async_trait;
use axum::Router;
use sea_orm_migration::MigrationTrait;
use tokio_util::sync::CancellationToken;
use toolkit::api::OpenApiRegistry;
use toolkit::api::canonical_error_middleware;
use toolkit::client_hub::ClientScope;
use toolkit::{DatabaseCapability, Gear, GearCtx, RestApiCapability};
use toolkit_db::DBProvider;
use tracing::{error, info, warn};
use crate::infra::db::repo::ChatEngineDb;
use chat_engine_sdk::plugin::ChatEngineBackendPlugin;
use crate::api::rest::routes::ChatEngineServices;
use crate::api::rest::{NoopWebhookEmitter, WebhookEmitter, WebhookEmitterAdapter};
use crate::config::ChatEngineConfig;
use crate::domain::export::NotImplementedExportStorage;
use crate::domain::service::webhook::WebhookEmitter as DomainWebhookEmitter;
use crate::domain::service::{
ExportService, IntelligenceService, MessageService, PluginService, ReactionService,
SearchService, SessionService, ShareUrlBuilder, VariantService,
};
use crate::infra::db::migrations::Migrator;
use crate::infra::db::repo::message_repo::SeaMessageRepo;
use crate::infra::db::repo::plugin_config_repo::SeaPluginConfigRepo;
use crate::infra::db::repo::reaction_repo::SeaReactionRepo;
use crate::infra::db::repo::session_repo::SeaSessionRepo;
use crate::infra::db::repo::session_type_repo::SeaSessionTypeRepo;
use crate::infra::leader::{LeaderElector, work_fn};
use crate::infra::llm_gateway::LlmGatewayPlugin;
use crate::infra::search::NotImplementedSearchBackend;
use crate::infra::webhook_compat::WebhookCompatPlugin;
pub const DEFAULT_WEBHOOK_COMPAT_INSTANCE_ID: &str = "gtx.cf.chat_engine.webhook_compat_plugin.v1~";
struct RuntimeState {
services: ChatEngineServices,
webhooks: Arc<dyn WebhookEmitter>,
intelligence: Arc<IntelligenceService>,
stream_buffer: Arc<dyn crate::domain::ports::StreamEventBuffer>,
config: Arc<ChatEngineConfig>,
leader: Arc<dyn LeaderElector>,
}
const STREAM_BUFFER_SWEEP_PERIOD: Duration = Duration::from_mins(5);
#[toolkit::gear(
name = "chat-engine",
capabilities = [db, rest, stateful],
client = chat_engine_sdk::ChatEngineBackendPlugin,
ctor = ChatEngineModule::new(),
lifecycle(entry = "serve", stop_timeout = "30s", await_ready)
)]
pub struct ChatEngineModule {
runtime: OnceLock<RuntimeState>,
}
impl Default for ChatEngineModule {
fn default() -> Self {
Self::new()
}
}
impl ChatEngineModule {
#[must_use]
pub fn new() -> Self {
Self {
runtime: OnceLock::new(),
}
}
fn runtime(&self) -> anyhow::Result<&RuntimeState> {
self.runtime
.get()
.ok_or_else(|| anyhow::anyhow!("ChatEngineModule not initialised"))
}
pub async fn serve(
self: Arc<Self>,
cancel: CancellationToken,
ready: toolkit::lifecycle::ReadySignal,
) -> anyhow::Result<()> {
let runtime = self.runtime()?;
let interval_hours = runtime.config.retention_cleanup_interval_hours;
let period = Duration::from_secs(interval_hours.saturating_mul(3600));
let intelligence = Arc::clone(&runtime.intelligence);
let stream_buffer = Arc::clone(&runtime.stream_buffer);
let leader = Arc::clone(&runtime.leader);
ready.notify();
info!(
interval_hours,
"chat-engine retention-cleanup task running (leader-gated)"
);
leader
.run_role(
"retention-cleanup",
cancel,
work_fn(move |cancel| {
let intelligence = Arc::clone(&intelligence);
let stream_buffer = Arc::clone(&stream_buffer);
async move {
let mut interval = tokio::time::interval(period);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
interval.tick().await;
let mut sweep = tokio::time::interval(STREAM_BUFFER_SWEEP_PERIOD);
sweep.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
sweep.tick().await;
loop {
tokio::select! {
() = cancel.cancelled() => {
info!(
"chat-engine retention-cleanup loop stopping \
(leadership lost or shutdown)"
);
return Ok(());
}
_ = interval.tick() => {
if let Err(err) =
run_retention_cleanup_tick(intelligence.as_ref()).await
{
error!(
error = %err,
"chat-engine retention-cleanup tick failed; continuing",
);
}
}
_ = sweep.tick() => {
run_stream_buffer_sweep_tick(stream_buffer.as_ref()).await;
}
}
}
}
}),
)
.await
}
}
async fn run_retention_cleanup_tick(intelligence: &IntelligenceService) -> anyhow::Result<()> {
let report = intelligence.run_retention_cleanup_all_tenants().await?;
info!(
sessions_scanned = report.sessions.len(),
sessions_skipped_locked = report.skipped_count(),
total_messages_deleted = report.total_messages_deleted(),
"chat-engine retention-cleanup tick completed"
);
Ok(())
}
async fn run_stream_buffer_sweep_tick(buffer: &dyn crate::domain::ports::StreamEventBuffer) {
match buffer.delete_expired(time::OffsetDateTime::now_utc()).await {
Ok(removed) if removed > 0 => {
info!(
removed,
"chat-engine stream-buffer TTL sweep removed expired events"
);
}
Ok(_) => {}
Err(err) => {
warn!(error = %err, "chat-engine stream-buffer TTL sweep failed; continuing");
}
}
}
#[allow(
clippy::unused_async,
reason = "async is needed when the k8s feature is enabled"
)]
async fn create_leader_elector() -> anyhow::Result<Arc<dyn LeaderElector>> {
#[cfg(feature = "k8s")]
{
use crate::infra::leader::k8s_lease::{K8sLeaseConfig, K8sLeaseElector};
use anyhow::Context as _;
let config = K8sLeaseConfig::from_env("chat-engine")
.context("k8s feature enabled: POD_NAMESPACE and POD_NAME are required")?;
let elector = K8sLeaseElector::from_default(config)
.await
.context("k8s feature enabled: kube client init failed")?;
info!("chat-engine: using k8s Lease leader election");
Ok(Arc::new(elector))
}
#[cfg(not(feature = "k8s"))]
{
info!("chat-engine: using noop leader election (single-process mode)");
Ok(crate::infra::leader::noop())
}
}
#[async_trait]
impl Gear for ChatEngineModule {
async fn init(&self, ctx: &GearCtx) -> anyhow::Result<()> {
info!("initialising {} module", Self::MODULE_NAME);
let cfg: ChatEngineConfig = ctx.config_or_default()?;
cfg.validate()
.map_err(|e| anyhow::anyhow!("invalid chat-engine config: {e}"))?;
let config = Arc::new(cfg);
let leader = create_leader_elector().await?;
let db_raw = ctx.db_required()?;
let db: Arc<ChatEngineDb> = Arc::new(DBProvider::new(db_raw.db()));
let sessions_repo: Arc<dyn crate::domain::ports::SessionRepo> =
Arc::new(SeaSessionRepo::new(Arc::clone(&db)));
let session_types_repo: Arc<dyn crate::domain::ports::SessionTypeRepo> =
Arc::new(SeaSessionTypeRepo::new(Arc::clone(&db)));
let messages_repo: Arc<dyn crate::domain::ports::MessageRepo> =
Arc::new(SeaMessageRepo::new(Arc::clone(&db)));
let plugin_config_repo: Arc<dyn crate::domain::ports::PluginConfigRepo> =
Arc::new(SeaPluginConfigRepo::new(Arc::clone(&db)));
let reactions_repo: Arc<dyn crate::domain::ports::ReactionRepo> =
Arc::new(SeaReactionRepo::new(Arc::clone(&db)));
let variants_repo: Arc<dyn crate::domain::service::VariantRepo> = Arc::new(
crate::infra::db::repo::variant_repo::SeaVariantRepo::new(Arc::clone(&db)),
);
let client_hub = ctx.client_hub();
let webhook_compat = Arc::new(
WebhookCompatPlugin::new(DEFAULT_WEBHOOK_COMPAT_INSTANCE_ID)
.map_err(|e| anyhow::anyhow!("failed to build webhook-compat plugin: {e}"))?,
);
client_hub.register_scoped::<dyn ChatEngineBackendPlugin>(
ClientScope::gts_id(DEFAULT_WEBHOOK_COMPAT_INSTANCE_ID),
webhook_compat.clone() as Arc<dyn ChatEngineBackendPlugin>,
);
if config.llm_gateway_base_url.is_some() {
warn!(
"llm-gateway plugin instantiation requested but production transport clients \
are not yet wired in this build; the plugin slot remains empty"
);
}
let _ = LlmGatewayPlugin::new;
let plugin_service = PluginService::new(client_hub.clone(), plugin_config_repo.clone());
let webhooks_rest: Arc<dyn WebhookEmitter> = Arc::new(NoopWebhookEmitter::default());
let webhooks_domain: Arc<dyn DomainWebhookEmitter> =
Arc::new(WebhookEmitterAdapter::new(webhooks_rest.clone()));
let plugin_deadline = Duration::from_secs(config.plugin_deadline_secs);
let sessions = Arc::new(
SessionService::new(
sessions_repo.clone(),
session_types_repo.clone(),
plugin_service.clone(),
webhooks_domain.clone(),
)
.with_plugin_timeout(plugin_deadline),
);
let stream_buffer: Arc<dyn crate::domain::ports::StreamEventBuffer> = Arc::new(
crate::infra::db::repo::stream_event_repo::SeaStreamEventBuffer::new(Arc::clone(&db)),
);
let messages = Arc::new(
MessageService::new(
sessions_repo.clone(),
session_types_repo.clone(),
messages_repo.clone(),
plugin_service.clone(),
)
.with_webhook_emitter(webhooks_domain.clone())
.with_streaming_buffer_size(config.ndjson_buffer_size)
.with_plugin_deadline(plugin_deadline)
.with_stream_buffer(Arc::clone(&stream_buffer)),
);
let variants = Arc::new(
VariantService::new(
sessions_repo.clone(),
session_types_repo.clone(),
messages_repo.clone(),
variants_repo.clone(),
plugin_service.clone(),
Arc::clone(&messages),
)
.with_plugin_timeout(plugin_deadline),
);
let reactions = Arc::new(ReactionService::new(
sessions_repo.clone(),
session_types_repo.clone(),
messages_repo.clone(),
reactions_repo.clone(),
plugin_service.clone(),
));
let search_backend: Arc<dyn crate::domain::service::SearchBackend> =
Arc::new(NotImplementedSearchBackend::new());
let search = Arc::new(SearchService::new(
sessions_repo.clone(),
messages_repo.clone(),
search_backend,
));
let intelligence = Arc::new(
IntelligenceService::new(
sessions_repo.clone(),
session_types_repo.clone(),
messages_repo.clone(),
plugin_service.clone(),
)
.with_buffer_size(config.summary_buffer_size)
.with_summary_deadline(plugin_deadline)
.with_retention_caps(
config.retention_max_sessions_per_tick,
config.retention_max_deletes_per_session,
),
);
let share_urls =
config
.share_base_url
.as_ref()
.map_or_else(ShareUrlBuilder::default, |base| ShareUrlBuilder {
base_url: base.clone(),
});
let export_storage = Arc::new(NotImplementedExportStorage);
let export = Arc::new(
ExportService::new(sessions_repo.clone(), messages_repo.clone(), export_storage)
.with_share_urls(share_urls),
);
let services = ChatEngineServices {
sessions,
messages,
variants,
reactions,
search,
intelligence: Arc::clone(&intelligence),
export,
};
let runtime = RuntimeState {
services,
webhooks: webhooks_rest,
intelligence,
stream_buffer,
config,
leader,
};
self.runtime
.set(runtime)
.map_err(|_| anyhow::anyhow!("chat-engine module already initialised"))?;
info!("{} module initialised", Self::MODULE_NAME);
Ok(())
}
}
impl DatabaseCapability for ChatEngineModule {
fn migrations(&self) -> Vec<Box<dyn MigrationTrait>> {
use sea_orm_migration::MigratorTrait;
Migrator::migrations()
}
}
impl RestApiCapability for ChatEngineModule {
fn register_rest(
&self,
_ctx: &GearCtx,
router: Router,
openapi: &dyn OpenApiRegistry,
) -> anyhow::Result<Router> {
let runtime = self.runtime()?;
let router = router.layer(axum::middleware::from_fn(canonical_error_middleware));
if !runtime.config.enable_search {
info!(
"chat-engine search endpoints disabled (enable_search=false); \
production search backends are still stubs",
);
}
let router = crate::api::rest::register_routes(
router,
openapi,
runtime.services.clone(),
Arc::clone(&runtime.webhooks),
Arc::clone(&runtime.stream_buffer),
runtime.config.enable_search,
);
Ok(router)
}
}