use std::time::{Duration, Instant};
use async_trait::async_trait;
use futures::stream::{self, BoxStream, StreamExt};
use tokio_util::sync::CancellationToken;
use uuid::Uuid;
use crate::error::PluginError;
use crate::models::{
Capability, CapabilityValue, HealthStatus, Message, StreamingEvent, TenantId, UserId,
};
pub type PluginStream = BoxStream<'static, Result<StreamingEvent, PluginError>>;
#[must_use]
pub fn empty_stream() -> PluginStream {
stream::empty().boxed()
}
#[must_use]
pub fn stream_from_events(events: Vec<StreamingEvent>) -> PluginStream {
stream::iter(events.into_iter().map(Ok)).boxed()
}
#[allow(clippy::module_name_repetitions)]
#[derive(Debug, Clone)]
pub struct SessionPluginCtx {
pub session_type_id: Uuid,
pub session_id: Option<Uuid>,
pub call_ctx: PluginCallContext,
}
#[allow(clippy::module_name_repetitions)]
#[derive(Clone)]
pub struct MessagePluginCtx {
pub session_id: Uuid,
pub message_id: Uuid,
pub messages: Vec<Message>,
pub call_ctx: PluginCallContext,
}
impl std::fmt::Debug for MessagePluginCtx {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let mut user = 0_usize;
let mut assistant = 0_usize;
let mut system = 0_usize;
for m in &self.messages {
match m.role {
crate::models::MessageRole::User => user += 1,
crate::models::MessageRole::Assistant => assistant += 1,
crate::models::MessageRole::System => system += 1,
}
}
let messages_summary = format_args!(
"<redacted: {} message(s); user={user}, assistant={assistant}, system={system}>",
self.messages.len()
)
.to_string();
f.debug_struct("MessagePluginCtx")
.field("session_id", &self.session_id)
.field("message_id", &self.message_id)
.field("messages", &messages_summary)
.field("call_ctx", &self.call_ctx)
.finish()
}
}
#[allow(clippy::module_name_repetitions)]
#[derive(Clone)]
pub struct PluginCallContext {
pub request_id: Uuid,
pub tenant_id: TenantId,
pub user_id: UserId,
pub plugin_instance_id: String,
pub session_type_id: Uuid,
pub plugin_config: Option<serde_json::Value>,
pub enabled_capabilities: Option<Vec<CapabilityValue>>,
pub deadline: Option<Instant>,
pub cancel: CancellationToken,
}
impl PluginCallContext {
#[must_use]
pub fn is_cancelled(&self) -> bool {
self.cancel.is_cancelled()
}
#[must_use]
pub fn remaining(&self) -> Option<Duration> {
self.deadline.map(|d| {
d.checked_duration_since(Instant::now())
.unwrap_or(Duration::ZERO)
})
}
}
impl std::fmt::Debug for PluginCallContext {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let plugin_config_redacted: Option<&'static str> =
self.plugin_config.as_ref().map(|_| "<redacted>");
f.debug_struct("PluginCallContext")
.field("request_id", &self.request_id)
.field("tenant_id", &self.tenant_id)
.field("user_id", &self.user_id)
.field("plugin_instance_id", &self.plugin_instance_id)
.field("session_type_id", &self.session_type_id)
.field("plugin_config", &plugin_config_redacted)
.field("enabled_capabilities", &self.enabled_capabilities)
.field("remaining", &self.remaining())
.field("cancelled", &self.is_cancelled())
.finish()
}
}
#[cfg(test)]
mod plugin_call_context_tests {
use super::{CancellationToken, Duration, Instant, PluginCallContext, TenantId, UserId};
use uuid::Uuid;
fn make_ctx() -> PluginCallContext {
PluginCallContext {
request_id: Uuid::nil(),
tenant_id: TenantId::new("t"),
user_id: UserId::new("u"),
plugin_instance_id: "p".into(),
session_type_id: Uuid::nil(),
plugin_config: None,
enabled_capabilities: None,
deadline: None,
cancel: CancellationToken::new(),
}
}
#[test]
fn debug_redacts_plugin_config_when_present() {
let mut ctx = make_ctx();
ctx.plugin_config = Some(serde_json::json!({"api_key": "super-secret-123"}));
let printed = format!("{ctx:?}");
assert!(printed.contains("<redacted>"), "got: {printed}");
assert!(
!printed.contains("super-secret-123"),
"secret leaked: {printed}"
);
}
#[test]
fn debug_prints_none_when_plugin_config_absent() {
let ctx = make_ctx();
let printed = format!("{ctx:?}");
assert!(printed.contains("plugin_config: None"), "got: {printed}");
assert!(!printed.contains("<redacted>"), "got: {printed}");
}
#[test]
fn is_cancelled_reflects_token_state() {
let ctx = make_ctx();
assert!(!ctx.is_cancelled());
ctx.cancel.cancel();
assert!(ctx.is_cancelled());
}
#[test]
fn remaining_is_none_when_no_deadline() {
let ctx = make_ctx();
assert!(ctx.remaining().is_none());
}
#[test]
fn remaining_returns_positive_duration_for_future_deadline() {
let mut ctx = make_ctx();
ctx.deadline = Some(Instant::now() + Duration::from_secs(10));
let r = ctx.remaining().expect("should be set");
assert!(r > Duration::from_secs(5) && r <= Duration::from_secs(10));
}
#[test]
fn remaining_is_zero_when_deadline_already_elapsed() {
let mut ctx = make_ctx();
ctx.deadline = Some(
Instant::now()
.checked_sub(Duration::from_secs(1))
.expect("monotonic clock is at least 1s past its reference"),
);
assert_eq!(ctx.remaining(), Some(Duration::ZERO));
}
#[test]
fn remaining_is_zero_at_exact_deadline() {
let mut ctx = make_ctx();
ctx.deadline = Some(Instant::now());
let r = ctx.remaining().expect("deadline is set");
assert!(r <= Duration::from_millis(1), "expected ~ZERO, got {r:?}");
}
}
#[cfg(test)]
mod message_plugin_ctx_debug_tests {
use super::{CancellationToken, MessagePluginCtx, PluginCallContext, TenantId, UserId};
use crate::models::{Message, MessageRole};
use time::OffsetDateTime;
use uuid::Uuid;
fn make_message(role: MessageRole, secret_text: &str) -> Message {
let now = OffsetDateTime::now_utc();
Message {
message_id: Uuid::nil(),
session_id: Uuid::nil(),
tenant_id: None,
user_id: None,
parent_message_id: None,
variant_index: 0,
is_active: true,
role,
parts: vec![crate::models::MessagePart::text(
Uuid::nil(),
Uuid::nil(),
0,
secret_text,
)],
file_ids: vec![],
metadata: None,
is_complete: true,
is_hidden_from_user: false,
is_hidden_from_backend: false,
created_at: now,
updated_at: now,
}
}
fn make_call_ctx() -> PluginCallContext {
PluginCallContext {
request_id: Uuid::nil(),
tenant_id: TenantId::new("t"),
user_id: UserId::new("u"),
plugin_instance_id: "p".into(),
session_type_id: Uuid::nil(),
plugin_config: None,
enabled_capabilities: None,
deadline: None,
cancel: CancellationToken::new(),
}
}
#[test]
fn debug_redacts_message_content_but_shows_per_role_summary() {
let ctx = MessagePluginCtx {
session_id: Uuid::nil(),
message_id: Uuid::nil(),
messages: vec![
make_message(MessageRole::System, "system-secret-prompt"),
make_message(MessageRole::User, "i had a heart attack last night"),
make_message(MessageRole::Assistant, "private-response-PII"),
make_message(MessageRole::User, "another sensitive question"),
],
call_ctx: make_call_ctx(),
};
let printed = format!("{ctx:?}");
assert!(
!printed.contains("system-secret-prompt"),
"system content leaked: {printed}"
);
assert!(
!printed.contains("heart attack"),
"user PII leaked: {printed}"
);
assert!(
!printed.contains("private-response-PII"),
"assistant content leaked: {printed}"
);
assert!(
!printed.contains("sensitive question"),
"user content leaked: {printed}"
);
assert!(printed.contains("4 message(s)"), "got: {printed}");
assert!(printed.contains("user=2"), "got: {printed}");
assert!(printed.contains("assistant=1"), "got: {printed}");
assert!(printed.contains("system=1"), "got: {printed}");
assert!(printed.contains("<redacted"), "got: {printed}");
}
#[test]
fn debug_shows_zero_counts_for_empty_history() {
let ctx = MessagePluginCtx {
session_id: Uuid::nil(),
message_id: Uuid::nil(),
messages: vec![],
call_ctx: make_call_ctx(),
};
let printed = format!("{ctx:?}");
assert!(printed.contains("0 message(s)"), "got: {printed}");
assert!(printed.contains("user=0"), "got: {printed}");
}
}
#[derive(Debug, Clone, Default)]
pub struct SessionPluginResponse {
pub capabilities: Vec<Capability>,
pub metadata: Option<serde_json::Value>,
}
impl From<Vec<Capability>> for SessionPluginResponse {
fn from(capabilities: Vec<Capability>) -> Self {
Self {
capabilities,
metadata: None,
}
}
}
#[async_trait]
pub trait ChatEngineBackendPlugin: Send + Sync {
async fn on_session_type_configured(
&self,
_ctx: SessionPluginCtx,
) -> Result<SessionPluginResponse, PluginError> {
Ok(SessionPluginResponse::default())
}
async fn on_session_created(
&self,
_ctx: SessionPluginCtx,
) -> Result<SessionPluginResponse, PluginError> {
Ok(SessionPluginResponse::default())
}
async fn on_session_updated(
&self,
_ctx: SessionPluginCtx,
) -> Result<SessionPluginResponse, PluginError> {
Ok(SessionPluginResponse::default())
}
async fn on_message(&self, _ctx: MessagePluginCtx) -> Result<PluginStream, PluginError> {
Ok(empty_stream())
}
async fn on_message_recreate(
&self,
_ctx: MessagePluginCtx,
) -> Result<PluginStream, PluginError> {
Ok(empty_stream())
}
async fn on_session_summary(
&self,
_ctx: SessionPluginCtx,
) -> Result<PluginStream, PluginError> {
Ok(empty_stream())
}
async fn health_check(&self) -> Result<HealthStatus, PluginError> {
Ok(HealthStatus::Healthy)
}
fn plugin_instance_id(&self) -> &str;
}