use crate::envelope::{Envelope, Outcome};
use crate::event::GatewayHealth;
use crate::frame::ResBody;
use crate::ops::{Delivery, HealthArgs, RegisterChannelArgs, Report, Welcome};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use serde_json::Value;
pub const HELLO_TIMEOUT_MS: u64 = 5_000;
pub const HELLO_REQUIRED_MESSAGE: &str = "hello required first";
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case")]
pub enum MountKind {
#[default]
Agent,
Gateway,
Admin,
}
#[derive(
Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize, JsonSchema,
)]
#[serde(rename_all = "snake_case")]
pub enum Capability {
Register,
Report,
Inject,
Recycle,
Probe,
Typing,
Conversations,
}
impl Capability {
pub fn as_str(self) -> &'static str {
match self {
Capability::Register => "register",
Capability::Report => "report",
Capability::Inject => "inject",
Capability::Recycle => "recycle",
Capability::Probe => "probe",
Capability::Typing => "typing",
Capability::Conversations => "conversations",
}
}
pub const ALL: [Capability; 7] = [
Capability::Register,
Capability::Report,
Capability::Inject,
Capability::Recycle,
Capability::Probe,
Capability::Typing,
Capability::Conversations,
];
}
impl std::fmt::Display for Capability {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", deny_unknown_fields)]
pub struct AgentMount {
pub role: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub session: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub task_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub pid: Option<u32>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", deny_unknown_fields)]
pub struct GatewayMount {
pub gateway: String,
pub platform: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", deny_unknown_fields)]
pub struct ClusterMount {
pub cluster: String,
pub role: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(untagged)]
pub enum Mount {
Agent(AgentMount),
Gateway(GatewayMount),
Cluster(ClusterMount),
Admin,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct HelloArgs {
pub protocol: u16,
pub plugin: String,
pub version: String,
pub kind: MountKind,
pub capabilities: Vec<Capability>,
#[serde(skip_serializing_if = "Option::is_none")]
pub mount: Option<Mount>,
}
impl HelloArgs {
pub fn has(&self, capability: Capability) -> bool {
self.capabilities.contains(&capability)
}
pub fn missing(&self, expected: &[Capability]) -> Vec<Capability> {
expected.iter().copied().filter(|c| !self.has(*c)).collect()
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub struct HelloAck {
pub protocol: u16,
pub role: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub session_id: Option<String>,
pub generation: u64,
pub prose: String,
pub server: ServerInfo,
pub host_capabilities: Vec<Capability>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct ServerInfo {
pub connected: bool,
pub cluster: String,
pub name: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case", tag = "op", content = "args")]
pub enum PluginOp {
Hello(HelloArgs),
Report(Report),
SessionRegister(SessionRegisterArgs),
AssignAck(AssignAckArgs),
Send(Box<Envelope>),
Deliver(Delivery),
RegisterChannel(RegisterChannelArgs),
Health(HealthArgs),
Typing(TypingArgs),
Detach(DetachArgs),
}
impl PluginOp {
pub fn name(&self) -> &'static str {
match self {
PluginOp::Hello(_) => "hello",
PluginOp::Report(_) => "report",
PluginOp::SessionRegister(_) => "session_register",
PluginOp::AssignAck(_) => "assign_ack",
PluginOp::Send(_) => "send",
PluginOp::Deliver(_) => "deliver",
PluginOp::RegisterChannel(_) => "register_channel",
PluginOp::Health(_) => "health",
PluginOp::Typing(_) => "typing",
PluginOp::Detach(_) => "detach",
}
}
pub fn is_pre_auth(self) -> bool {
matches!(self, PluginOp::Hello(_))
}
pub fn is_gateway(self) -> bool {
matches!(
self,
PluginOp::Deliver(_)
| PluginOp::RegisterChannel(_)
| PluginOp::Health(_)
| PluginOp::Typing(_)
)
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case", tag = "op", content = "args")]
pub enum HostOp {
Welcome(HelloAck),
Assign(AssignArgs),
RenderSend(RenderSendArgs),
Probe(Value),
Recycle(RecycleArgs),
ConfigGet(ConfigGetArgs),
Bye(ByeNotice),
}
impl HostOp {
pub fn name(&self) -> &'static str {
match self {
HostOp::Welcome(_) => "welcome",
HostOp::Assign(_) => "assign",
HostOp::RenderSend(_) => "render_send",
HostOp::Probe(_) => "probe",
HostOp::Recycle(_) => "recycle",
HostOp::ConfigGet(_) => "config_get",
HostOp::Bye(_) => "bye",
}
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub struct AssignArgs {
pub envelope: Box<Envelope>,
pub prose: String,
pub task_id: String,
pub generation: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub parent: Option<Box<Envelope>>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct AssignAckArgs {
pub task_id: String,
pub accepted: bool,
#[serde(skip_serializing_if = "Option::is_none")]
pub reason: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct SessionRegisterArgs {
pub session_id: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub pid: Option<u32>,
pub generation: u64,
#[serde(skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub task_id: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub struct RenderSendArgs {
pub envelope: Box<Envelope>,
pub conversation: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub gateway_ref: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub reply_to: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct RecycleArgs {
pub task_id: String,
pub reason: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub outcome: Option<Outcome>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct ConfigGetArgs {
pub key: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct TypingArgs {
pub conversation: String,
pub on: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct DetachArgs {
pub reason: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct ByeNotice {
pub reason: String,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, JsonSchema)]
#[serde(untagged)]
pub enum AdapterMsg {
Plugin(PluginOp),
Host(HostOp),
Res(ResBody),
}
impl AdapterMsg {
pub fn direction(&self) -> MsgDirection {
match self {
AdapterMsg::Plugin(_) => MsgDirection::ToHost,
AdapterMsg::Host(_) => MsgDirection::ToPlugin,
AdapterMsg::Res(body) => {
if body.ok {
MsgDirection::Response
} else {
MsgDirection::ErrorResponse
}
}
}
}
pub fn op_name(&self) -> Option<&str> {
match self {
AdapterMsg::Plugin(op) => Some(op.name()),
AdapterMsg::Host(op) => Some(op.name()),
AdapterMsg::Res(_) => None,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum MsgDirection {
ToHost,
ToPlugin,
Response,
ErrorResponse,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, JsonSchema, Default)]
#[serde(rename_all = "snake_case", default)]
pub struct GatewayBinding {
pub gateway: String,
pub channel: String,
pub conversation: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub title: Option<String>,
}
impl From<GatewayHealth> for HealthArgs {
fn from(state: GatewayHealth) -> Self {
HealthArgs {
state: state.as_str().to_string(),
detail: None,
uptime_s: 0,
}
}
}
pub type WelcomeSlice = Welcome;
#[cfg(test)]
mod tests {
use super::*;
use crate::envelope::{Body, MsgKind, Principal, new_envelope, new_task_id};
use crate::ops::SessionProjection;
fn note(text: &str) -> Envelope {
new_envelope(
MsgKind::Note,
Principal::role("planner"),
Principal::gateway("tg1", "telegram", Some("42".into())),
Body::text(text),
None,
)
.expect("note")
}
#[test]
fn hello_matches_the_documented_shape() {
let hello = PluginOp::Hello(HelloArgs {
protocol: crate::PROTOCOL_VERSION,
plugin: "onlyne-agent-pi".into(),
version: "1.0.0".into(),
kind: MountKind::Agent,
capabilities: vec![
Capability::Register,
Capability::Report,
Capability::Inject,
Capability::Recycle,
],
mount: Some(Mount::Agent(AgentMount {
role: "planner".into(),
session: Some("8b1c".into()),
task_id: None,
pid: Some(4212),
})),
});
let value = serde_json::to_value(&hello).expect("encode");
assert_eq!(value["op"], "hello");
assert_eq!(value["args"]["kind"], "agent");
assert_eq!(value["args"]["capabilities"][0], "register");
assert_eq!(value["args"]["mount"]["role"], "planner");
assert_eq!(value["args"]["mount"]["session"], "8b1c");
assert!(
value["args"]["mount"].get("kind").is_none(),
"the mount object is flat; `kind` sits beside it in args"
);
let back: PluginOp = serde_json::from_value(value).expect("decode");
assert_eq!(back, hello);
}
#[test]
fn mounts_disambiguate_by_field_set_in_declaration_order() {
let cluster = Mount::Cluster(ClusterMount {
cluster: "cluster-b".into(),
role: "cluster-b".into(),
});
let value = serde_json::to_value(&cluster).expect("encode cluster mount");
assert_eq!(
value,
serde_json::json!({"cluster": "cluster-b", "role": "cluster-b"})
);
let back: Mount = serde_json::from_value(value).expect("decode cluster mount");
assert_eq!(
back, cluster,
"a cluster mount must not decode as an agent mount"
);
let gateway = Mount::Gateway(GatewayMount {
gateway: "gw1".into(),
platform: "telegram".into(),
});
let value = serde_json::to_value(&gateway).expect("encode gateway mount");
assert_eq!(
value,
serde_json::json!({"gateway": "gw1", "platform": "telegram"})
);
assert_eq!(
serde_json::from_value::<Mount>(value).expect("decode gateway mount"),
gateway
);
let admin: Mount =
serde_json::from_value(serde_json::json!(null)).expect("decode admin mount");
assert_eq!(admin, Mount::Admin);
assert_eq!(
serde_json::to_value(Mount::Admin).expect("encode admin mount"),
serde_json::Value::Null
);
serde_json::from_value::<Mount>(serde_json::json!({"cluster": "b"}))
.expect_err("a mount missing its identifying fields matches no variant");
serde_json::from_value::<Mount>(serde_json::json!({}))
.expect_err("an empty mount names nothing and is refused");
}
#[test]
fn every_capability_the_plan_declares_round_trips() {
for name in ["register", "report", "inject", "recycle", "probe", "typing"] {
let capability: Capability = serde_json::from_value(serde_json::json!(name))
.unwrap_or_else(|e| panic!("{name} is a declared capability but refused: {e}"));
assert_eq!(capability.as_str(), name, "{name} does not round-trip");
}
}
#[test]
fn capability_gap_is_reported_in_declaration_order() {
let hello = HelloArgs {
protocol: crate::PROTOCOL_VERSION,
plugin: "onlyne-cli-agent".into(),
version: "1.0.0".into(),
kind: MountKind::Agent,
capabilities: vec![Capability::Report],
mount: None,
};
assert_eq!(
hello.missing(&[Capability::Report, Capability::Recycle, Capability::Inject]),
vec![Capability::Recycle, Capability::Inject]
);
assert!(!hello.has(Capability::Probe));
}
#[test]
fn untagged_adapter_msg_reads_both_directions() {
let to_host = serde_json::json!({"op":"send","args":note("hi")});
let msg: AdapterMsg = serde_json::from_value(to_host).expect("plugin op decodes");
assert_eq!(msg.direction(), MsgDirection::ToHost);
assert_eq!(msg.op_name(), Some("send"));
let to_plugin = serde_json::json!({"op":"probe","args":null});
let msg: AdapterMsg = serde_json::from_value(to_plugin).expect("host op decodes");
assert_eq!(msg.direction(), MsgDirection::ToPlugin);
assert_eq!(msg.op_name(), Some("probe"));
let reply = serde_json::json!({"ok":false,"error":{"code":"invalid","message":"x"}});
let msg: AdapterMsg = serde_json::from_value(reply).expect("res decodes");
assert_eq!(msg.direction(), MsgDirection::ErrorResponse);
assert_eq!(msg.op_name(), None);
}
#[test]
fn a_frame_that_is_not_an_op_is_rejected() {
let err = serde_json::from_value::<AdapterMsg>(serde_json::json!({"type":"noise"}))
.expect_err("unrecognised");
assert!(
err.to_string().contains("data did not match"),
"err = {err}"
);
}
#[test]
fn gateway_only_ops_are_marked_and_preauth_is_hello_only() {
assert!(PluginOp::Health(HealthArgs::default()).is_gateway());
assert!(
PluginOp::Deliver(Delivery {
msg_id: "m".into(),
envelope: Box::new(note("x")),
})
.is_gateway()
);
assert!(!PluginOp::Send(Box::new(note("y"))).is_gateway());
assert!(PluginOp::Hello(HelloArgs::default()).is_pre_auth());
assert!(!PluginOp::Send(Box::new(note("y"))).is_pre_auth());
}
#[test]
fn assign_carries_prose_alongside_the_envelope() {
let args = AssignArgs {
envelope: Box::new(note("do it")),
prose: "Read the incoming task".into(),
task_id: new_task_id(),
generation: 1,
parent: None,
};
let value = serde_json::to_value(HostOp::Assign(args.clone())).expect("encode");
assert_eq!(value["op"], "assign");
assert_eq!(value["args"]["prose"], "Read the incoming task");
assert!(value["args"].get("parent").is_none());
let back: HostOp = serde_json::from_value(value).expect("decode");
assert_eq!(back, HostOp::Assign(args));
}
#[test]
fn report_frames_use_the_lifecycle_kind() {
let op = PluginOp::Report(Report::Heartbeat {
task_id: new_task_id(),
generation: 1,
seq: 14,
observed: serde_json::json!({"state": "running"}),
cluster_ref: None,
});
let value = serde_json::to_value(&op).expect("encode");
assert_eq!(value["op"], "report");
assert_eq!(value["args"]["kind"], "heartbeat");
assert_eq!(value["args"]["data"]["seq"], 14);
let back: PluginOp = serde_json::from_value(value).expect("decode");
assert_eq!(back, op);
}
#[test]
fn session_register_accepts_a_missing_task_binding() {
let args = SessionRegisterArgs {
session_id: "8b1c".into(),
pid: Some(4212),
generation: 1,
title: Some("swarm:planner:8b1c".into()),
task_id: None,
};
let value = serde_json::to_value(PluginOp::SessionRegister(args)).expect("encode");
assert_eq!(value["op"], "session_register");
assert!(value["args"].get("task_id").is_none());
}
#[test]
fn detach_and_bye_record_their_reasons() {
let detach = PluginOp::Detach(DetachArgs {
reason: "operator".into(),
});
assert_eq!(
serde_json::to_value(&detach).expect("encode")["args"]["reason"],
"operator"
);
let bye = HostOp::Bye(ByeNotice {
reason: "shutdown".into(),
});
assert_eq!(bye.name(), "bye");
assert_eq!(
serde_json::to_value(&bye).expect("encode")["args"]["reason"],
"shutdown"
);
}
#[test]
fn health_maps_from_gateway_health() {
let args: HealthArgs = GatewayHealth::Reconnecting.into();
assert_eq!(args.state, "reconnecting");
assert_eq!(args.uptime_s, 0);
assert!(args.detail.is_none());
}
#[test]
fn projection_travels_inside_the_welcome_slice() {
let welcome = Welcome {
cluster: "cluster-a".into(),
server: "srv".into(),
role: "planner".into(),
admin: false,
max_sessions: 3,
reuse: true,
prose: "Read the incoming task".into(),
spec_hash: "abc".into(),
aggregate: Some("cluster-b".into()),
allowed_targets: vec!["builder".into()],
allowed_senders: vec!["*".into()],
session_command: Some(vec!["pi".into(), "--session-id".into(), "{session}".into()]),
timeout_ready_ms: Some(30_000),
timeout_running_ms: Some(120_000),
timeout_idle_ms: Some(60_000),
intent_attempts: Some(3),
intent_backoff_ms: Some(vec![1000, 2000, 4000]),
relay_required: Some(vec!["writer".into()]),
relay_count: None,
seq: 41,
};
assert_eq!(welcome.role, "planner");
assert_eq!(welcome.max_sessions, 3);
assert_eq!(
SessionProjection::default_working().lifecycle,
LifecycleMarker::created()
);
}
#[test]
fn untagged_adapter_msg_resolves_plugin_host_and_response() {
let plugin = AdapterMsg::Plugin(PluginOp::Report(Report::Heartbeat {
task_id: new_task_id(),
generation: 1,
seq: 14,
observed: serde_json::json!({"state": "running"}),
cluster_ref: None,
}));
let value = serde_json::to_value(&plugin).expect("encode plugin op");
assert_eq!(value["op"], "report");
let back: AdapterMsg = serde_json::from_value(value).expect("decode plugin op");
assert_eq!(back, plugin);
assert_eq!(back.direction(), MsgDirection::ToHost);
assert_eq!(back.op_name(), Some("report"));
let host = AdapterMsg::Host(HostOp::Recycle(RecycleArgs {
task_id: new_task_id(),
reason: "retire".into(),
outcome: Some(Outcome::Done),
}));
let value = serde_json::to_value(&host).expect("encode host op");
assert_eq!(value["op"], "recycle");
let back: AdapterMsg = serde_json::from_value(value).expect("decode host op");
assert_eq!(back, host);
assert_eq!(back.direction(), MsgDirection::ToPlugin);
assert_eq!(back.op_name(), Some("recycle"));
let res = AdapterMsg::Res(ResBody::ok(serde_json::json!({"state": "in_flight"})));
let value = serde_json::to_value(&res).expect("encode res body");
assert_eq!(value["ok"], true);
let back: AdapterMsg = serde_json::from_value(value).expect("decode res body");
assert_eq!(back, res);
assert_eq!(back.direction(), MsgDirection::Response);
assert_eq!(back.op_name(), None);
}
struct LifecycleMarker;
impl LifecycleMarker {
fn created() -> crate::event::Lifecycle {
crate::event::Lifecycle::Created
}
}
}