use std::sync::Arc;
use async_trait::async_trait;
use nexo_tool_meta::admin::agent_events::AgentEventKind;
use nexo_tool_meta::admin::processing::{PendingInbound, ProcessingControlState, ProcessingScope};
use crate::agent::admin_rpc::channel_outbound::{
ChannelOutboundDispatcher, ChannelOutboundError, OutboundMessage,
};
use crate::agent::admin_rpc::dispatcher::AdminRpcError;
use crate::agent::admin_rpc::transcript_appender::{
TranscriptAppender, TranscriptEntry, TranscriptRole,
};
use crate::agent::agent_events::AgentEventEmitter;
#[async_trait]
pub trait ProcessingControlStore: Send + Sync + std::fmt::Debug {
async fn get(&self, scope: &ProcessingScope) -> anyhow::Result<ProcessingControlState>;
async fn set(
&self,
scope: ProcessingScope,
state: ProcessingControlState,
) -> anyhow::Result<bool>;
async fn clear(&self, scope: &ProcessingScope) -> anyhow::Result<bool>;
async fn push_pending(
&self,
_scope: &ProcessingScope,
_inbound: PendingInbound,
) -> anyhow::Result<(usize, u32)> {
Err(anyhow::anyhow!(
"push_pending not implemented for this store"
))
}
async fn drain_pending(&self, _scope: &ProcessingScope) -> anyhow::Result<Vec<PendingInbound>> {
Ok(Vec::new())
}
async fn pending_depth(&self, _scope: &ProcessingScope) -> anyhow::Result<usize> {
Ok(0)
}
}
pub fn err_not_implemented(method: &str, detail: &str) -> AdminRpcError {
AdminRpcError::MethodNotFound(format!("not_implemented: {method} — {detail}"))
}
pub async fn state(
store: &dyn ProcessingControlStore,
params: serde_json::Value,
) -> crate::agent::admin_rpc::dispatcher::AdminRpcResult {
use nexo_tool_meta::admin::processing::{ProcessingStateParams, ProcessingStateResponse};
let p: ProcessingStateParams = match serde_json::from_value(params) {
Ok(p) => p,
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams(e.to_string()),
)
}
};
match store.get(&p.scope).await {
Ok(state) => crate::agent::admin_rpc::dispatcher::AdminRpcResult::ok(
serde_json::to_value(ProcessingStateResponse { state })
.unwrap_or(serde_json::Value::Null),
),
Err(e) => crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal(format!("processing.state read: {e}")),
),
}
}
pub async fn pause(
store: &dyn ProcessingControlStore,
emitter: Option<&Arc<dyn AgentEventEmitter>>,
params: serde_json::Value,
) -> crate::agent::admin_rpc::dispatcher::AdminRpcResult {
use nexo_tool_meta::admin::processing::{ProcessingAck, ProcessingPauseParams};
let p: ProcessingPauseParams = match serde_json::from_value(params) {
Ok(p) => p,
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams(e.to_string()),
)
}
};
if !p.scope.is_v0_supported() {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(err_not_implemented(
"nexo/admin/processing/pause",
"v0 routes only Conversation scope",
));
}
let prev_state = match store.get(&p.scope).await {
Ok(s) => s,
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal(format!("processing.pause get: {e}")),
)
}
};
if matches!(prev_state, ProcessingControlState::PausedByOperator { .. }) {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::ok(
serde_json::to_value(ProcessingAck {
changed: false,
correlation_id: uuid::Uuid::new_v4(),
transcript_stamped: None,
drained_pending: None,
})
.unwrap_or(serde_json::Value::Null),
);
}
let at_ms = now_epoch_ms();
let new_state = ProcessingControlState::PausedByOperator {
scope: p.scope.clone(),
paused_at_ms: at_ms,
operator_token_hash: p.operator_token_hash,
reason: p.reason,
};
let new_state_for_emit = new_state.clone();
let changed = match store.set(p.scope.clone(), new_state).await {
Ok(c) => c,
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal(format!("processing.pause set: {e}")),
)
}
};
if changed {
if let Some(em) = emitter {
em.emit(AgentEventKind::ProcessingStateChanged {
agent_id: p.scope.agent_id().to_string(),
scope: p.scope.clone(),
prev_state,
new_state: new_state_for_emit,
at_ms,
tenant_id: None,
})
.await;
}
}
crate::agent::admin_rpc::dispatcher::AdminRpcResult::ok(
serde_json::to_value(ProcessingAck {
changed,
correlation_id: uuid::Uuid::new_v4(),
transcript_stamped: None,
drained_pending: None,
})
.unwrap_or(serde_json::Value::Null),
)
}
pub async fn resume(
store: &dyn ProcessingControlStore,
emitter: Option<&Arc<dyn AgentEventEmitter>>,
appender: Option<&dyn TranscriptAppender>,
params: serde_json::Value,
) -> crate::agent::admin_rpc::dispatcher::AdminRpcResult {
use nexo_tool_meta::admin::processing::{
ProcessingAck, ProcessingResumeParams, PROCESSING_SUMMARY_MAX_LEN,
};
let p: ProcessingResumeParams = match serde_json::from_value(params) {
Ok(p) => p,
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams(e.to_string()),
)
}
};
if !p.scope.is_v0_supported() {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(err_not_implemented(
"nexo/admin/processing/resume",
"v0 routes only Conversation scope",
));
}
if let Some(summary) = &p.summary_for_agent {
if p.session_id.is_none() {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams("session_id_required_with_summary".into()),
);
}
let trimmed = summary.trim();
if trimmed.is_empty() {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams("empty_summary".into()),
);
}
if summary.chars().count() > PROCESSING_SUMMARY_MAX_LEN {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams("summary_too_long".into()),
);
}
}
let prev_state = match store.get(&p.scope).await {
Ok(s) => s,
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal(format!("processing.resume get: {e}")),
)
}
};
let at_ms = now_epoch_ms();
let changed = match store.clear(&p.scope).await {
Ok(c) => c,
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal(format!("processing.resume clear: {e}")),
)
}
};
if changed {
if let Some(em) = emitter {
em.emit(AgentEventKind::ProcessingStateChanged {
agent_id: p.scope.agent_id().to_string(),
scope: p.scope.clone(),
prev_state,
new_state: ProcessingControlState::AgentActive,
at_ms,
tenant_id: None,
})
.await;
}
}
let transcript_stamped = match (p.summary_for_agent.as_ref(), p.session_id, appender) {
(Some(summary), Some(session_id), Some(app)) => {
let entry = TranscriptEntry {
role: TranscriptRole::System,
content: format!("[operator_summary] {}", summary.trim()),
source_plugin: "intervention:summary".into(),
sender_id: Some(format!("operator:{}", p.operator_token_hash)),
message_id: None,
};
match app.append(p.scope.agent_id(), session_id, entry).await {
Ok(()) => Some(true),
Err(e) => {
tracing::warn!(
error = %e,
agent = %p.scope.agent_id(),
session_id = %session_id,
"operator summary stamp failed; resume already cleared",
);
Some(false)
}
}
}
(None, _, _) => None,
_ => Some(false),
};
let drained = match store.drain_pending(&p.scope).await {
Ok(v) => v,
Err(e) => {
tracing::warn!(
error = %e,
agent = %p.scope.agent_id(),
"drain_pending failed; resume continues without replay",
);
Vec::new()
}
};
if !drained.is_empty() {
if let (Some(session_id), Some(app)) = (p.session_id, appender) {
for it in &drained {
let entry = TranscriptEntry {
role: TranscriptRole::User,
content: it.body.clone(),
source_plugin: it.source_plugin.clone(),
sender_id: Some(it.from_contact_id.clone()),
message_id: it.message_id,
};
if let Err(e) = app.append(p.scope.agent_id(), session_id, entry).await {
tracing::warn!(
error = %e,
agent = %p.scope.agent_id(),
session_id = %session_id,
"pending inbound stamp failed; replay best-effort",
);
}
}
}
}
let drained_pending = if drained.is_empty() {
None
} else {
Some(drained.len() as u32)
};
crate::agent::admin_rpc::dispatcher::AdminRpcResult::ok(
serde_json::to_value(ProcessingAck {
changed,
correlation_id: uuid::Uuid::new_v4(),
transcript_stamped,
drained_pending,
})
.unwrap_or(serde_json::Value::Null),
)
}
pub async fn intervention(
store: &dyn ProcessingControlStore,
outbound: Option<&dyn ChannelOutboundDispatcher>,
appender: Option<&dyn TranscriptAppender>,
params: serde_json::Value,
) -> crate::agent::admin_rpc::dispatcher::AdminRpcResult {
use nexo_tool_meta::admin::processing::{
InterventionAction, ProcessingAck, ProcessingInterventionParams,
};
let p: ProcessingInterventionParams = match serde_json::from_value(params) {
Ok(p) => p,
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams(e.to_string()),
)
}
};
if !p.scope.is_v0_supported() {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(err_not_implemented(
"nexo/admin/processing/intervention",
"v0 routes only Conversation scope",
));
}
if !p.action.is_v0_supported() {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(err_not_implemented(
"nexo/admin/processing/intervention",
"v0 routes only Reply action",
));
}
match store.get(&p.scope).await {
Ok(ProcessingControlState::PausedByOperator { .. }) => {}
Ok(_) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams(format!("scope_not_paused: {}", p.scope.agent_id())),
)
}
Err(e) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal(format!("processing.intervention read: {e}")),
)
}
}
let (outbound_message_id, reply_body, reply_channel) = match &p.action {
InterventionAction::Reply {
channel,
account_id,
to,
body,
msg_kind,
attachments,
reply_to_msg_id,
} => {
let Some(disp) = outbound else {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal("channel_outbound dispatcher not configured".into()),
);
};
let msg = OutboundMessage {
channel: channel.clone(),
account_id: account_id.clone(),
to: to.clone(),
body: body.clone(),
msg_kind: msg_kind.clone(),
attachments: attachments.clone(),
reply_to_msg_id: reply_to_msg_id.clone(),
};
let omid = match disp.send(msg).await {
Ok(ack) => ack.outbound_message_id,
Err(ChannelOutboundError::ChannelUnavailable(name)) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal(format!("channel_unavailable: {name}")),
)
}
Err(ChannelOutboundError::InvalidParams(msg)) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::InvalidParams(msg),
)
}
Err(ChannelOutboundError::Transport(msg)) => {
return crate::agent::admin_rpc::dispatcher::AdminRpcResult::err(
AdminRpcError::Internal(format!("transport: {msg}")),
)
}
};
(omid, body.clone(), channel.clone())
}
_ => unreachable!("is_v0_supported gate above limited to Reply"),
};
let transcript_stamped = match (p.session_id, appender) {
(Some(session_id), Some(app)) => {
let entry = TranscriptEntry {
role: TranscriptRole::Assistant,
content: reply_body,
source_plugin: format!("intervention:{reply_channel}"),
sender_id: Some(format!("operator:{}", p.operator_token_hash)),
message_id: outbound_message_id
.as_deref()
.and_then(|s| uuid::Uuid::parse_str(s).ok()),
};
match app.append(p.scope.agent_id(), session_id, entry).await {
Ok(()) => Some(true),
Err(e) => {
tracing::warn!(
error = %e,
agent = %p.scope.agent_id(),
session_id = %session_id,
"transcript stamp failed; reply already sent",
);
Some(false)
}
}
}
_ => Some(false),
};
crate::agent::admin_rpc::dispatcher::AdminRpcResult::ok(
serde_json::to_value(ProcessingAck {
changed: true,
correlation_id: uuid::Uuid::new_v4(),
transcript_stamped,
drained_pending: None,
})
.unwrap_or(serde_json::Value::Null),
)
}
fn now_epoch_ms() -> u64 {
use std::time::SystemTime;
SystemTime::now()
.duration_since(SystemTime::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Mutex;
use nexo_tool_meta::admin::processing::{
InterventionAction, ProcessingControlState, ProcessingScope,
};
use crate::agent::admin_rpc::channel_outbound::{
ChannelOutboundError, OutboundAck, OutboundMessage,
};
#[derive(Debug, Default)]
struct CapturingOutbound {
sent: std::sync::Mutex<Vec<OutboundMessage>>,
respond_with_id: Option<String>,
}
#[async_trait]
impl ChannelOutboundDispatcher for CapturingOutbound {
async fn send(&self, msg: OutboundMessage) -> Result<OutboundAck, ChannelOutboundError> {
self.sent.lock().unwrap().push(msg);
Ok(OutboundAck {
outbound_message_id: self.respond_with_id.clone(),
})
}
}
#[derive(Debug, Default)]
struct MockStore {
rows: Mutex<std::collections::HashMap<ProcessingScope, ProcessingControlState>>,
pending: Mutex<
std::collections::HashMap<ProcessingScope, std::collections::VecDeque<PendingInbound>>,
>,
}
#[async_trait]
impl ProcessingControlStore for MockStore {
async fn get(&self, scope: &ProcessingScope) -> anyhow::Result<ProcessingControlState> {
Ok(self
.rows
.lock()
.unwrap()
.get(scope)
.cloned()
.unwrap_or(ProcessingControlState::AgentActive))
}
async fn set(
&self,
scope: ProcessingScope,
state: ProcessingControlState,
) -> anyhow::Result<bool> {
let mut rows = self.rows.lock().unwrap();
let prev = rows.insert(scope, state.clone());
Ok(prev != Some(state))
}
async fn clear(&self, scope: &ProcessingScope) -> anyhow::Result<bool> {
Ok(self.rows.lock().unwrap().remove(scope).is_some())
}
async fn push_pending(
&self,
scope: &ProcessingScope,
inbound: PendingInbound,
) -> anyhow::Result<(usize, u32)> {
let mut p = self.pending.lock().unwrap();
let q = p
.entry(scope.clone())
.or_insert_with(std::collections::VecDeque::new);
q.push_back(inbound);
Ok((q.len(), 0))
}
async fn drain_pending(
&self,
scope: &ProcessingScope,
) -> anyhow::Result<Vec<PendingInbound>> {
Ok(self
.pending
.lock()
.unwrap()
.remove(scope)
.map(|q| q.into_iter().collect())
.unwrap_or_default())
}
}
fn convo() -> ProcessingScope {
ProcessingScope::Conversation {
agent_id: "ana".into(),
channel: "whatsapp".into(),
account_id: "acc".into(),
contact_id: "wa.55".into(),
mcp_channel_source: None,
}
}
#[tokio::test]
async fn pause_and_state_round_trip() {
let store = MockStore::default();
let pause_params = serde_json::json!({
"scope": convo(),
"operator_token_hash": "abcdef0123456789",
"reason": "escalated",
});
let result = pause(&store, None, pause_params).await;
assert!(result.result.is_some(), "pause ok");
let ack: nexo_tool_meta::admin::processing::ProcessingAck =
serde_json::from_value(result.result.unwrap()).unwrap();
assert!(ack.changed, "first pause changes state");
let state_params = serde_json::json!({ "scope": convo() });
let state_result = state(&store, state_params).await;
let resp: nexo_tool_meta::admin::processing::ProcessingStateResponse =
serde_json::from_value(state_result.result.unwrap()).unwrap();
assert!(matches!(
resp.state,
ProcessingControlState::PausedByOperator { .. }
));
}
#[tokio::test]
async fn pause_idempotent_returns_changed_false_on_second_call() {
let store = MockStore::default();
let pause_params = serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
});
let _ = pause(&store, None, pause_params.clone()).await;
let result = pause(&store, None, pause_params).await;
let ack: nexo_tool_meta::admin::processing::ProcessingAck =
serde_json::from_value(result.result.unwrap()).unwrap();
assert!(!ack.changed, "second pause is a no-op");
}
#[tokio::test]
async fn pause_rejects_non_v0_scope() {
let store = MockStore::default();
let result = pause(
&store,
None,
serde_json::json!({
"scope": ProcessingScope::Agent { agent_id: "ana".into() },
"operator_token_hash": "h",
}),
)
.await;
match result.error.expect("error") {
AdminRpcError::MethodNotFound(m) => {
assert!(m.contains("not_implemented"), "got: {m}");
}
other => panic!("expected MethodNotFound, got {other:?}"),
}
}
#[tokio::test]
async fn resume_clears_state_and_intervention_after_resume_is_rejected() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let _ = resume(
&store,
None,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let outbound = CapturingOutbound::default();
let attempt = intervention(
&store,
Some(&outbound),
None,
serde_json::json!({
"scope": convo(),
"action": InterventionAction::Reply {
channel: "whatsapp".into(),
account_id: "acc".into(),
to: "wa.55".into(),
body: "hi".into(),
msg_kind: "text".into(),
attachments: vec![],
reply_to_msg_id: None,
},
"operator_token_hash": "h",
}),
)
.await;
match attempt.error.expect("error") {
AdminRpcError::InvalidParams(m) => {
assert!(m.contains("scope_not_paused"), "got: {m}");
}
other => panic!("expected InvalidParams, got {other:?}"),
}
assert!(outbound.sent.lock().unwrap().is_empty());
}
#[tokio::test]
async fn intervention_dispatches_reply_through_outbound_when_scope_is_paused() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let outbound = CapturingOutbound {
respond_with_id: Some("provider-id-9".into()),
..Default::default()
};
let result = intervention(
&store,
Some(&outbound),
None,
serde_json::json!({
"scope": convo(),
"action": {
"kind": "reply",
"channel": "whatsapp",
"account_id": "acc",
"to": "wa.55",
"body": "hello operator",
"msg_kind": "text",
},
"operator_token_hash": "h",
}),
)
.await;
assert!(result.result.is_some(), "ok: {result:?}");
let sent = outbound.sent.lock().unwrap();
assert_eq!(sent.len(), 1);
assert_eq!(sent[0].channel, "whatsapp");
assert_eq!(sent[0].account_id, "acc");
assert_eq!(sent[0].to, "wa.55");
assert_eq!(sent[0].body, "hello operator");
assert_eq!(sent[0].msg_kind, "text");
}
#[tokio::test]
async fn intervention_returns_internal_when_outbound_missing() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let result = intervention(
&store,
None,
None,
serde_json::json!({
"scope": convo(),
"action": {
"kind": "reply",
"channel": "whatsapp",
"account_id": "acc",
"to": "wa.55",
"body": "hi",
"msg_kind": "text",
},
"operator_token_hash": "h",
}),
)
.await;
match result.error.expect("error") {
AdminRpcError::Internal(m) => {
assert!(m.contains("channel_outbound"), "got: {m}");
}
other => panic!("expected Internal, got {other:?}"),
}
}
#[derive(Debug, Default)]
struct RecordingAppender {
captured: std::sync::Mutex<Vec<(String, uuid::Uuid, super::TranscriptEntry)>>,
fail: bool,
}
#[async_trait]
impl super::TranscriptAppender for RecordingAppender {
async fn append(
&self,
agent_id: &str,
session_id: uuid::Uuid,
entry: super::TranscriptEntry,
) -> anyhow::Result<()> {
if self.fail {
return Err(anyhow::anyhow!("synthetic disk full"));
}
self.captured
.lock()
.unwrap()
.push((agent_id.into(), session_id, entry));
Ok(())
}
}
fn paused_with_session() -> uuid::Uuid {
uuid::Uuid::parse_str("11111111-1111-4111-8111-111111111111").unwrap()
}
#[tokio::test]
async fn intervention_stamps_transcript_when_session_and_appender_both_set() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let outbound = CapturingOutbound {
respond_with_id: Some("550e8400-e29b-41d4-a716-446655440000".into()),
..Default::default()
};
let appender = RecordingAppender::default();
let session_id = paused_with_session();
let result = intervention(
&store,
Some(&outbound),
Some(&appender),
serde_json::json!({
"scope": convo(),
"action": {
"kind": "reply",
"channel": "whatsapp",
"account_id": "acc",
"to": "wa.55",
"body": "ya te resuelvo",
"msg_kind": "text",
},
"operator_token_hash": "tokhash",
"session_id": session_id,
}),
)
.await;
assert!(result.error.is_none(), "{result:?}");
assert_eq!(outbound.sent.lock().unwrap().len(), 1);
let captured = appender.captured.lock().unwrap();
assert_eq!(captured.len(), 1);
let (agent_id, sid, entry) = &captured[0];
assert_eq!(agent_id, "ana");
assert_eq!(*sid, session_id);
assert!(matches!(entry.role, super::TranscriptRole::Assistant));
assert_eq!(entry.content, "ya te resuelvo");
assert_eq!(entry.source_plugin, "intervention:whatsapp");
assert_eq!(entry.sender_id.as_deref(), Some("operator:tokhash"));
assert!(entry.message_id.is_some(), "outbound_message_id threaded");
let v = result.result.unwrap();
assert_eq!(v["transcript_stamped"], true);
}
#[tokio::test]
async fn intervention_skips_stamp_when_session_id_absent() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let outbound = CapturingOutbound::default();
let appender = RecordingAppender::default();
let result = intervention(
&store,
Some(&outbound),
Some(&appender),
serde_json::json!({
"scope": convo(),
"action": {
"kind": "reply",
"channel": "whatsapp",
"account_id": "acc",
"to": "wa.55",
"body": "hi",
"msg_kind": "text",
},
"operator_token_hash": "h",
}),
)
.await;
assert_eq!(outbound.sent.lock().unwrap().len(), 1);
assert!(appender.captured.lock().unwrap().is_empty());
let v = result.result.unwrap();
assert_eq!(v["transcript_stamped"], false);
}
#[tokio::test]
async fn intervention_skips_stamp_when_appender_unwired() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let outbound = CapturingOutbound::default();
let result = intervention(
&store,
Some(&outbound),
None, serde_json::json!({
"scope": convo(),
"action": {
"kind": "reply",
"channel": "whatsapp",
"account_id": "acc",
"to": "wa.55",
"body": "hi",
"msg_kind": "text",
},
"operator_token_hash": "h",
"session_id": paused_with_session(),
}),
)
.await;
assert_eq!(outbound.sent.lock().unwrap().len(), 1);
let v = result.result.unwrap();
assert_eq!(v["transcript_stamped"], false);
}
#[tokio::test]
async fn intervention_degrades_when_appender_returns_err() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let outbound = CapturingOutbound::default();
let appender = RecordingAppender {
fail: true,
..Default::default()
};
let result = intervention(
&store,
Some(&outbound),
Some(&appender),
serde_json::json!({
"scope": convo(),
"action": {
"kind": "reply",
"channel": "whatsapp",
"account_id": "acc",
"to": "wa.55",
"body": "hi",
"msg_kind": "text",
},
"operator_token_hash": "h",
"session_id": paused_with_session(),
}),
)
.await;
assert!(result.error.is_none(), "{result:?}");
assert_eq!(outbound.sent.lock().unwrap().len(), 1);
let v = result.result.unwrap();
assert_eq!(v["transcript_stamped"], false);
}
#[tokio::test]
async fn intervention_stamp_omits_message_id_when_outbound_returns_none() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let outbound = CapturingOutbound {
respond_with_id: None, ..Default::default()
};
let appender = RecordingAppender::default();
let result = intervention(
&store,
Some(&outbound),
Some(&appender),
serde_json::json!({
"scope": convo(),
"action": {
"kind": "reply",
"channel": "telegram",
"account_id": "tg.bot",
"to": "tg.55",
"body": "hola",
"msg_kind": "text",
},
"operator_token_hash": "h",
"session_id": paused_with_session(),
}),
)
.await;
assert!(result.error.is_none());
let captured = appender.captured.lock().unwrap();
assert_eq!(captured.len(), 1);
let entry = &captured[0].2;
assert!(entry.message_id.is_none());
assert_eq!(entry.source_plugin, "intervention:telegram");
}
#[tokio::test]
async fn intervention_rejects_non_v0_action() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let outbound = CapturingOutbound::default();
let result = intervention(
&store,
Some(&outbound),
None,
serde_json::json!({
"scope": convo(),
"action": {
"kind": "skip_item",
"item_id": "x",
"reason": "y",
},
"operator_token_hash": "h",
}),
)
.await;
match result.error.expect("error") {
AdminRpcError::MethodNotFound(m) => {
assert!(m.contains("not_implemented"));
}
other => panic!("expected MethodNotFound, got {other:?}"),
}
}
#[tokio::test]
async fn resume_injects_summary_as_system_entry_with_prefix() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let appender = RecordingAppender::default();
let session_id = paused_with_session();
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "tokhash",
"session_id": session_id,
"summary_for_agent": " cliente confirmó dirección ",
}),
)
.await;
assert!(result.error.is_none(), "{result:?}");
let captured = appender.captured.lock().unwrap();
assert_eq!(captured.len(), 1);
let (agent_id, sid, entry) = &captured[0];
assert_eq!(agent_id, "ana");
assert_eq!(*sid, session_id);
assert!(matches!(entry.role, super::TranscriptRole::System));
assert_eq!(
entry.content,
"[operator_summary] cliente confirmó dirección"
);
assert_eq!(entry.source_plugin, "intervention:summary");
assert_eq!(entry.sender_id.as_deref(), Some("operator:tokhash"));
let v = result.result.unwrap();
assert_eq!(v["transcript_stamped"], true);
assert_eq!(v["changed"], true);
}
#[tokio::test]
async fn resume_rejects_summary_without_session_id() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let appender = RecordingAppender::default();
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"summary_for_agent": "ok",
}),
)
.await;
match result.error.expect("error") {
AdminRpcError::InvalidParams(m) => {
assert!(m.contains("session_id_required_with_summary"), "got: {m}");
}
other => panic!("expected InvalidParams, got {other:?}"),
}
assert!(matches!(
store.get(&convo()).await,
Ok(ProcessingControlState::PausedByOperator { .. })
));
assert!(appender.captured.lock().unwrap().is_empty());
}
#[tokio::test]
async fn resume_rejects_empty_summary() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let result = resume(
&store,
None,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": paused_with_session(),
"summary_for_agent": " \t \n",
}),
)
.await;
match result.error.expect("error") {
AdminRpcError::InvalidParams(m) => assert!(m.contains("empty_summary")),
other => panic!("expected InvalidParams, got {other:?}"),
}
}
#[tokio::test]
async fn resume_rejects_summary_over_4096_chars() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let huge = "a".repeat(4097);
let result = resume(
&store,
None,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": paused_with_session(),
"summary_for_agent": huge,
}),
)
.await;
match result.error.expect("error") {
AdminRpcError::InvalidParams(m) => assert!(m.contains("summary_too_long")),
other => panic!("expected InvalidParams, got {other:?}"),
}
}
#[tokio::test]
async fn resume_without_summary_skips_injection_and_returns_none() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let appender = RecordingAppender::default();
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
assert!(result.error.is_none());
assert!(appender.captured.lock().unwrap().is_empty());
let v = result.result.unwrap();
assert!(v.get("transcript_stamped").is_none());
assert_eq!(v["changed"], true);
}
#[tokio::test]
async fn resume_logs_and_proceeds_when_appender_errs() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let appender = RecordingAppender {
fail: true,
..Default::default()
};
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": paused_with_session(),
"summary_for_agent": "anything",
}),
)
.await;
assert!(result.error.is_none(), "{result:?}");
let v = result.result.unwrap();
assert_eq!(v["transcript_stamped"], false);
assert_eq!(v["changed"], true);
assert!(matches!(
store.get(&convo()).await,
Ok(ProcessingControlState::AgentActive)
));
}
fn pending(seq: u64) -> PendingInbound {
PendingInbound {
message_id: None,
from_contact_id: format!("wa.55+{seq}"),
body: format!("msg {seq}"),
timestamp_ms: 1_700_000_000_000 + seq,
source_plugin: "whatsapp".into(),
}
}
#[tokio::test]
async fn resume_drains_pending_into_transcript_user_entries() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
for i in 0..3 {
store.push_pending(&convo(), pending(i)).await.unwrap();
}
let appender = RecordingAppender::default();
let session_id = paused_with_session();
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": session_id,
}),
)
.await;
assert!(result.error.is_none(), "{result:?}");
let v = result.result.unwrap();
assert_eq!(v["drained_pending"], 3);
let captured = appender.captured.lock().unwrap();
assert_eq!(captured.len(), 3);
for (i, (_, sid, entry)) in captured.iter().enumerate() {
assert_eq!(*sid, session_id);
assert!(matches!(entry.role, super::TranscriptRole::User));
assert_eq!(entry.content, format!("msg {i}"));
assert_eq!(entry.source_plugin, "whatsapp");
assert_eq!(
entry.sender_id.as_deref(),
Some(format!("wa.55+{i}").as_str())
);
}
}
#[tokio::test]
async fn resume_with_empty_queue_returns_drained_pending_none() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let appender = RecordingAppender::default();
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": paused_with_session(),
}),
)
.await;
assert!(result.error.is_none());
let v = result.result.unwrap();
assert!(v.get("drained_pending").is_none(), "no field when empty");
assert!(appender.captured.lock().unwrap().is_empty());
}
#[tokio::test]
async fn resume_drains_only_target_scope() {
let store = MockStore::default();
let other = ProcessingScope::Conversation {
agent_id: "ana".into(),
channel: "whatsapp".into(),
account_id: "acc".into(),
contact_id: "wa.99".into(),
mcp_channel_source: None,
};
let _ = pause(
&store,
None,
serde_json::json!({ "scope": convo(), "operator_token_hash": "h" }),
)
.await;
let _ = pause(
&store,
None,
serde_json::json!({ "scope": other, "operator_token_hash": "h" }),
)
.await;
store.push_pending(&convo(), pending(0)).await.unwrap();
store.push_pending(&convo(), pending(1)).await.unwrap();
let other_scope = ProcessingScope::Conversation {
agent_id: "ana".into(),
channel: "whatsapp".into(),
account_id: "acc".into(),
contact_id: "wa.99".into(),
mcp_channel_source: None,
};
store.push_pending(&other_scope, pending(99)).await.unwrap();
let appender = RecordingAppender::default();
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": paused_with_session(),
}),
)
.await;
assert!(result.error.is_none());
let v = result.result.unwrap();
assert_eq!(v["drained_pending"], 2);
let other_remaining = store.drain_pending(&other_scope).await.unwrap();
assert_eq!(other_remaining.len(), 1);
}
#[tokio::test]
async fn resume_drains_but_skips_stamp_when_session_id_missing() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({ "scope": convo(), "operator_token_hash": "h" }),
)
.await;
store.push_pending(&convo(), pending(0)).await.unwrap();
store.push_pending(&convo(), pending(1)).await.unwrap();
let appender = RecordingAppender::default();
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
assert!(result.error.is_none());
let v = result.result.unwrap();
assert_eq!(v["drained_pending"], 2);
assert!(appender.captured.lock().unwrap().is_empty());
assert!(store.drain_pending(&convo()).await.unwrap().is_empty());
}
#[tokio::test]
async fn resume_drains_alongside_summary_injection() {
let store = MockStore::default();
let _ = pause(
&store,
None,
serde_json::json!({ "scope": convo(), "operator_token_hash": "h" }),
)
.await;
store.push_pending(&convo(), pending(0)).await.unwrap();
let appender = RecordingAppender::default();
let session_id = paused_with_session();
let result = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": session_id,
"summary_for_agent": "cliente confirmó dirección",
}),
)
.await;
assert!(result.error.is_none());
let v = result.result.unwrap();
assert_eq!(v["transcript_stamped"], true);
assert_eq!(v["drained_pending"], 1);
let captured = appender.captured.lock().unwrap();
assert_eq!(captured.len(), 2);
assert!(matches!(captured[0].2.role, super::TranscriptRole::System));
assert!(captured[0].2.content.starts_with("[operator_summary]"));
assert!(matches!(captured[1].2.role, super::TranscriptRole::User));
assert_eq!(captured[1].2.content, "msg 0");
}
#[tokio::test]
async fn pause_resume_cycles_dont_drag_along_pending_entries() {
let store = MockStore::default();
let appender = RecordingAppender::default();
let session_id = paused_with_session();
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
store.push_pending(&convo(), pending(0)).await.unwrap();
store.push_pending(&convo(), pending(1)).await.unwrap();
let r1 = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": session_id,
}),
)
.await;
assert!(r1.error.is_none());
assert_eq!(r1.result.unwrap()["drained_pending"], 2);
assert_eq!(appender.captured.lock().unwrap().len(), 2);
assert!(store.drain_pending(&convo()).await.unwrap().is_empty());
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
store.push_pending(&convo(), pending(99)).await.unwrap();
let r2 = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": session_id,
}),
)
.await;
assert!(r2.error.is_none());
let v2 = r2.result.unwrap();
assert_eq!(v2["drained_pending"], 1, "no drag-along from cycle 1");
let captured = appender.captured.lock().unwrap();
assert_eq!(captured.len(), 3);
assert_eq!(captured[2].2.content, "msg 99");
drop(captured);
let _ = pause(
&store,
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let r3 = resume(
&store,
None,
Some(&appender),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
"session_id": session_id,
}),
)
.await;
assert!(r3.error.is_none());
let v3 = r3.result.unwrap();
assert!(
v3.get("drained_pending").is_none(),
"empty drain MUST omit the field, got: {v3}",
);
assert_eq!(appender.captured.lock().unwrap().len(), 3);
assert!(matches!(
store.get(&convo()).await.unwrap(),
ProcessingControlState::AgentActive
));
assert!(store.drain_pending(&convo()).await.unwrap().is_empty());
}
#[derive(Debug, Default)]
struct CapturingEmitter {
events: std::sync::Mutex<Vec<AgentEventKind>>,
}
#[async_trait]
impl crate::agent::agent_events::AgentEventEmitter for CapturingEmitter {
async fn emit(&self, event: AgentEventKind) {
self.events.lock().unwrap().push(event);
}
}
#[tokio::test]
async fn pause_emits_processing_state_changed_with_prev_active() {
let store = MockStore::default();
let captured: Arc<CapturingEmitter> = Arc::new(CapturingEmitter::default());
let emitter: Arc<dyn crate::agent::agent_events::AgentEventEmitter> = captured.clone();
let _ = pause(
&store,
Some(&emitter),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "abcdef0123456789",
"reason": "investigando",
}),
)
.await;
let evs = captured.events.lock().unwrap();
assert_eq!(evs.len(), 1, "first pause must emit once");
match &evs[0] {
AgentEventKind::ProcessingStateChanged {
agent_id,
prev_state,
new_state,
..
} => {
assert_eq!(agent_id, "ana");
assert!(matches!(prev_state, ProcessingControlState::AgentActive));
assert!(matches!(
new_state,
ProcessingControlState::PausedByOperator { .. }
));
}
other => panic!("unexpected emit: {other:?}"),
}
}
#[tokio::test]
async fn pause_skips_emit_when_already_paused() {
let store = MockStore::default();
let captured: Arc<CapturingEmitter> = Arc::new(CapturingEmitter::default());
let emitter: Arc<dyn crate::agent::agent_events::AgentEventEmitter> = captured.clone();
let params = serde_json::json!({
"scope": convo(),
"operator_token_hash": "abcdef0123456789",
});
let _ = pause(&store, Some(&emitter), params.clone()).await;
let _ = pause(&store, Some(&emitter), params).await;
let evs = captured.events.lock().unwrap();
assert_eq!(
evs.len(),
1,
"idempotent re-pause must not double-emit (subscribers would see phantom)",
);
}
#[tokio::test]
async fn resume_emits_processing_state_changed_with_prev_paused() {
let store = MockStore::default();
let captured: Arc<CapturingEmitter> = Arc::new(CapturingEmitter::default());
let emitter: Arc<dyn crate::agent::agent_events::AgentEventEmitter> = captured.clone();
let _ = pause(
&store,
Some(&emitter),
serde_json::json!({
"scope": convo(),
"operator_token_hash": "abcdef0123456789",
}),
)
.await;
let _ = resume(
&store,
Some(&emitter),
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "abcdef0123456789",
}),
)
.await;
let evs = captured.events.lock().unwrap();
assert_eq!(evs.len(), 2, "pause + resume each emit once");
match &evs[1] {
AgentEventKind::ProcessingStateChanged {
prev_state,
new_state,
..
} => {
assert!(matches!(
prev_state,
ProcessingControlState::PausedByOperator { .. }
));
assert!(matches!(new_state, ProcessingControlState::AgentActive));
}
other => panic!("unexpected emit: {other:?}"),
}
}
#[tokio::test]
async fn resume_skips_emit_when_already_active() {
let store = MockStore::default();
let captured: Arc<CapturingEmitter> = Arc::new(CapturingEmitter::default());
let emitter: Arc<dyn crate::agent::agent_events::AgentEventEmitter> = captured.clone();
let _ = resume(
&store,
Some(&emitter),
None,
serde_json::json!({
"scope": convo(),
"operator_token_hash": "h",
}),
)
.await;
let evs = captured.events.lock().unwrap();
assert_eq!(
evs.len(),
0,
"no-op resume must not emit (subscribers would see phantom transitions)",
);
}
}