use std::sync::atomic::Ordering;
use std::sync::Arc;
use super::super::*;
use super::hub_ext::*;
use super::pipeline::*;
use super::queue::*;
use super::quote::*;
use crate::hub::{AdminConfig, InMemoryQueue};
use crate::ilink::types::{MessageItem, TextItem, WeixinMessage};
use crate::store::Store;
async fn make_state_with_client() -> (Arc<HubState>, String) {
let (state, vtoken, _mock) =
make_state_with_client_and_mock(crate::hub::tests::MockUpstream::returning_ok()).await;
(state, vtoken)
}
async fn make_state_with_client_and_mock(
upstream: Arc<dyn crate::ilink::UpstreamSink>,
) -> (Arc<HubState>, String, Arc<dyn crate::ilink::UpstreamSink>) {
let store = Store::connect("sqlite::memory:")
.await
.expect("in-memory store");
let mock_ref = Arc::clone(&upstream);
let queue = Arc::new(InMemoryQueue::new());
let (_tx, shutdown_rx) = tokio::sync::watch::channel(false);
let state = HubState::new(
upstream,
Arc::new(store),
queue,
shutdown_rx,
"test-relay-secret".to_string(),
AdminConfig::from_env(),
);
let (_, vtoken, _) =
state
.clients
.registry
.write()
.await
.register("test-backend".to_string(), None, None);
(state, vtoken, mock_ref)
}
fn make_quote_msg(
from_user_id: &str,
ref_create_time_ms: i64,
quoted_text: Option<&str>,
) -> WeixinMessage {
let mut mi_obj = serde_json::Map::new();
mi_obj.insert(
"create_time_ms".to_string(),
serde_json::Value::Number(ref_create_time_ms.into()),
);
if let Some(t) = quoted_text {
mi_obj.insert("text_item".to_string(), serde_json::json!({"text": t}));
}
let extra = serde_json::json!({
"ref_msg": {
"message_item": serde_json::Value::Object(mi_obj)
}
});
let item = MessageItem {
item_type: Some(1),
text_item: Some(TextItem {
text: Some("this is the follow-up reply".to_string()),
}),
extra,
..Default::default()
};
WeixinMessage {
from_user_id: Some(from_user_id.to_string()),
item_list: Some(Arc::new(vec![item])),
..Default::default()
}
}
fn make_quote_msg_with_id(from_user_id: &str, ref_msg_id: i64) -> WeixinMessage {
let extra = serde_json::json!({
"ref_msg": {
"message_item": {
"create_time_ms": 1_000_000i64,
"msg_id": ref_msg_id.to_string(),
}
}
});
let item = MessageItem {
item_type: Some(1),
text_item: Some(TextItem {
text: Some("this is the follow-up reply".to_string()),
}),
extra,
..Default::default()
};
WeixinMessage {
from_user_id: Some(from_user_id.to_string()),
item_list: Some(Arc::new(vec![item])),
..Default::default()
}
}
#[tokio::test]
async fn quote_reply_l0_msg_id_exact_routing() {
let (state, vtoken) = make_state_with_client().await;
let peer_user_id = "peer:user1";
let session_name = "at-20260713-090000000";
let ilink_msg_id = 1_783_912_191_000_000_123i64;
state
.store
.save_message_with_msg_id(
"vctx-test",
Some(&vtoken),
session_name,
peer_user_id,
"assistant",
"exact-match reply",
Some(ilink_msg_id),
)
.await
.expect("save message with msg_id");
let msg = make_quote_msg_with_id("user1", ilink_msg_id);
let result = resolve_quote_from_msg_id(&state, &msg).await;
assert!(result.is_some(), "L0 msg_id lookup must find the exact row");
match result.unwrap() {
QuoteOrigin::Client {
session_name: sn,
vtoken: vt,
..
} => {
assert_eq!(sn, Some(session_name.to_string()));
assert_eq!(vt, vtoken);
}
other => panic!("expected QuoteOrigin::Client, got {other:?}"),
}
}
#[tokio::test]
async fn quote_reply_l0_msg_id_miss_returns_none() {
let (state, _vtoken) = make_state_with_client().await;
let msg = make_quote_msg_with_id("user1", 9_000_000_000_000_000_000);
let result = resolve_quote_from_msg_id(&state, &msg).await;
assert!(result.is_none(), "unknown msg_id must not resolve via L0");
}
#[tokio::test]
async fn quote_reply_l0_msg_id_scoped_to_peer() {
let (state, vtoken) = make_state_with_client().await;
let ilink_msg_id = 1_783_912_191_000_000_456i64;
state
.store
.save_message_with_msg_id(
"vctx-other",
Some(&vtoken),
"at-20260713-090000000",
"peer:other-user",
"assistant",
"other-peer reply",
Some(ilink_msg_id),
)
.await
.expect("save message with msg_id");
let msg = make_quote_msg_with_id("user1", ilink_msg_id);
let result = resolve_quote_from_msg_id(&state, &msg).await;
assert!(result.is_none(), "L0 must not match across peer scopes");
}
#[tokio::test]
async fn at_mention_quote_reply_l1_timestamp_routing() {
let (state, vtoken) = make_state_with_client().await;
let peer_user_id = "peer:user1";
let session_name = "at-20260704-103000000";
state
.store
.save_message(
"vctx-test",
Some(&vtoken),
session_name,
peer_user_id,
"assistant",
"Hello from @mention session",
)
.await
.expect("save message");
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as i64;
let msg = make_quote_msg("user1", now_ms, None);
let result = resolve_quote_from_timestamp(&state, &msg).await;
assert!(
result.is_some(),
"L1 timestamp lookup must find the at-mention session"
);
match result.unwrap() {
QuoteOrigin::Client {
session_name: sn,
vtoken: vt,
..
} => {
assert_eq!(sn, Some(session_name.to_string()));
assert_eq!(vt, vtoken);
}
other => panic!("expected QuoteOrigin::Client, got {other:?}"),
}
}
#[tokio::test]
async fn at_mention_quote_reply_l2_content_routing() {
let (state, vtoken) = make_state_with_client().await;
let peer_user_id = "peer:user1";
let session_name = "at-20260704-103000000";
let content = "Hello from @mention session — this content prefix will match the quote";
state
.store
.save_message(
"vctx-test",
Some(&vtoken),
session_name,
peer_user_id,
"assistant",
content,
)
.await
.expect("save message");
let old_ms = 1_000_000i64;
let msg = make_quote_msg("user1", old_ms, Some(content));
let result = resolve_quote_from_db(&state, &msg).await;
assert!(
result.is_some(),
"L2 content lookup must find the at-mention session"
);
match result.unwrap() {
QuoteOrigin::Client {
session_name: sn, ..
} => assert_eq!(sn, Some(session_name.to_string())),
other => panic!("expected QuoteOrigin::Client, got {other:?}"),
}
}
#[tokio::test]
async fn at_mention_quote_reply_l3_footer_routing() {
let (state, _vtoken) = make_state_with_client().await;
let (_, ilink_claude_vtoken, _) =
state
.clients
.registry
.write()
.await
.register("ilink-claude".to_string(), None, None);
let session_name = "at-20260704-103000";
let quoted_text = format!("Some reply text\n\n---\nilink-claude · {session_name}");
let old_ms = 1_000_000i64;
let msg = make_quote_msg("user1", old_ms, Some("ed_text));
let result = resolve_quote_from_footer(&state, &msg).await;
assert!(
result.is_some(),
"L3 footer lookup must find the ilink-claude client"
);
match result.unwrap() {
QuoteOrigin::Client {
vtoken: vt,
session_name: sn,
..
} => {
assert_eq!(
vt, ilink_claude_vtoken,
"vtoken must match the registered ilink-claude client"
);
assert_eq!(
sn,
Some(session_name.to_string()),
"session_name must be parsed from the footer"
);
}
other => panic!("expected QuoteOrigin::Client, got {other:?}"),
}
}
#[tokio::test]
async fn like_injection_protection_in_content_lookup() {
let (state, vtoken) = make_state_with_client().await;
let peer_user_id = "peer:user1";
state
.store
.save_message(
"vctx-test",
Some(&vtoken),
"session-percent",
peer_user_id,
"assistant",
"100%bonus plan",
)
.await
.expect("save message: percent row");
state
.store
.save_message(
"vctx-test",
Some(&vtoken),
"session-x",
peer_user_id,
"assistant",
"100xplan",
)
.await
.expect("save message: x row");
let msg = make_quote_msg("user1", 1_000_000, Some("100%bonus plan"));
let result = resolve_quote_from_db(&state, &msg).await;
assert!(
result.is_some(),
"LIKE lookup must find the percent-content row"
);
match result.unwrap() {
QuoteOrigin::Client {
session_name: sn, ..
} => assert_eq!(
sn,
Some("session-percent".to_string()),
"must return the percent row, not the wildcard-hit x row"
),
other => panic!("expected QuoteOrigin::Client, got {other:?}"),
}
}
#[tokio::test]
async fn at_mention_quote_reply_l3_footer_persona_routing() {
let (state, _vtoken) = make_state_with_client().await;
let (_, ilink_claude_vtoken, _) =
state
.clients
.registry
.write()
.await
.register("ilink-claude".to_string(), None, None);
let vctx = state
.store
.find_or_create_vctx("user1", None, "ctx-at1")
.await
.expect("find_or_create_vctx");
state
.store
.set_backend_session(
&vctx,
&ilink_claude_vtoken,
"at-20260704-103000",
"cli-uuid",
)
.await
.expect("set_backend_session");
let quoted_text = "body\n\n---\nat-20260704-103000";
let msg = make_quote_msg("user1", 1_000_000, Some(quoted_text));
let result = resolve_quote_from_footer(&state, &msg).await;
assert!(
result.is_some(),
"persona-mode footer slow path must resolve the session to a client"
);
match result.unwrap() {
QuoteOrigin::Client {
vtoken: vt,
session_name: sn,
..
} => {
assert_eq!(
vt, ilink_claude_vtoken,
"vtoken must match the registered ilink-claude client"
);
assert_eq!(
sn,
Some("at-20260704-103000".to_string()),
"session_name must be the at-key from the footer"
);
}
other => panic!("expected QuoteOrigin::Client, got {other:?}"),
}
}
#[tokio::test]
async fn at_mention_quote_reply_l3_footer_unregistered_name_returns_none() {
let (state, _vtoken) = make_state_with_client().await;
let quoted_text = "body\n\n---\nghost-client · at-20260704-103000";
let msg = make_quote_msg("user1", 1_000_000, Some(quoted_text));
let result = resolve_quote_from_footer(&state, &msg).await;
assert!(
result.is_none(),
"unregistered backend name must yield None, not panic"
);
}
#[tokio::test]
async fn quote_reply_l3_full_cascade_fallback() {
let (state, _vtoken) = make_state_with_client().await;
let (_, ilink_claude_vtoken, _) =
state
.clients
.registry
.write()
.await
.register("ilink-claude".to_string(), None, None);
let old_ms = 1_000_000i64;
let quoted_text = "body\n\n---\nilink-claude · at-20260704-103000";
let msg = make_quote_msg("user1", old_ms, Some(quoted_text));
let from_ts = resolve_quote_from_timestamp(&state, &msg).await;
assert!(from_ts.is_none(), "L1 must miss on empty DB");
let from_db = resolve_quote_from_db(&state, &msg).await;
assert!(from_db.is_none(), "L2 must miss on empty DB");
let from_footer = resolve_quote_from_footer(&state, &msg).await;
assert!(
from_footer.is_some(),
"L3 footer lookup must succeed after L1 and L2 both miss"
);
match from_footer.unwrap() {
QuoteOrigin::Client {
vtoken: vt,
name: n,
session_name: sn,
..
} => {
assert_eq!(vt, ilink_claude_vtoken, "vtoken must match ilink-claude");
assert_eq!(n, "ilink-claude", "name must be ilink-claude");
assert_eq!(
sn,
Some("at-20260704-103000".to_string()),
"session_name must be parsed from footer"
);
}
other => panic!("expected QuoteOrigin::Client, got {other:?}"),
}
}
#[tokio::test]
async fn at_mention_quote_reply_l3_footer_with_label_routing() {
let (state, _vtoken) = make_state_with_client().await;
let (_, ilink_claude_vtoken, _) =
state
.clients
.registry
.write()
.await
.register("ilink-claude".to_string(), None, None);
let quoted_text = "body\n\n---\nilink-claude · office · at-20260704-103000";
let msg = make_quote_msg("user1", 1_000_000, Some(quoted_text));
let result = resolve_quote_from_footer(&state, &msg).await;
assert!(
result.is_some(),
"three-part footer must resolve the named client"
);
match result.unwrap() {
QuoteOrigin::Client {
vtoken: vt,
name: n,
session_name: sn,
..
} => {
assert_eq!(
vt, ilink_claude_vtoken,
"vtoken must match ilink-claude (first footer segment)"
);
assert_eq!(n, "ilink-claude", "name must be the first footer segment");
assert_eq!(
sn,
Some("at-20260704-103000".to_string()),
"session_name must be the last footer segment"
);
}
other => panic!("expected QuoteOrigin::Client, got {other:?}"),
}
}
#[test]
fn build_no_backend_reply_non_command_returns_no_backend_online() {
let result = build_no_backend_reply(Some("hello, tell me about Rust"));
assert_eq!(
result,
messages::NO_BACKEND_ONLINE,
"non-command text must produce NO_BACKEND_ONLINE"
);
}
#[test]
fn build_no_backend_reply_none_returns_no_backend_online() {
let result = build_no_backend_reply(None);
assert_eq!(result, messages::NO_BACKEND_ONLINE);
}
#[test]
fn build_no_backend_reply_command_returns_unrecognized_command() {
let result = build_no_backend_reply(Some("/unknown_cmd"));
assert_eq!(
result,
messages::UNRECOGNIZED_COMMAND,
"slash command without a backend must return UNRECOGNIZED_COMMAND"
);
}
#[tokio::test]
async fn push_to_queue_pub_increments_dispatched_metric() {
let queue: Arc<dyn MessageQueue> = Arc::new(InMemoryQueue::new());
let metrics = Metrics::default();
let msg = WeixinMessage::default();
let before = metrics.messages_dispatched.load(Ordering::Relaxed);
push_to_queue_pub(&queue, &metrics, "vhub_test", msg).await;
assert_eq!(
metrics.messages_dispatched.load(Ordering::Relaxed),
before + 1,
"push_to_queue_pub must increment messages_dispatched by 1"
);
}
#[tokio::test]
async fn resolve_quote_from_footer_session_prefix_uses_session_path() {
let (state, vtoken) = make_state_with_client().await;
let session_key = "session-alpha";
let vctx = state
.store
.find_or_create_vctx("user1", None, "real-ctx-session-prefix")
.await
.expect("vctx");
state
.store
.set_backend_session(&vctx, &vtoken, session_key, "cli-sid")
.await
.expect("backend session");
let quoted_text = "body\n\n---\nsession-alpha · at-decoy";
let msg = make_quote_msg("user1", 1_000_000, Some(quoted_text));
let result = resolve_quote_from_footer(&state, &msg).await;
assert!(
result.is_some(),
"session- prefix alone must resolve via slow path (|| not &&)"
);
match result.unwrap() {
QuoteOrigin::Client {
vtoken: vt,
session_name,
..
} => {
assert_eq!(vt, vtoken);
assert_eq!(session_name.as_deref(), Some(session_key));
}
other => panic!("expected Client origin, got {other:?}"),
}
}
#[tokio::test]
async fn build_hub_ext_for_vctx_session_override_returns_non_empty_session_id() {
let (state, vtoken) = make_state_with_client().await;
let vctx = "vctx-test-hub-ext";
let session_name = "my-session";
let session_value = "claude-session-abc123";
state
.store
.set_backend_session(vctx, &vtoken, session_name, session_value)
.await
.expect("set_backend_session must succeed");
let hub_ext =
build_hub_ext_for_vctx(&state.store, vctx, &vtoken, Some(session_name.to_string())).await;
assert!(
hub_ext.is_some(),
"hub_ext must be Some when session override is provided"
);
let ext = hub_ext.unwrap();
assert_eq!(
ext.session_id.as_deref(),
Some(session_value),
"session_id must equal the stored session value (not-empty guard must work)"
);
}
#[tokio::test]
async fn dispatch_message_empty_context_skips_forward() {
let mock = crate::hub::tests::MockUpstream::returning_ok();
let (state, vtoken, mock_ref) = make_state_with_client_and_mock(mock).await;
{
let mut router = state.routing.router.lock().await;
router.set_route("user-empty-ctx", vtoken);
}
let before_dropped = state.metrics.messages_dropped.load(Ordering::Relaxed);
let before_dispatched = state.metrics.messages_dispatched.load(Ordering::Relaxed);
let msg = WeixinMessage {
context_token: Some(String::new()),
from_user_id: Some("user-empty-ctx".into()),
item_list: Some(Arc::new(vec![MessageItem {
item_type: Some(1),
text_item: Some(TextItem {
text: Some("hello".into()),
}),
..Default::default()
}])),
..Default::default()
};
dispatch_message(Arc::clone(&state), msg).await;
assert_eq!(
mock_ref.polls_ok(),
0,
"empty context must not call upstream send_message"
);
assert_eq!(
state.metrics.messages_dispatched.load(Ordering::Relaxed),
before_dispatched,
"empty context must not push to queue"
);
let _ = before_dropped;
}
#[tokio::test]
async fn dispatch_broadcast_empty_context_skips_no_backend_reply() {
let mock = crate::hub::tests::MockUpstream::returning_ok();
let store = Store::connect("sqlite::memory:")
.await
.expect("in-memory store");
let mock_ref = Arc::clone(&mock);
let queue = Arc::new(InMemoryQueue::new());
let (_tx, shutdown_rx) = tokio::sync::watch::channel(false);
let state = HubState::new(
mock,
Arc::new(store),
queue,
shutdown_rx,
"test-relay-secret".to_string(),
AdminConfig::from_env(),
);
let msg = WeixinMessage {
context_token: Some(String::new()),
from_user_id: Some("user-bcast".into()),
item_list: Some(Arc::new(vec![MessageItem {
item_type: Some(1),
text_item: Some(TextItem {
text: Some("hello world".into()),
}),
..Default::default()
}])),
..Default::default()
};
handle_broadcast(Arc::clone(&state), msg).await;
assert_eq!(
mock_ref.polls_ok(),
0,
"empty context on no-backend path must not send reply"
);
}
#[tokio::test]
async fn dispatch_broadcast_no_backend_sends_reply_with_context() {
let mock = crate::hub::tests::MockUpstream::returning_ok();
let store = Store::connect("sqlite::memory:")
.await
.expect("in-memory store");
let mock_ref = Arc::clone(&mock);
let queue = Arc::new(InMemoryQueue::new());
let (_tx, shutdown_rx) = tokio::sync::watch::channel(false);
let state = HubState::new(
mock,
Arc::new(store),
queue,
shutdown_rx,
"test-relay-secret".to_string(),
AdminConfig::from_env(),
);
let msg = WeixinMessage {
context_token: Some("real-ctx-no-backend".into()),
from_user_id: Some("user-nobackend".into()),
item_list: Some(Arc::new(vec![MessageItem {
item_type: Some(1),
text_item: Some(TextItem {
text: Some("hello world".into()),
}),
..Default::default()
}])),
..Default::default()
};
handle_broadcast(Arc::clone(&state), msg).await;
assert_eq!(
mock_ref.polls_ok(),
1,
"no-backend path with context must send fallback reply"
);
}
#[tokio::test]
async fn dispatch_broadcast_online_empty_context_skips_queue() {
let (state, vtoken) = make_state_with_client().await;
state.clients.registry.write().await.mark_online(&vtoken);
let before = state.metrics.messages_dispatched.load(Ordering::Relaxed);
let msg = WeixinMessage {
context_token: Some(String::new()),
from_user_id: Some("user-bcast-online".into()),
item_list: Some(Arc::new(vec![MessageItem {
item_type: Some(1),
text_item: Some(TextItem {
text: Some("hello broadcast".into()),
}),
..Default::default()
}])),
..Default::default()
};
handle_broadcast(Arc::clone(&state), msg).await;
assert_eq!(
state.metrics.messages_dispatched.load(Ordering::Relaxed),
before,
"empty context must not push_shared to online clients"
);
let drained = state.clients.queue.drain(&vtoken).await.expect("drain");
assert!(
drained.is_empty(),
"queue must stay empty when broadcast context is empty"
);
}
#[tokio::test]
async fn dispatch_broadcast_online_pushes_shared_to_queue() {
let (state, vtoken) = make_state_with_client().await;
state.clients.registry.write().await.mark_online(&vtoken);
let before = state.metrics.messages_dispatched.load(Ordering::Relaxed);
let msg = WeixinMessage {
context_token: Some("real-ctx-bcast".into()),
from_user_id: Some("user-bcast-push".into()),
item_list: Some(Arc::new(vec![MessageItem {
item_type: Some(1),
text_item: Some(TextItem {
text: Some("hello broadcast".into()),
}),
..Default::default()
}])),
..Default::default()
};
handle_broadcast(Arc::clone(&state), msg).await;
assert_eq!(
state.metrics.messages_dispatched.load(Ordering::Relaxed),
before + 1,
"broadcast must call push_shared_to_queue for each online client"
);
let drained = state.clients.queue.drain(&vtoken).await.expect("drain");
assert_eq!(
drained.len(),
1,
"exactly one shared message must be queued"
);
assert!(
drained[0]
.context_token
.as_deref()
.is_some_and(|c| !c.is_empty()),
"queued message must carry a non-empty vctx"
);
}
#[tokio::test]
async fn dispatch_at_mention_empty_context_skips_queue() {
let (state, vtoken) = make_state_with_client().await;
state.clients.registry.write().await.mark_online(&vtoken);
let before = state.metrics.messages_dispatched.load(Ordering::Relaxed);
let msg = WeixinMessage {
context_token: Some(String::new()),
from_user_id: Some("user-at".into()),
item_list: Some(Arc::new(vec![MessageItem {
item_type: Some(1),
text_item: Some(TextItem {
text: Some("@test-backend please help".into()),
}),
..Default::default()
}])),
..Default::default()
};
dispatch_message(Arc::clone(&state), msg).await;
assert_eq!(
state.metrics.messages_dispatched.load(Ordering::Relaxed),
before,
"@mention with empty context must not push to queue"
);
let drained = state.clients.queue.drain(&vtoken).await.expect("drain");
assert!(
drained.is_empty(),
"queue must stay empty for empty-ctx @mention"
);
}
#[tokio::test]
async fn resolve_quote_from_timestamp_rejects_empty_vtoken() {
let (state, _vtoken) = make_state_with_client().await;
let peer_user_id = "peer:user-empty-vt";
state
.store
.save_message(
"vctx-empty-vt",
Some(""),
"at-empty-vtoken",
peer_user_id,
"assistant",
"reply with empty vtoken",
)
.await
.expect("save");
let now_ms = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_millis() as i64;
let msg = make_quote_msg("user-empty-vt", now_ms, None);
let result = resolve_quote_from_timestamp(&state, &msg).await;
assert!(
result.is_none(),
"empty vtoken row must not become QuoteOrigin::Client"
);
}
#[tokio::test]
async fn resolve_quote_from_db_rejects_empty_vtoken() {
let (state, _vtoken) = make_state_with_client().await;
let peer_user_id = "peer:user-empty-vt-db";
let body = "unique-empty-vtoken-prefix-content-xyz";
state
.store
.save_message(
"vctx-empty-vt-db",
Some(""),
"at-empty-vtoken-db",
peer_user_id,
"assistant",
body,
)
.await
.expect("save");
let msg = make_quote_msg("user-empty-vt-db", 1_000_000, Some(body));
let result = resolve_quote_from_db(&state, &msg).await;
assert!(
result.is_none(),
"empty vtoken row must not become QuoteOrigin::Client via DB lookup"
);
}