use std::collections::HashMap;
use std::sync::Arc;
use async_trait::async_trait;
use serde_json::Value;
use nexo_tool_meta::admin::credentials::{
CredentialRegisterInput, CredentialRegisterResponse, CredentialSummary,
CredentialValidationOutcome, CredentialsListFilter, CredentialsListResponse,
CredentialsRevokeParams, CredentialsRevokeResponse,
};
use super::agents::YamlPatcher;
use crate::agent::admin_rpc::dispatcher::{AdminRpcError, AdminRpcResult};
#[async_trait]
pub trait ChannelCredentialPersister: Send + Sync {
fn channel(&self) -> &str;
fn validate_shape(
&self,
payload: &Value,
metadata: &HashMap<String, Value>,
) -> Result<(), AdminRpcError>;
async fn persist(
&self,
instance: Option<&str>,
payload: &Value,
metadata: &HashMap<String, Value>,
) -> Result<(), AdminRpcError>;
async fn revoke(&self, instance: Option<&str>) -> Result<bool, AdminRpcError>;
async fn probe(
&self,
_instance: Option<&str>,
_payload: &Value,
_metadata: &HashMap<String, Value>,
) -> CredentialValidationOutcome {
CredentialValidationOutcome {
probed: false,
healthy: false,
detail: None,
reason_code: Some(nexo_tool_meta::admin::credentials::reason_code::NOT_PROBED.into()),
}
}
}
pub type PersisterRegistry = HashMap<String, Arc<dyn ChannelCredentialPersister>>;
pub trait CredentialStore: Send + Sync {
fn list_credentials(&self) -> anyhow::Result<Vec<(String, Option<String>)>>;
fn write_credential(
&self,
channel: &str,
instance: Option<&str>,
payload: &Value,
) -> anyhow::Result<()>;
fn delete_credential(&self, channel: &str, instance: Option<&str>) -> anyhow::Result<bool>;
}
pub fn list(store: &dyn CredentialStore, yaml: &dyn YamlPatcher, params: Value) -> AdminRpcResult {
let filter: CredentialsListFilter = serde_json::from_value(params).unwrap_or_default();
let creds = match store.list_credentials() {
Ok(c) => c,
Err(e) => {
return AdminRpcResult::err(AdminRpcError::Internal(format!(
"credential store read: {e}"
)));
}
};
let agent_ids = match yaml.list_agent_ids() {
Ok(a) => a,
Err(e) => {
return AdminRpcResult::err(AdminRpcError::Internal(format!("yaml read: {e}")));
}
};
let mut summaries: Vec<CredentialSummary> = creds
.into_iter()
.filter(|(c, _)| match &filter.channel_filter {
Some(f) => c == f,
None => true,
})
.map(|(channel, instance)| {
let mut agents_using: Vec<String> = agent_ids
.iter()
.filter(|aid| agent_has_binding(yaml, aid, &channel, instance.as_deref()))
.cloned()
.collect();
agents_using.sort();
CredentialSummary {
channel,
instance,
agent_ids: agents_using,
}
})
.collect();
summaries.sort_by(|a, b| {
a.channel
.cmp(&b.channel)
.then_with(|| a.instance.cmp(&b.instance))
});
let response = CredentialsListResponse {
credentials: summaries,
};
AdminRpcResult::ok(serde_json::to_value(response).unwrap_or(Value::Null))
}
pub async fn register(
store: &dyn CredentialStore,
yaml: &dyn YamlPatcher,
persister: Option<Arc<dyn ChannelCredentialPersister>>,
params: Value,
reload_signal: &(dyn Fn() + Send + Sync),
) -> AdminRpcResult {
let input: CredentialRegisterInput = match serde_json::from_value(params) {
Ok(i) => i,
Err(e) => return AdminRpcResult::err(AdminRpcError::InvalidParams(e.to_string())),
};
if input.channel.is_empty() {
return AdminRpcResult::err(AdminRpcError::InvalidParams("channel is empty".into()));
}
if let Some(p) = persister.as_ref() {
if let Err(e) = p.validate_shape(&input.payload, &input.metadata) {
return AdminRpcResult::err(e);
}
}
if let Err(e) =
store.write_credential(&input.channel, input.instance.as_deref(), &input.payload)
{
return AdminRpcResult::err(AdminRpcError::Internal(format!("credential write: {e}")));
}
let mut persister_persisted = false;
if let Some(p) = persister.as_ref() {
if let Err(e) = p
.persist(input.instance.as_deref(), &input.payload, &input.metadata)
.await
{
return AdminRpcResult::err(e);
}
persister_persisted = true;
}
for agent_id in &input.agent_ids {
if let Err(e) = bind_agent(yaml, agent_id, &input.channel, input.instance.as_deref()) {
return AdminRpcResult::err(AdminRpcError::Internal(format!(
"agent `{agent_id}` bind failed: {e}"
)));
}
}
if !input.agent_ids.is_empty() || persister_persisted {
reload_signal();
}
let validation = match persister.as_ref() {
Some(p) => Some(
p.probe(input.instance.as_deref(), &input.payload, &input.metadata)
.await,
),
None => None,
};
let response = CredentialRegisterResponse {
summary: CredentialSummary {
channel: input.channel,
instance: input.instance,
agent_ids: input.agent_ids,
},
validation,
};
AdminRpcResult::ok(serde_json::to_value(response).unwrap_or(Value::Null))
}
pub async fn revoke(
store: &dyn CredentialStore,
yaml: &dyn YamlPatcher,
persister: Option<Arc<dyn ChannelCredentialPersister>>,
params: Value,
reload_signal: &(dyn Fn() + Send + Sync),
) -> AdminRpcResult {
let p: CredentialsRevokeParams = match serde_json::from_value(params) {
Ok(p) => p,
Err(e) => return AdminRpcResult::err(AdminRpcError::InvalidParams(e.to_string())),
};
let agent_ids = match yaml.list_agent_ids() {
Ok(a) => a,
Err(e) => {
return AdminRpcResult::err(AdminRpcError::Internal(format!("yaml read: {e}")));
}
};
let mut unbound = Vec::new();
for agent_id in &agent_ids {
if !agent_has_binding(yaml, agent_id, &p.channel, p.instance.as_deref()) {
continue;
}
if let Err(e) = unbind_agent(yaml, agent_id, &p.channel, p.instance.as_deref()) {
return AdminRpcResult::err(AdminRpcError::Internal(format!(
"agent `{agent_id}` unbind failed: {e}"
)));
}
unbound.push(agent_id.clone());
}
let mut persister_revoked = false;
if let Some(pp) = persister.as_ref() {
match pp.revoke(p.instance.as_deref()).await {
Ok(removed_in_persister) => persister_revoked = removed_in_persister,
Err(e) => return AdminRpcResult::err(e),
}
}
let removed = match store.delete_credential(&p.channel, p.instance.as_deref()) {
Ok(r) => r,
Err(e) => {
return AdminRpcResult::err(AdminRpcError::Internal(format!("credential delete: {e}")));
}
};
if removed || !unbound.is_empty() || persister_revoked {
reload_signal();
}
let response = CredentialsRevokeResponse {
removed,
unbound_agents: unbound,
};
AdminRpcResult::ok(serde_json::to_value(response).unwrap_or(Value::Null))
}
fn agent_has_binding(
yaml: &dyn YamlPatcher,
agent_id: &str,
channel: &str,
instance: Option<&str>,
) -> bool {
let Some(Value::Array(bindings)) = yaml
.read_agent_field(agent_id, "inbound_bindings")
.ok()
.flatten()
else {
return false;
};
bindings
.iter()
.any(|b| binding_matches(b, channel, instance))
}
fn binding_matches(b: &Value, channel: &str, instance: Option<&str>) -> bool {
let Some(plugin) = b.get("plugin").and_then(Value::as_str) else {
return false;
};
if plugin != channel {
return false;
}
let bind_instance = b.get("instance").and_then(Value::as_str);
bind_instance == instance
}
fn bind_agent(
yaml: &dyn YamlPatcher,
agent_id: &str,
channel: &str,
instance: Option<&str>,
) -> anyhow::Result<()> {
let mut bindings: Vec<Value> = match yaml.read_agent_field(agent_id, "inbound_bindings")? {
Some(Value::Array(arr)) => arr,
_ => Vec::new(),
};
if bindings
.iter()
.any(|b| binding_matches(b, channel, instance))
{
return Ok(());
}
let mut new_binding = serde_json::Map::new();
new_binding.insert("plugin".into(), Value::String(channel.into()));
if let Some(i) = instance {
new_binding.insert("instance".into(), Value::String(i.into()));
}
bindings.push(Value::Object(new_binding));
yaml.upsert_agent_field(agent_id, "inbound_bindings", Value::Array(bindings))?;
Ok(())
}
fn unbind_agent(
yaml: &dyn YamlPatcher,
agent_id: &str,
channel: &str,
instance: Option<&str>,
) -> anyhow::Result<()> {
let bindings: Vec<Value> = match yaml.read_agent_field(agent_id, "inbound_bindings")? {
Some(Value::Array(arr)) => arr,
_ => return Ok(()),
};
let filtered: Vec<Value> = bindings
.into_iter()
.filter(|b| !binding_matches(b, channel, instance))
.collect();
yaml.upsert_agent_field(agent_id, "inbound_bindings", Value::Array(filtered))?;
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use std::collections::HashMap;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
#[derive(Debug, Default)]
struct MockStore {
creds: Mutex<HashMap<(String, Option<String>), Value>>,
}
impl CredentialStore for MockStore {
fn list_credentials(&self) -> anyhow::Result<Vec<(String, Option<String>)>> {
Ok(self.creds.lock().unwrap().keys().cloned().collect())
}
fn write_credential(
&self,
channel: &str,
instance: Option<&str>,
payload: &Value,
) -> anyhow::Result<()> {
self.creds.lock().unwrap().insert(
(channel.to_string(), instance.map(String::from)),
payload.clone(),
);
Ok(())
}
fn delete_credential(&self, channel: &str, instance: Option<&str>) -> anyhow::Result<bool> {
Ok(self
.creds
.lock()
.unwrap()
.remove(&(channel.to_string(), instance.map(String::from)))
.is_some())
}
}
#[derive(Debug, Default)]
struct MockYaml {
agents: Mutex<HashMap<String, HashMap<String, Value>>>,
}
impl MockYaml {
fn with_agents(ids: &[&str]) -> Arc<Self> {
let me = Arc::new(Self::default());
for id in ids {
me.set(id, "model.provider", Value::String("minimax".into()));
me.set(id, "inbound_bindings", Value::Array(vec![]));
}
me
}
fn set(&self, agent_id: &str, dotted: &str, value: Value) {
self.agents
.lock()
.unwrap()
.entry(agent_id.to_string())
.or_default()
.insert(dotted.to_string(), value);
}
fn bindings(&self, agent_id: &str) -> Vec<Value> {
self.agents
.lock()
.unwrap()
.get(agent_id)
.and_then(|m| m.get("inbound_bindings").cloned())
.and_then(|v| v.as_array().cloned())
.unwrap_or_default()
}
}
impl YamlPatcher for MockYaml {
fn list_agent_ids(&self) -> anyhow::Result<Vec<String>> {
let mut ids: Vec<String> = self.agents.lock().unwrap().keys().cloned().collect();
ids.sort();
Ok(ids)
}
fn read_agent_field(&self, agent_id: &str, dotted: &str) -> anyhow::Result<Option<Value>> {
Ok(self
.agents
.lock()
.unwrap()
.get(agent_id)
.and_then(|m| m.get(dotted).cloned()))
}
fn upsert_agent_field(
&self,
agent_id: &str,
dotted: &str,
value: Value,
) -> anyhow::Result<()> {
self.set(agent_id, dotted, value);
Ok(())
}
fn remove_agent(&self, agent_id: &str) -> anyhow::Result<()> {
self.agents.lock().unwrap().remove(agent_id);
Ok(())
}
}
fn reload_counter() -> (Arc<AtomicUsize>, impl Fn()) {
let count = Arc::new(AtomicUsize::new(0));
let counter = Arc::clone(&count);
(count, move || {
counter.fetch_add(1, Ordering::Relaxed);
})
}
#[tokio::test]
async fn credentials_register_links_many_agents() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana", "carlos"]);
let (count, reload) = reload_counter();
let result = register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "whatsapp",
"instance": "shared",
"agent_ids": ["ana", "carlos"],
"payload": { "token": "wa.X" }
}),
&reload,
)
.await;
let response: CredentialRegisterResponse =
serde_json::from_value(result.result.unwrap()).unwrap();
assert_eq!(
response.summary.agent_ids,
vec!["ana".to_string(), "carlos".into()]
);
assert!(response.validation.is_none());
for agent in ["ana", "carlos"] {
let bindings = yaml.bindings(agent);
assert_eq!(bindings.len(), 1);
assert_eq!(bindings[0]["plugin"], "whatsapp");
assert_eq!(bindings[0]["instance"], "shared");
}
assert_eq!(count.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn credentials_register_idempotent_skips_duplicate_binding() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (_count, reload) = reload_counter();
let params = serde_json::json!({
"channel": "whatsapp",
"instance": "personal",
"agent_ids": ["ana"],
"payload": { "token": "v1" }
});
register(&*store, &*yaml, None, params.clone(), &reload).await;
register(&*store, &*yaml, None, params, &reload).await;
let bindings = yaml.bindings("ana");
assert_eq!(bindings.len(), 1);
}
#[tokio::test]
async fn credentials_register_with_empty_agent_ids_writes_store_only() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (count, reload) = reload_counter();
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "whatsapp",
"agent_ids": [],
"payload": { "x": 1 }
}),
&reload,
)
.await;
assert_eq!(yaml.bindings("ana").len(), 0);
assert_eq!(count.load(Ordering::Relaxed), 0);
assert_eq!(store.list_credentials().unwrap().len(), 1);
}
#[tokio::test]
async fn credentials_revoke_removes_bindings_from_all_agents() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana", "carlos", "bob"]);
let (_, reload) = reload_counter();
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "whatsapp",
"instance": "shared",
"agent_ids": ["ana", "carlos"],
"payload": {}
}),
&reload,
)
.await;
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "whatsapp",
"instance": "personal",
"agent_ids": ["bob"],
"payload": {}
}),
&reload,
)
.await;
let result = revoke(
&*store,
&*yaml,
None,
serde_json::json!({ "channel": "whatsapp", "instance": "shared" }),
&reload,
)
.await;
let response: CredentialsRevokeResponse =
serde_json::from_value(result.result.unwrap()).unwrap();
assert!(response.removed);
let mut unbound = response.unbound_agents.clone();
unbound.sort();
assert_eq!(unbound, vec!["ana".to_string(), "carlos".into()]);
assert_eq!(yaml.bindings("bob").len(), 1);
assert_eq!(yaml.bindings("ana").len(), 0);
assert_eq!(yaml.bindings("carlos").len(), 0);
}
#[tokio::test]
async fn credentials_revoke_unknown_credential_is_idempotent() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (_, reload) = reload_counter();
let result = revoke(
&*store,
&*yaml,
None,
serde_json::json!({ "channel": "whatsapp", "instance": "ghost" }),
&reload,
)
.await;
let response: CredentialsRevokeResponse =
serde_json::from_value(result.result.unwrap()).unwrap();
assert!(!response.removed);
assert!(response.unbound_agents.is_empty());
}
#[tokio::test]
async fn credentials_list_includes_only_filesystem_creds_with_agent_index() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana", "carlos"]);
let (_, reload) = reload_counter();
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "whatsapp",
"instance": "shared",
"agent_ids": ["ana", "carlos"],
"payload": {}
}),
&reload,
)
.await;
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "whatsapp",
"instance": "orphan",
"agent_ids": [],
"payload": {}
}),
&reload,
)
.await;
let result = list(&*store, &*yaml, Value::Null);
let response: CredentialsListResponse =
serde_json::from_value(result.result.unwrap()).unwrap();
assert_eq!(response.credentials.len(), 2);
let shared = response
.credentials
.iter()
.find(|c| c.instance.as_deref() == Some("shared"))
.unwrap();
let mut shared_agents = shared.agent_ids.clone();
shared_agents.sort();
assert_eq!(shared_agents, vec!["ana".to_string(), "carlos".into()]);
let orphan = response
.credentials
.iter()
.find(|c| c.instance.as_deref() == Some("orphan"))
.unwrap();
assert!(orphan.agent_ids.is_empty());
}
#[tokio::test]
async fn credentials_list_channel_filter() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (_, reload) = reload_counter();
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "whatsapp",
"agent_ids": ["ana"],
"payload": {}
}),
&reload,
)
.await;
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "future_telegram",
"agent_ids": [],
"payload": {}
}),
&reload,
)
.await;
let result = list(
&*store,
&*yaml,
serde_json::json!({ "channel_filter": "whatsapp" }),
);
let response: CredentialsListResponse =
serde_json::from_value(result.result.unwrap()).unwrap();
assert_eq!(response.credentials.len(), 1);
assert_eq!(response.credentials[0].channel, "whatsapp");
}
use crate::agent::admin_rpc::dispatcher::AdminRpcDispatcher;
use nexo_tool_meta::admin::credentials::reason_code;
struct MockPersister {
channel: String,
calls: Mutex<Vec<String>>,
fail_validate: bool,
fail_persist: bool,
fail_revoke: bool,
probe_unhealthy: bool,
}
impl MockPersister {
fn new(channel: &str) -> Self {
Self {
channel: channel.into(),
calls: Mutex::new(Vec::new()),
fail_validate: false,
fail_persist: false,
fail_revoke: false,
probe_unhealthy: false,
}
}
fn record(&self, call: &str) {
self.calls.lock().unwrap().push(call.into());
}
}
#[async_trait]
impl ChannelCredentialPersister for MockPersister {
fn channel(&self) -> &str {
&self.channel
}
fn validate_shape(
&self,
_payload: &Value,
_metadata: &HashMap<String, Value>,
) -> Result<(), AdminRpcError> {
self.record("validate_shape");
if self.fail_validate {
return Err(AdminRpcError::InvalidParams("mock validate".into()));
}
Ok(())
}
async fn persist(
&self,
_instance: Option<&str>,
_payload: &Value,
_metadata: &HashMap<String, Value>,
) -> Result<(), AdminRpcError> {
self.record("persist");
if self.fail_persist {
return Err(AdminRpcError::Internal("mock persist".into()));
}
Ok(())
}
async fn revoke(&self, _instance: Option<&str>) -> Result<bool, AdminRpcError> {
self.record("revoke");
if self.fail_revoke {
return Err(AdminRpcError::Internal("mock revoke".into()));
}
Ok(true)
}
async fn probe(
&self,
_instance: Option<&str>,
_payload: &Value,
_metadata: &HashMap<String, Value>,
) -> CredentialValidationOutcome {
self.record("probe");
if self.probe_unhealthy {
return CredentialValidationOutcome {
probed: true,
healthy: false,
detail: Some("mock unhealthy".into()),
reason_code: Some(reason_code::CONNECTIVITY_FAILED.into()),
};
}
CredentialValidationOutcome {
probed: true,
healthy: true,
detail: None,
reason_code: Some(reason_code::OK.into()),
}
}
}
#[test]
fn dispatcher_registry_starts_empty() {
let d = AdminRpcDispatcher::new();
assert!(d.persister_for("telegram").is_none());
assert!(d.persister_for("email").is_none());
}
#[test]
fn dispatcher_registry_registers_and_retrieves_persister() {
let mut d = AdminRpcDispatcher::new();
let mock = Arc::new(MockPersister::new("telegram"));
d.register_persister(mock.clone());
let fetched = d.persister_for("telegram").expect("registered");
assert_eq!(fetched.channel(), "telegram");
let _ = futures::executor::block_on(fetched.persist(None, &Value::Null, &HashMap::new()));
assert_eq!(
mock.calls.lock().unwrap().as_slice(),
&["persist".to_string()]
);
}
#[test]
fn dispatcher_registry_two_distinct_channels_coexist() {
let mut d = AdminRpcDispatcher::new();
d.register_persister(Arc::new(MockPersister::new("telegram")));
d.register_persister(Arc::new(MockPersister::new("email")));
assert!(d.persister_for("telegram").is_some());
assert!(d.persister_for("email").is_some());
assert!(d.persister_for("whatsapp").is_none());
}
#[test]
#[should_panic(expected = "duplicate ChannelCredentialPersister registration")]
fn dispatcher_registry_panics_on_duplicate_channel() {
let mut d = AdminRpcDispatcher::new();
d.register_persister(Arc::new(MockPersister::new("telegram")));
d.register_persister(Arc::new(MockPersister::new("telegram")));
}
#[tokio::test]
async fn persister_default_probe_returns_not_probed() {
struct MinimalPersister;
#[async_trait]
impl ChannelCredentialPersister for MinimalPersister {
fn channel(&self) -> &str {
"minimal"
}
fn validate_shape(
&self,
_: &Value,
_: &HashMap<String, Value>,
) -> Result<(), AdminRpcError> {
Ok(())
}
async fn persist(
&self,
_: Option<&str>,
_: &Value,
_: &HashMap<String, Value>,
) -> Result<(), AdminRpcError> {
Ok(())
}
async fn revoke(&self, _: Option<&str>) -> Result<bool, AdminRpcError> {
Ok(false)
}
}
let p = MinimalPersister;
let outcome = p.probe(None, &Value::Null, &HashMap::new()).await;
assert!(!outcome.probed);
assert!(!outcome.healthy);
assert_eq!(
outcome.reason_code.as_deref(),
Some(reason_code::NOT_PROBED)
);
}
#[tokio::test]
async fn credentials_register_empty_channel_returns_invalid_params() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&[]);
let (_, reload) = reload_counter();
let result = register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "",
"agent_ids": [],
"payload": {}
}),
&reload,
)
.await;
let err = result.error.expect("error");
assert!(matches!(err, AdminRpcError::InvalidParams(_)));
}
#[tokio::test]
async fn register_with_persister_runs_validate_persist_probe_in_order() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (count, reload) = reload_counter();
let persister = Arc::new(MockPersister::new("telegram"));
let result = register(
&*store,
&*yaml,
Some(persister.clone()),
serde_json::json!({
"channel": "telegram",
"instance": "kate",
"agent_ids": ["ana"],
"payload": { "token": "tg.X" },
"metadata": { "polling": { "enabled": true } }
}),
&reload,
)
.await;
let response: CredentialRegisterResponse =
serde_json::from_value(result.result.expect("ok")).unwrap();
assert_eq!(response.summary.channel, "telegram");
let v = response.validation.expect("validation present");
assert!(v.probed && v.healthy);
assert_eq!(v.reason_code.as_deref(), Some(reason_code::OK));
let calls = persister.calls.lock().unwrap().clone();
assert_eq!(
calls,
vec![
"validate_shape".to_string(),
"persist".into(),
"probe".into()
]
);
assert_eq!(store.list_credentials().unwrap().len(), 1);
assert_eq!(yaml.bindings("ana").len(), 1);
assert_eq!(count.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn register_aborts_when_validate_shape_fails_no_opaque_write() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (count, reload) = reload_counter();
let mut p = MockPersister::new("telegram");
p.fail_validate = true;
let persister = Arc::new(p);
let result = register(
&*store,
&*yaml,
Some(persister.clone()),
serde_json::json!({
"channel": "telegram",
"instance": "kate",
"agent_ids": ["ana"],
"payload": {}
}),
&reload,
)
.await;
assert!(matches!(
result.error.unwrap(),
AdminRpcError::InvalidParams(_)
));
assert_eq!(store.list_credentials().unwrap().len(), 0);
assert!(yaml.bindings("ana").is_empty());
assert_eq!(count.load(Ordering::Relaxed), 0);
assert_eq!(
persister.calls.lock().unwrap().as_slice(),
&["validate_shape".to_string()]
);
}
#[tokio::test]
async fn register_persister_persist_error_keeps_opaque_blob_skips_bind() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (count, reload) = reload_counter();
let mut p = MockPersister::new("telegram");
p.fail_persist = true;
let persister = Arc::new(p);
let result = register(
&*store,
&*yaml,
Some(persister.clone()),
serde_json::json!({
"channel": "telegram",
"instance": "kate",
"agent_ids": ["ana"],
"payload": { "token": "x" }
}),
&reload,
)
.await;
assert!(matches!(result.error.unwrap(), AdminRpcError::Internal(_)));
assert_eq!(store.list_credentials().unwrap().len(), 1);
assert!(yaml.bindings("ana").is_empty());
assert_eq!(count.load(Ordering::Relaxed), 0);
let calls = persister.calls.lock().unwrap().clone();
assert_eq!(calls, vec!["validate_shape".to_string(), "persist".into()]);
}
#[tokio::test]
async fn register_unhealthy_probe_does_not_abort_register() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (count, reload) = reload_counter();
let mut p = MockPersister::new("telegram");
p.probe_unhealthy = true;
let persister = Arc::new(p);
let result = register(
&*store,
&*yaml,
Some(persister.clone()),
serde_json::json!({
"channel": "telegram",
"instance": "kate",
"agent_ids": ["ana"],
"payload": { "token": "x" }
}),
&reload,
)
.await;
let response: CredentialRegisterResponse =
serde_json::from_value(result.result.expect("ok")).unwrap();
let v = response.validation.expect("validation present");
assert!(v.probed);
assert!(!v.healthy);
assert_eq!(
v.reason_code.as_deref(),
Some(reason_code::CONNECTIVITY_FAILED)
);
assert_eq!(store.list_credentials().unwrap().len(), 1);
assert_eq!(yaml.bindings("ana").len(), 1);
assert_eq!(count.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn register_no_persister_falls_back_to_opaque_path() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (_count, reload) = reload_counter();
let result = register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "ghost_channel",
"agent_ids": ["ana"],
"payload": {}
}),
&reload,
)
.await;
let response: CredentialRegisterResponse =
serde_json::from_value(result.result.unwrap()).unwrap();
assert!(response.validation.is_none());
assert_eq!(store.list_credentials().unwrap().len(), 1);
}
#[tokio::test]
async fn register_with_persister_fires_reload_even_without_agent_ids() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (count, reload) = reload_counter();
let persister = Arc::new(MockPersister::new("telegram"));
register(
&*store,
&*yaml,
Some(persister),
serde_json::json!({
"channel": "telegram",
"instance": "future",
"agent_ids": [],
"payload": { "token": "x" }
}),
&reload,
)
.await;
assert_eq!(count.load(Ordering::Relaxed), 1);
}
#[tokio::test]
async fn revoke_calls_persister_before_opaque_delete() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (_, reload) = reload_counter();
let persister = Arc::new(MockPersister::new("telegram"));
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "telegram",
"instance": "kate",
"agent_ids": ["ana"],
"payload": {}
}),
&reload,
)
.await;
let result = revoke(
&*store,
&*yaml,
Some(persister.clone()),
serde_json::json!({ "channel": "telegram", "instance": "kate" }),
&reload,
)
.await;
let response: CredentialsRevokeResponse =
serde_json::from_value(result.result.unwrap()).unwrap();
assert!(response.removed);
assert_eq!(
persister.calls.lock().unwrap().as_slice(),
&["revoke".to_string()]
);
}
#[tokio::test]
async fn revoke_persister_revoke_error_propagates_internal() {
let store = Arc::new(MockStore::default());
let yaml = MockYaml::with_agents(&["ana"]);
let (_, reload) = reload_counter();
let mut p = MockPersister::new("telegram");
p.fail_revoke = true;
let persister = Arc::new(p);
register(
&*store,
&*yaml,
None,
serde_json::json!({
"channel": "telegram",
"instance": "kate",
"agent_ids": ["ana"],
"payload": {}
}),
&reload,
)
.await;
let result = revoke(
&*store,
&*yaml,
Some(persister),
serde_json::json!({ "channel": "telegram", "instance": "kate" }),
&reload,
)
.await;
assert!(matches!(result.error.unwrap(), AdminRpcError::Internal(_)));
assert_eq!(store.list_credentials().unwrap().len(), 1);
}
}