use std::collections::BTreeMap;
use serde::{Deserialize, Serialize};
pub mod artifacts;
pub mod control;
pub mod credentials;
pub mod execution;
pub mod permissions;
pub mod session;
pub mod streaming;
pub mod subscriptions;
pub mod telemetry;
pub use artifacts::{
ArtifactFetchPayload, ArtifactPutPayload, ArtifactRef, ArtifactRefPayload,
ArtifactReleasePayload,
};
pub use control::{
AckPayload, BackpressurePayload, CancelAcceptedPayload, CancelPayload, CancelRefusedPayload,
CancelTargetKind, InterruptPayload, NackPayload, ResumePayload,
};
pub use credentials::{CredentialId, CredentialScheme, ProvisionedCredential};
pub use execution::{
AgentDelegatePayload, AgentHandoffPayload, AgentRef, AgentRefParseError, JobAcceptedPayload,
JobCancelledPayload, JobCheckpointPayload, JobCompletedPayload, JobFailedPayload,
JobHeartbeatPayload, JobProgressPayload, JobResultChunkPayload, JobSchedulePayload,
JobStartedPayload, JobState, ResultChunkAssembler, ResultChunkEncoding, ResultChunkError,
ToolErrorPayload, ToolInvokePayload, ToolResultPayload, WorkflowCompletePayload,
WorkflowStartPayload,
};
pub use permissions::{
CostBudget, CostBudgetAmount, CostBudgetParseError, LeaseExtendedPayload, LeaseGrantedPayload,
LeaseRefreshPayload, LeaseRequest, LeaseRevokedPayload, LeaseSubsetViolation, ModelUse,
PermissionDenyPayload, PermissionGrantPayload, PermissionRequestPayload, TrustLevel,
};
pub use session::{
AuthScheme, ClientIdentity, Credentials, JobListEntry, RuntimeIdentity, SessionAcceptedPayload,
SessionAckPayload, SessionAuthenticatePayload, SessionChallengePayload, SessionClosePayload,
SessionEvictedPayload, SessionJobsPayload, SessionLease, SessionListJobsFilter,
SessionListJobsPayload, SessionOpenPayload, SessionPingPayload, SessionPongPayload,
SessionRefreshPayload, SessionRejectedPayload, SessionUnauthenticatedPayload,
};
pub use streaming::{
StreamChunkPayload, StreamClosePayload, StreamErrorPayload, StreamKind, StreamOpenPayload,
};
pub use subscriptions::{
JobSubscribePayload, JobSubscribedPayload, JobUnsubscribePayload, SubscribeAcceptedPayload,
SubscribeClosedPayload, SubscribeEventPayload, SubscribePayload, SubscriptionFilter,
SubscriptionSince, UnsubscribePayload,
};
pub use telemetry::TraceSpanPayload;
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
#[serde(default)]
pub struct Capabilities {
#[serde(skip_serializing_if = "Option::is_none")]
pub streaming: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub durable_jobs: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub checkpoints: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub binary_streams: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub agent_handoff: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub model_use: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub provisioned_credentials: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub artifacts: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub subscriptions: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub scheduled_jobs: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub interrupt: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub anonymous: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub heartbeat_recovery: Option<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub binary_encoding: Vec<String>,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub extensions: Vec<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub artifact_retention: Option<ArtifactRetention>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub agents: Option<AgentInventory>,
#[serde(flatten)]
pub extra: BTreeMap<String, serde_json::Value>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ArtifactRetention {
pub default_seconds: u64,
pub max_seconds: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct AgentInventoryEntry {
pub name: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub versions: Vec<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub default: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(untagged)]
pub enum AgentInventory {
Flat(Vec<String>),
Rich(Vec<AgentInventoryEntry>),
}
impl AgentInventory {
#[must_use]
pub fn into_rich(self) -> Vec<AgentInventoryEntry> {
match self {
Self::Flat(names) => names
.into_iter()
.map(|name| AgentInventoryEntry {
name,
versions: vec![],
default: None,
})
.collect(),
Self::Rich(entries) => entries,
}
}
#[must_use]
pub fn satisfies(&self, agent: &execution::AgentRef) -> bool {
match self {
Self::Flat(names) => names.iter().any(|n| n == &agent.name),
Self::Rich(entries) => entries.iter().any(|e| {
e.name == agent.name
&& agent
.version
.as_ref()
.is_none_or(|v| e.versions.iter().any(|known| known == v))
}),
}
}
}
impl Capabilities {
#[must_use]
pub const fn has(&self, name: CapabilityName) -> bool {
match name {
CapabilityName::Streaming => matches!(self.streaming, Some(true)),
CapabilityName::DurableJobs => matches!(self.durable_jobs, Some(true)),
CapabilityName::Checkpoints => matches!(self.checkpoints, Some(true)),
CapabilityName::BinaryStreams => matches!(self.binary_streams, Some(true)),
CapabilityName::AgentHandoff => matches!(self.agent_handoff, Some(true)),
CapabilityName::ModelUse => matches!(self.model_use, Some(true)),
CapabilityName::ProvisionedCredentials => {
matches!(self.provisioned_credentials, Some(true))
}
CapabilityName::Artifacts => matches!(self.artifacts, Some(true)),
CapabilityName::Subscriptions => matches!(self.subscriptions, Some(true)),
CapabilityName::ScheduledJobs => matches!(self.scheduled_jobs, Some(true)),
CapabilityName::Interrupt => matches!(self.interrupt, Some(true)),
CapabilityName::Anonymous => matches!(self.anonymous, Some(true)),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum CapabilityName {
Streaming,
DurableJobs,
Checkpoints,
BinaryStreams,
AgentHandoff,
ModelUse,
ProvisionedCredentials,
Artifacts,
Subscriptions,
ScheduledJobs,
Interrupt,
Anonymous,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "type", content = "payload")]
#[non_exhaustive]
pub enum MessageType {
#[serde(rename = "session.open")]
SessionOpen(SessionOpenPayload),
#[serde(rename = "session.challenge")]
SessionChallenge(SessionChallengePayload),
#[serde(rename = "session.authenticate")]
SessionAuthenticate(SessionAuthenticatePayload),
#[serde(rename = "session.accepted")]
SessionAccepted(SessionAcceptedPayload),
#[serde(rename = "session.unauthenticated")]
SessionUnauthenticated(SessionUnauthenticatedPayload),
#[serde(rename = "session.rejected")]
SessionRejected(SessionRejectedPayload),
#[serde(rename = "session.refresh")]
SessionRefresh(SessionRefreshPayload),
#[serde(rename = "session.evicted")]
SessionEvicted(SessionEvictedPayload),
#[serde(rename = "session.close")]
SessionClose(SessionClosePayload),
#[serde(rename = "session.ping")]
SessionPing(SessionPingPayload),
#[serde(rename = "session.pong")]
SessionPong(SessionPongPayload),
#[serde(rename = "session.ack")]
SessionAck(SessionAckPayload),
#[serde(rename = "session.list_jobs")]
SessionListJobs(SessionListJobsPayload),
#[serde(rename = "session.jobs")]
SessionJobs(SessionJobsPayload),
#[serde(rename = "ping")]
Ping(PingPayload),
#[serde(rename = "pong")]
Pong(PongPayload),
#[serde(rename = "ack")]
Ack(AckPayload),
#[serde(rename = "nack")]
Nack(NackPayload),
#[serde(rename = "cancel")]
Cancel(CancelPayload),
#[serde(rename = "cancel.accepted")]
CancelAccepted(CancelAcceptedPayload),
#[serde(rename = "cancel.refused")]
CancelRefused(CancelRefusedPayload),
#[serde(rename = "interrupt")]
Interrupt(InterruptPayload),
#[serde(rename = "resume")]
Resume(ResumePayload),
#[serde(rename = "backpressure")]
Backpressure(BackpressurePayload),
#[serde(rename = "tool.invoke")]
ToolInvoke(ToolInvokePayload),
#[serde(rename = "tool.result")]
ToolResult(ToolResultPayload),
#[serde(rename = "tool.error")]
ToolError(ToolErrorPayload),
#[serde(rename = "job.accepted")]
JobAccepted(JobAcceptedPayload),
#[serde(rename = "job.started")]
JobStarted(JobStartedPayload),
#[serde(rename = "job.progress")]
JobProgress(JobProgressPayload),
#[serde(rename = "job.heartbeat")]
JobHeartbeat(JobHeartbeatPayload),
#[serde(rename = "job.completed")]
JobCompleted(JobCompletedPayload),
#[serde(rename = "job.failed")]
JobFailed(JobFailedPayload),
#[serde(rename = "job.cancelled")]
JobCancelled(JobCancelledPayload),
#[serde(rename = "job.result_chunk")]
JobResultChunk(JobResultChunkPayload),
#[serde(rename = "stream.open")]
StreamOpen(StreamOpenPayload),
#[serde(rename = "stream.chunk")]
StreamChunk(StreamChunkPayload),
#[serde(rename = "stream.close")]
StreamClose(StreamClosePayload),
#[serde(rename = "stream.error")]
StreamError(StreamErrorPayload),
#[serde(rename = "permission.request")]
PermissionRequest(PermissionRequestPayload),
#[serde(rename = "permission.grant")]
PermissionGrant(PermissionGrantPayload),
#[serde(rename = "permission.deny")]
PermissionDeny(PermissionDenyPayload),
#[serde(rename = "lease.granted")]
LeaseGranted(LeaseGrantedPayload),
#[serde(rename = "lease.extended")]
LeaseExtended(LeaseExtendedPayload),
#[serde(rename = "lease.revoked")]
LeaseRevoked(LeaseRevokedPayload),
#[serde(rename = "lease.refresh")]
LeaseRefresh(LeaseRefreshPayload),
#[serde(rename = "subscribe")]
Subscribe(SubscribePayload),
#[serde(rename = "subscribe.accepted")]
SubscribeAccepted(SubscribeAcceptedPayload),
#[serde(rename = "subscribe.event")]
SubscribeEvent(SubscribeEventPayload),
#[serde(rename = "unsubscribe")]
Unsubscribe(UnsubscribePayload),
#[serde(rename = "subscribe.closed")]
SubscribeClosed(SubscribeClosedPayload),
#[serde(rename = "job.subscribe")]
JobSubscribe(JobSubscribePayload),
#[serde(rename = "job.subscribed")]
JobSubscribed(JobSubscribedPayload),
#[serde(rename = "job.unsubscribe")]
JobUnsubscribe(JobUnsubscribePayload),
#[serde(rename = "artifact.put")]
ArtifactPut(ArtifactPutPayload),
#[serde(rename = "artifact.fetch")]
ArtifactFetch(ArtifactFetchPayload),
#[serde(rename = "artifact.ref")]
ArtifactRef(ArtifactRefPayload),
#[serde(rename = "artifact.release")]
ArtifactRelease(ArtifactReleasePayload),
#[serde(rename = "event.emit")]
EventEmit(EventEmitPayload),
#[serde(rename = "log")]
Log(LogPayload),
#[serde(rename = "metric")]
Metric(MetricPayload),
#[serde(rename = "trace.span")]
TraceSpan(TraceSpanPayload),
}
impl MessageType {
#[must_use]
pub const fn type_name(&self) -> &'static str {
match self {
Self::SessionOpen(_) => "session.open",
Self::SessionChallenge(_) => "session.challenge",
Self::SessionAuthenticate(_) => "session.authenticate",
Self::SessionAccepted(_) => "session.accepted",
Self::SessionUnauthenticated(_) => "session.unauthenticated",
Self::SessionRejected(_) => "session.rejected",
Self::SessionRefresh(_) => "session.refresh",
Self::SessionEvicted(_) => "session.evicted",
Self::SessionClose(_) => "session.close",
Self::SessionPing(_) => "session.ping",
Self::SessionPong(_) => "session.pong",
Self::SessionAck(_) => "session.ack",
Self::SessionListJobs(_) => "session.list_jobs",
Self::SessionJobs(_) => "session.jobs",
Self::Ping(_) => "ping",
Self::Pong(_) => "pong",
Self::Ack(_) => "ack",
Self::Nack(_) => "nack",
Self::Cancel(_) => "cancel",
Self::CancelAccepted(_) => "cancel.accepted",
Self::CancelRefused(_) => "cancel.refused",
Self::Interrupt(_) => "interrupt",
Self::Resume(_) => "resume",
Self::Backpressure(_) => "backpressure",
Self::ToolInvoke(_) => "tool.invoke",
Self::ToolResult(_) => "tool.result",
Self::ToolError(_) => "tool.error",
Self::JobAccepted(_) => "job.accepted",
Self::JobStarted(_) => "job.started",
Self::JobProgress(_) => "job.progress",
Self::JobHeartbeat(_) => "job.heartbeat",
Self::JobCompleted(_) => "job.completed",
Self::JobFailed(_) => "job.failed",
Self::JobCancelled(_) => "job.cancelled",
Self::JobResultChunk(_) => "job.result_chunk",
Self::StreamOpen(_) => "stream.open",
Self::StreamChunk(_) => "stream.chunk",
Self::StreamClose(_) => "stream.close",
Self::StreamError(_) => "stream.error",
Self::PermissionRequest(_) => "permission.request",
Self::PermissionGrant(_) => "permission.grant",
Self::PermissionDeny(_) => "permission.deny",
Self::LeaseGranted(_) => "lease.granted",
Self::LeaseExtended(_) => "lease.extended",
Self::LeaseRevoked(_) => "lease.revoked",
Self::LeaseRefresh(_) => "lease.refresh",
Self::Subscribe(_) => "subscribe",
Self::SubscribeAccepted(_) => "subscribe.accepted",
Self::SubscribeEvent(_) => "subscribe.event",
Self::Unsubscribe(_) => "unsubscribe",
Self::SubscribeClosed(_) => "subscribe.closed",
Self::JobSubscribe(_) => "job.subscribe",
Self::JobSubscribed(_) => "job.subscribed",
Self::JobUnsubscribe(_) => "job.unsubscribe",
Self::ArtifactPut(_) => "artifact.put",
Self::ArtifactFetch(_) => "artifact.fetch",
Self::ArtifactRef(_) => "artifact.ref",
Self::ArtifactRelease(_) => "artifact.release",
Self::EventEmit(_) => "event.emit",
Self::Log(_) => "log",
Self::Metric(_) => "metric",
Self::TraceSpan(_) => "trace.span",
}
}
#[must_use]
pub const fn is_handshake(&self) -> bool {
matches!(
self,
Self::SessionOpen(_)
| Self::SessionChallenge(_)
| Self::SessionAuthenticate(_)
| Self::SessionAccepted(_)
| Self::SessionUnauthenticated(_)
| Self::SessionRejected(_)
)
}
#[must_use]
pub const fn is_countable_event(&self) -> bool {
!matches!(
self,
Self::SessionOpen(_)
| Self::SessionChallenge(_)
| Self::SessionAuthenticate(_)
| Self::SessionAccepted(_)
| Self::SessionUnauthenticated(_)
| Self::SessionRejected(_)
| Self::SessionRefresh(_)
| Self::SessionEvicted(_)
| Self::SessionClose(_)
| Self::SessionPing(_)
| Self::SessionPong(_)
| Self::SessionAck(_)
| Self::SessionListJobs(_)
| Self::SessionJobs(_)
| Self::Ping(_)
| Self::Pong(_)
)
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct PingPayload {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub nonce: Option<String>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct PongPayload {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub nonce: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct EventEmitPayload {
pub name: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub data: Option<serde_json::Value>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum LogLevel {
Trace,
Debug,
Info,
Warn,
Error,
Critical,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LogPayload {
pub level: LogLevel,
pub message: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub attributes: Option<serde_json::Value>,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct MetricPayload {
pub name: String,
pub value: f64,
pub unit: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub dims: Option<serde_json::Value>,
}
pub mod standard_names {
pub const TOKENS_USED: &str = "tokens.used";
pub const COST_USD: &str = "cost.usd";
pub const GPU_SECONDS: &str = "gpu.seconds";
pub const TOOL_INVOCATIONS: &str = "tool.invocations";
pub const LATENCY_MS: &str = "latency.ms";
pub const BYTES_IN: &str = "bytes.in";
pub const BYTES_OUT: &str = "bytes.out";
pub const ERRORS_TOTAL: &str = "errors.total";
}
pub const SUBSCRIPTION_BACKFILL_COMPLETE: &str = "subscription.backfill_complete";
#[cfg(test)]
#[allow(
clippy::expect_used,
clippy::unwrap_used,
clippy::panic,
clippy::missing_panics_doc
)]
mod tests {
use super::*;
#[test]
fn ping_round_trips_through_serde() {
let m = MessageType::Ping(PingPayload {
nonce: Some("abc".into()),
});
let json = serde_json::to_string(&m).expect("serialize");
let back: MessageType = serde_json::from_str(&json).expect("deserialize");
assert_eq!(m, back);
}
#[test]
fn session_ping_round_trips_through_serde() {
let now = chrono::Utc::now();
let m = MessageType::SessionPing(SessionPingPayload {
nonce: "p_01J".into(),
sent_at: now,
});
let json = serde_json::to_string(&m).expect("serialize");
let back: MessageType = serde_json::from_str(&json).expect("deserialize");
assert_eq!(m, back);
}
#[test]
fn session_ack_round_trips_through_serde() {
let m = MessageType::SessionAck(SessionAckPayload {
last_processed_seq: 1827,
});
let json = serde_json::to_string(&m).expect("serialize");
let back: MessageType = serde_json::from_str(&json).expect("deserialize");
assert_eq!(m, back);
let v: serde_json::Value = serde_json::from_str(&json).expect("value");
assert_eq!(
v,
serde_json::json!({
"type": "session.ack",
"payload": { "last_processed_seq": 1827 },
})
);
}
#[test]
fn session_list_jobs_round_trips_through_serde() {
let m = MessageType::SessionListJobs(SessionListJobsPayload {
filter: Some(SessionListJobsFilter {
status: vec!["running".into()],
agent: Some("echo".into()),
created_after: None,
created_before: None,
}),
limit: Some(50),
cursor: None,
});
let json = serde_json::to_string(&m).expect("serialize");
let back: MessageType = serde_json::from_str(&json).expect("deserialize");
assert_eq!(m, back);
}
#[test]
fn session_jobs_round_trips_through_serde() {
let now = chrono::Utc::now();
let m = MessageType::SessionJobs(SessionJobsPayload {
request_id: "msg_x".into(),
jobs: vec![JobListEntry {
job_id: crate::ids::JobId::new(),
agent: "echo@1.0.0".into(),
status: "running".into(),
parent_job_id: None,
created_at: now,
trace_id: None,
last_event_seq: 0,
}],
next_cursor: None,
});
let json = serde_json::to_string(&m).expect("serialize");
let back: MessageType = serde_json::from_str(&json).expect("deserialize");
assert_eq!(m, back);
}
#[test]
fn countable_event_classification_excludes_session_control() {
let now = chrono::Utc::now();
assert!(!MessageType::SessionPing(SessionPingPayload {
nonce: "n".into(),
sent_at: now,
})
.is_countable_event());
assert!(!MessageType::SessionAck(SessionAckPayload {
last_processed_seq: 0,
})
.is_countable_event());
assert!(!MessageType::Ping(PingPayload::default()).is_countable_event());
assert!(
MessageType::JobAccepted(crate::messages::JobAcceptedPayload {
job_id: crate::ids::JobId::new(),
credentials: vec![],
lease: None,
})
.is_countable_event()
);
}
#[test]
fn session_pong_round_trips_through_serde() {
let now = chrono::Utc::now();
let m = MessageType::SessionPong(SessionPongPayload {
ping_nonce: "p_01J".into(),
received_at: now,
});
let json = serde_json::to_string(&m).expect("serialize");
let back: MessageType = serde_json::from_str(&json).expect("deserialize");
assert_eq!(m, back);
}
#[test]
fn ping_wire_shape_matches_rfc() {
let m = MessageType::Ping(PingPayload { nonce: None });
let json = serde_json::to_value(&m).expect("serialize");
assert_eq!(json, serde_json::json!({"type": "ping", "payload": {}}));
}
#[test]
fn metric_with_standard_name_round_trips() {
let m = MessageType::Metric(MetricPayload {
name: standard_names::TOKENS_USED.into(),
value: 1432.0,
unit: "tokens".into(),
dims: Some(serde_json::json!({"model": "claude-3.5", "kind": "input"})),
});
let json = serde_json::to_string(&m).expect("serialize");
let back: MessageType = serde_json::from_str(&json).expect("deserialize");
assert_eq!(m, back);
}
#[test]
fn type_name_matches_wire_discriminator() {
let cases = [
(MessageType::Ping(PingPayload::default()), "ping"),
(MessageType::Pong(PongPayload::default()), "pong"),
(
MessageType::EventEmit(EventEmitPayload {
name: "x".into(),
data: None,
}),
"event.emit",
),
];
for (m, expected) in cases {
assert_eq!(m.type_name(), expected);
}
}
#[test]
fn capabilities_default_is_empty() {
let c = Capabilities::default();
assert!(!c.has(CapabilityName::Streaming));
assert!(!c.has(CapabilityName::Anonymous));
}
#[test]
fn capabilities_with_flat_agents_round_trips() {
let json = serde_json::json!({
"agents": ["code-refactor", "web-research"],
});
let c: Capabilities = serde_json::from_value(json.clone()).expect("deserialize");
match c.agents.as_ref().expect("agents present") {
AgentInventory::Flat(names) => {
assert_eq!(
names,
&vec!["code-refactor".to_owned(), "web-research".into()]
);
}
AgentInventory::Rich(_) => panic!("expected flat shape"),
}
let re = serde_json::to_value(&c).expect("serialize");
assert_eq!(re["agents"], json["agents"]);
}
#[test]
fn capabilities_with_rich_agents_round_trips() {
let json = serde_json::json!({
"agents": [
{ "name": "code-refactor", "versions": ["1.0.0", "2.0.0"], "default": "2.0.0" },
{ "name": "indexer", "versions": ["0.9.0"] },
],
});
let c: Capabilities = serde_json::from_value(json).expect("deserialize");
let inv = c.agents.as_ref().expect("agents present");
match inv {
AgentInventory::Rich(entries) => {
assert_eq!(entries.len(), 2);
assert_eq!(entries[0].name, "code-refactor");
assert_eq!(
entries[0].versions,
vec!["1.0.0".to_owned(), "2.0.0".into()]
);
assert_eq!(entries[0].default.as_deref(), Some("2.0.0"));
}
AgentInventory::Flat(_) => panic!("expected rich shape"),
}
}
#[test]
fn agent_inventory_satisfies_resolution_rules() {
let flat = AgentInventory::Flat(vec!["echo".into()]);
assert!(flat.satisfies(&crate::messages::AgentRef::parse("echo").expect("parse")));
assert!(flat.satisfies(&crate::messages::AgentRef::parse("echo@1.0.0").expect("parse")));
let rich = AgentInventory::Rich(vec![AgentInventoryEntry {
name: "echo".into(),
versions: vec!["1.0.0".into(), "2.0.0".into()],
default: Some("2.0.0".into()),
}]);
assert!(rich.satisfies(&crate::messages::AgentRef::parse("echo").expect("parse")));
assert!(rich.satisfies(&crate::messages::AgentRef::parse("echo@1.0.0").expect("parse")));
assert!(!rich.satisfies(&crate::messages::AgentRef::parse("echo@9.9.9").expect("parse")));
assert!(!rich.satisfies(&crate::messages::AgentRef::parse("other").expect("parse")));
}
#[test]
fn capabilities_round_trip_with_extra_fields() {
let json = serde_json::json!({
"streaming": true,
"extensions": ["arcpx.example.v1"],
"totally_made_up": true,
});
let c: Capabilities = serde_json::from_value(json).expect("deserialize");
assert!(c.has(CapabilityName::Streaming));
assert_eq!(c.extensions, vec!["arcpx.example.v1"]);
assert!(c.extra.contains_key("totally_made_up"));
}
#[test]
fn unknown_type_fails_deserialize() {
let bad = "{\"type\":\"never.heard.of.it\",\"payload\":{}}";
let result: Result<MessageType, _> = serde_json::from_str(bad);
assert!(result.is_err());
}
#[test]
fn handshake_messages_classified_correctly() {
assert!(MessageType::SessionOpen(SessionOpenPayload {
auth: Credentials {
scheme: AuthScheme::None,
token: None
},
client: ClientIdentity {
kind: "test".into(),
version: "0".into(),
fingerprint: None,
principal: None
},
capabilities: Capabilities::default(),
})
.is_handshake());
assert!(!MessageType::Ping(PingPayload::default()).is_handshake());
}
#[test]
#[allow(clippy::too_many_lines)]
fn type_name_covers_every_variant() {
let now = chrono::Utc::now();
let cases: Vec<(MessageType, &'static str)> = vec![
(
MessageType::SessionOpen(SessionOpenPayload {
auth: Credentials {
scheme: AuthScheme::None,
token: None,
},
client: ClientIdentity {
kind: "t".into(),
version: "0".into(),
fingerprint: None,
principal: None,
},
capabilities: Capabilities::default(),
}),
"session.open",
),
(
MessageType::SessionChallenge(SessionChallengePayload {
challenge: "x".into(),
}),
"session.challenge",
),
(
MessageType::SessionAuthenticate(SessionAuthenticatePayload {
response: "x".into(),
}),
"session.authenticate",
),
(
MessageType::SessionAccepted(SessionAcceptedPayload {
session_id: crate::ids::SessionId::new(),
runtime: crate::messages::RuntimeIdentity {
kind: "rt".into(),
version: "0".into(),
fingerprint: None,
trust_level: None,
},
capabilities: Capabilities::default(),
lease: None,
}),
"session.accepted",
),
(
MessageType::SessionUnauthenticated(SessionUnauthenticatedPayload {
code: crate::error::ErrorCode::Unauthenticated,
message: "x".into(),
}),
"session.unauthenticated",
),
(
MessageType::SessionRejected(SessionRejectedPayload {
code: crate::error::ErrorCode::Unauthenticated,
message: "x".into(),
}),
"session.rejected",
),
(
MessageType::SessionRefresh(SessionRefreshPayload {
deadline: now,
challenge: None,
}),
"session.refresh",
),
(
MessageType::SessionEvicted(SessionEvictedPayload {
code: crate::error::ErrorCode::Cancelled,
reason: "x".into(),
}),
"session.evicted",
),
(
MessageType::SessionClose(SessionClosePayload::default()),
"session.close",
),
(
MessageType::SessionPing(SessionPingPayload {
nonce: "n".into(),
sent_at: now,
}),
"session.ping",
),
(
MessageType::SessionPong(SessionPongPayload {
ping_nonce: "n".into(),
received_at: now,
}),
"session.pong",
),
(
MessageType::SessionAck(SessionAckPayload {
last_processed_seq: 0,
}),
"session.ack",
),
(
MessageType::SessionListJobs(SessionListJobsPayload::default()),
"session.list_jobs",
),
(
MessageType::SessionJobs(SessionJobsPayload {
request_id: "r".into(),
jobs: vec![],
next_cursor: None,
}),
"session.jobs",
),
(MessageType::Ping(PingPayload::default()), "ping"),
(MessageType::Pong(PongPayload::default()), "pong"),
(MessageType::Ack(AckPayload { note: None }), "ack"),
(
MessageType::Nack(NackPayload {
code: crate::error::ErrorCode::Unknown,
message: "x".into(),
details: None,
}),
"nack",
),
(
MessageType::Cancel(CancelPayload {
target: CancelTargetKind::Job,
target_id: "x".into(),
reason: None,
deadline_ms: None,
}),
"cancel",
),
(
MessageType::CancelAccepted(CancelAcceptedPayload { target_id: None }),
"cancel.accepted",
),
(
MessageType::CancelRefused(CancelRefusedPayload {
target_id: "x".into(),
reason: "x".into(),
}),
"cancel.refused",
),
(
MessageType::Interrupt(InterruptPayload {
target: CancelTargetKind::Job,
target_id: "x".into(),
prompt: "x".into(),
}),
"interrupt",
),
(MessageType::Resume(ResumePayload::default()), "resume"),
(
MessageType::Backpressure(BackpressurePayload {
desired_rate_per_second: None,
buffer_remaining_bytes: None,
reason: None,
}),
"backpressure",
),
(
MessageType::ToolInvoke(ToolInvokePayload::new("x", serde_json::json!({}))),
"tool.invoke",
),
(
MessageType::ToolResult(ToolResultPayload {
value: None,
result_ref: None,
}),
"tool.result",
),
(
MessageType::ToolError(ToolErrorPayload {
code: crate::error::ErrorCode::Internal,
retryable: None,
message: "x".into(),
details: None,
}),
"tool.error",
),
(
MessageType::JobAccepted(JobAcceptedPayload {
job_id: crate::ids::JobId::new(),
credentials: vec![],
lease: None,
}),
"job.accepted",
),
(
MessageType::JobStarted(JobStartedPayload { description: None }),
"job.started",
),
(
MessageType::JobProgress(JobProgressPayload {
percent: None,
message: None,
}),
"job.progress",
),
(
MessageType::JobHeartbeat(JobHeartbeatPayload {
sequence: 1,
deadline_ms: None,
state: JobState::Running,
}),
"job.heartbeat",
),
(
MessageType::JobCompleted(JobCompletedPayload {
value: None,
result_ref: None,
result_id: None,
result_size: None,
summary: None,
}),
"job.completed",
),
(
MessageType::JobResultChunk(JobResultChunkPayload {
result_id: "r".into(),
chunk_seq: 0,
data: "x".into(),
encoding: crate::messages::ResultChunkEncoding::Utf8,
more: false,
}),
"job.result_chunk",
),
(
MessageType::JobFailed(JobFailedPayload {
code: crate::error::ErrorCode::Internal,
retryable: None,
message: "x".into(),
details: None,
}),
"job.failed",
),
(
MessageType::JobCancelled(JobCancelledPayload { reason: None }),
"job.cancelled",
),
(
MessageType::StreamOpen(StreamOpenPayload {
kind: StreamKind::Text,
content_type: None,
encoding: None,
}),
"stream.open",
),
(
MessageType::StreamChunk(StreamChunkPayload {
sequence: 0,
data: serde_json::json!(""),
content_type: None,
sha256: None,
redacted: false,
role: None,
}),
"stream.chunk",
),
(
MessageType::StreamClose(StreamClosePayload::default()),
"stream.close",
),
(
MessageType::StreamError(StreamErrorPayload {
code: crate::error::ErrorCode::Internal,
message: "x".into(),
}),
"stream.error",
),
(
MessageType::PermissionRequest(PermissionRequestPayload {
permission: "p".into(),
resource: "r".into(),
operation: "o".into(),
reason: None,
requested_lease_seconds: None,
}),
"permission.request",
),
(
MessageType::PermissionGrant(PermissionGrantPayload {
permission: "p".into(),
resource: "r".into(),
operation: "o".into(),
lease_seconds: 1,
}),
"permission.grant",
),
(
MessageType::PermissionDeny(PermissionDenyPayload {
permission: "p".into(),
reason: "x".into(),
}),
"permission.deny",
),
(
MessageType::LeaseGranted(LeaseGrantedPayload {
lease_id: crate::ids::LeaseId::new(),
permission: "p".into(),
resource: "r".into(),
operation: "o".into(),
expires_at: now,
}),
"lease.granted",
),
(
MessageType::LeaseExtended(LeaseExtendedPayload {
lease_id: crate::ids::LeaseId::new(),
expires_at: now,
}),
"lease.extended",
),
(
MessageType::LeaseRevoked(LeaseRevokedPayload {
lease_id: crate::ids::LeaseId::new(),
reason: "x".into(),
}),
"lease.revoked",
),
(
MessageType::LeaseRefresh(LeaseRefreshPayload {
lease_id: crate::ids::LeaseId::new(),
additional_seconds: 1,
}),
"lease.refresh",
),
(
MessageType::Subscribe(SubscribePayload::default()),
"subscribe",
),
(
MessageType::SubscribeAccepted(SubscribeAcceptedPayload {
subscription_id: crate::ids::SubscriptionId::new(),
}),
"subscribe.accepted",
),
(
MessageType::SubscribeEvent(SubscribeEventPayload {
event: serde_json::json!({}),
}),
"subscribe.event",
),
(
MessageType::Unsubscribe(UnsubscribePayload {
subscription_id: crate::ids::SubscriptionId::new(),
}),
"unsubscribe",
),
(
MessageType::SubscribeClosed(SubscribeClosedPayload {
subscription_id: crate::ids::SubscriptionId::new(),
code: crate::error::ErrorCode::Cancelled,
reason: "x".into(),
}),
"subscribe.closed",
),
(
MessageType::JobSubscribe(JobSubscribePayload {
job_id: crate::ids::JobId::new(),
from_event_seq: None,
history: false,
}),
"job.subscribe",
),
(
MessageType::JobSubscribed(JobSubscribedPayload {
job_id: crate::ids::JobId::new(),
current_status: "running".into(),
agent: "echo".into(),
parent_job_id: None,
trace_id: None,
subscribed_from: 0,
replayed: false,
}),
"job.subscribed",
),
(
MessageType::JobUnsubscribe(JobUnsubscribePayload {
job_id: crate::ids::JobId::new(),
}),
"job.unsubscribe",
),
(
MessageType::ArtifactPut(ArtifactPutPayload {
media_type: "x".into(),
data: String::new(),
sha256: None,
retain_seconds: None,
}),
"artifact.put",
),
(
MessageType::ArtifactFetch(ArtifactFetchPayload {
artifact_id: crate::ids::ArtifactId::new(),
}),
"artifact.fetch",
),
(
MessageType::ArtifactRef(ArtifactRefPayload {
artifact: ArtifactRef {
artifact_id: crate::ids::ArtifactId::new(),
uri: "arcp://x".into(),
media_type: "x".into(),
size: 0,
sha256: None,
expires_at: None,
},
}),
"artifact.ref",
),
(
MessageType::ArtifactRelease(ArtifactReleasePayload {
artifact_id: crate::ids::ArtifactId::new(),
}),
"artifact.release",
),
(
MessageType::EventEmit(EventEmitPayload {
name: "x".into(),
data: None,
}),
"event.emit",
),
(
MessageType::Log(LogPayload {
level: LogLevel::Info,
message: "x".into(),
attributes: None,
}),
"log",
),
(
MessageType::Metric(MetricPayload {
name: "x".into(),
value: 0.0,
unit: "u".into(),
dims: None,
}),
"metric",
),
(
MessageType::TraceSpan(TraceSpanPayload {
name: "x".into(),
trace_id: crate::ids::TraceId::new("t").expect("non-empty"),
span_id: crate::ids::SpanId::new("s").expect("non-empty"),
parent_span_id: None,
start_time: now,
end_time: now,
attributes: None,
}),
"trace.span",
),
];
for (msg, expected) in &cases {
assert_eq!(msg.type_name(), *expected);
}
assert_eq!(cases.len(), 62);
}
}