pub mod dto;
pub mod error;
pub mod handlers;
pub mod routes;
pub(crate) mod stream_reader;
pub use routes::register_routes;
use std::convert::Infallible;
use async_trait::async_trait;
use axum::body::Body;
use axum::http::{HeaderValue, StatusCode, header};
use axum::response::Response;
use futures::stream::StreamExt;
use uuid::Uuid;
use crate::domain::error::ChatEngineError;
use crate::domain::service::webhook::{
NoopWebhookEmitter as DomainNoopWebhookEmitter, WebhookEmitter as DomainWebhookEmitter,
WebhookEvent,
};
pub(crate) fn sse_delta_stream_response(
stream: crate::domain::service::message_service::SendMessageStream,
) -> Response {
use crate::domain::stream_delta::DeltaProjector;
use futures::stream;
let wire = stream
.scan(DeltaProjector::new(), |proj, evt| {
std::future::ready(Some(stream::iter(proj.project(evt))))
})
.flatten();
let body_stream = wire.map(|w| std::result::Result::<_, Infallible>::Ok(sse_frame(&w)));
Response::builder()
.status(StatusCode::OK)
.header(
header::CONTENT_TYPE,
HeaderValue::from_static("text/event-stream"),
)
.header(header::CACHE_CONTROL, HeaderValue::from_static("no-cache"))
.header("x-accel-buffering", HeaderValue::from_static("no"))
.body(Body::from_stream(body_stream))
.unwrap_or_else(|err| {
tracing::error!(error = %err, "failed to build SSE stream response");
let mut resp = Response::new(Body::empty());
*resp.status_mut() = StatusCode::INTERNAL_SERVER_ERROR;
resp
})
}
fn sse_frame(evt: &crate::domain::stream_delta::WireStreamEvent) -> Vec<u8> {
let data = serde_json::to_string(evt).unwrap_or_else(|err| {
tracing::error!(error = %err, "failed to serialize wire delta event");
r#"{"type":"message.error","error":"internal serialization failure"}"#.to_string()
});
format!(
"id: {}\nevent: {}\ndata: {}\n\n",
evt.seq(),
evt.event_name(),
data
)
.into_bytes()
}
#[async_trait]
pub trait WebhookEmitter: Send + Sync {
async fn emit_session_created(
&self,
session_id: Uuid,
tenant_id: &str,
user_id: &str,
session_type_id: Option<Uuid>,
) -> Result<(), ChatEngineError>;
async fn emit_message_new(
&self,
session_id: Uuid,
message_id: Uuid,
tenant_id: &str,
user_id: &str,
) -> Result<(), ChatEngineError>;
async fn emit_message_recreate(
&self,
session_id: Uuid,
message_id: Uuid,
tenant_id: &str,
user_id: &str,
) -> Result<(), ChatEngineError>;
async fn emit_message_aborted(
&self,
session_id: Uuid,
message_id: Uuid,
tenant_id: &str,
user_id: &str,
) -> Result<(), ChatEngineError>;
async fn emit_session_deleted(
&self,
session_id: Uuid,
tenant_id: &str,
user_id: &str,
hard: bool,
) -> Result<(), ChatEngineError>;
async fn emit_session_summary(
&self,
session_id: Uuid,
tenant_id: &str,
user_id: &str,
) -> Result<(), ChatEngineError>;
async fn emit_session_type_health_check(
&self,
session_type_id: Uuid,
) -> Result<(), ChatEngineError>;
}
#[derive(Debug, Default, Clone)]
pub struct NoopWebhookEmitter {
inner: DomainNoopWebhookEmitter,
}
#[async_trait]
impl WebhookEmitter for NoopWebhookEmitter {
async fn emit_session_created(
&self,
session_id: Uuid,
tenant_id: &str,
user_id: &str,
session_type_id: Option<Uuid>,
) -> Result<(), ChatEngineError> {
self.inner
.emit(WebhookEvent::SessionCreated {
session_id,
tenant_id: tenant_id.into(),
user_id: user_id.into(),
session_type_id,
})
.await
}
async fn emit_message_new(
&self,
session_id: Uuid,
message_id: Uuid,
tenant_id: &str,
user_id: &str,
) -> Result<(), ChatEngineError> {
tracing::debug!(
event = "message.new",
%session_id,
%message_id,
tenant_id,
user_id,
"webhook emitter (noop) \u{2014} event swallowed",
);
Ok(())
}
async fn emit_message_recreate(
&self,
session_id: Uuid,
message_id: Uuid,
tenant_id: &str,
user_id: &str,
) -> Result<(), ChatEngineError> {
tracing::debug!(
event = "message.recreate",
%session_id,
%message_id,
tenant_id,
user_id,
"webhook emitter (noop) \u{2014} event swallowed",
);
Ok(())
}
async fn emit_message_aborted(
&self,
session_id: Uuid,
message_id: Uuid,
tenant_id: &str,
user_id: &str,
) -> Result<(), ChatEngineError> {
tracing::debug!(
event = "message.aborted",
%session_id,
%message_id,
tenant_id,
user_id,
"webhook emitter (noop) \u{2014} event swallowed",
);
Ok(())
}
async fn emit_session_deleted(
&self,
session_id: Uuid,
tenant_id: &str,
user_id: &str,
hard: bool,
) -> Result<(), ChatEngineError> {
let event = if hard {
WebhookEvent::SessionHardDeleted {
session_id,
tenant_id: tenant_id.into(),
user_id: user_id.into(),
}
} else {
WebhookEvent::SessionSoftDeleted {
session_id,
tenant_id: tenant_id.into(),
user_id: user_id.into(),
}
};
self.inner.emit(event).await
}
async fn emit_session_summary(
&self,
session_id: Uuid,
tenant_id: &str,
user_id: &str,
) -> Result<(), ChatEngineError> {
tracing::debug!(
event = "session.summary",
%session_id,
tenant_id,
user_id,
"webhook emitter (noop) \u{2014} event swallowed",
);
Ok(())
}
async fn emit_session_type_health_check(
&self,
session_type_id: Uuid,
) -> Result<(), ChatEngineError> {
tracing::debug!(
event = "session_type.health_check",
%session_type_id,
"webhook emitter (noop) \u{2014} event swallowed",
);
Ok(())
}
}
pub struct WebhookEmitterAdapter<E: WebhookEmitter + ?Sized> {
inner: std::sync::Arc<E>,
}
impl<E> WebhookEmitterAdapter<E>
where
E: WebhookEmitter + ?Sized,
{
#[must_use]
pub fn new(inner: std::sync::Arc<E>) -> Self {
Self { inner }
}
}
#[async_trait]
impl<E> DomainWebhookEmitter for WebhookEmitterAdapter<E>
where
E: WebhookEmitter + ?Sized + Send + Sync,
{
async fn emit(&self, event: WebhookEvent) -> Result<(), ChatEngineError> {
match event {
WebhookEvent::SessionCreated {
session_id,
tenant_id,
user_id,
session_type_id,
} => {
self.inner
.emit_session_created(session_id, &tenant_id, &user_id, session_type_id)
.await
}
WebhookEvent::SessionArchived {
session_id,
tenant_id,
user_id,
}
| WebhookEvent::SessionRestored {
session_id,
tenant_id,
user_id,
} => {
tracing::debug!(
%session_id,
tenant_id,
user_id,
"lifecycle event (archived/restored) \u{2014} no dedicated webhook method",
);
Ok(())
}
WebhookEvent::SessionSoftDeleted {
session_id,
tenant_id,
user_id,
} => {
self.inner
.emit_session_deleted(session_id, &tenant_id, &user_id, false)
.await
}
WebhookEvent::SessionHardDeleted {
session_id,
tenant_id,
user_id,
} => {
self.inner
.emit_session_deleted(session_id, &tenant_id, &user_id, true)
.await
}
WebhookEvent::MessageDeleted {
session_id,
message_id,
tenant_id,
user_id,
..
} => {
self.inner
.emit_message_aborted(session_id, message_id, &tenant_id, &user_id)
.await
}
}
}
}
#[cfg(test)]
#[path = "mod_tests.rs"]
mod mod_tests;