feat(rust): add authenticated node sidecar bridge (#116863)

Adds the bounded authenticated sidecar bridge on the merged Rust node runtime. The unique sidecar delta passed focused Rust and process proof, and exact-head ClawSweeper review found no actionable code or security issue.
This commit is contained in:
Gio Della-Libera 2026-09-16 11:30:47 -07:00 • committed by GitHub
parent b2883dfac6
commit ce4f1d711b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
13 changed files with 5371 additions and 2 deletions

27
crates/Cargo.lock generated
View file

@ -56,6 +56,12 @@ dependencies = [
"rand_core",
]
[[package]]
name = "cmov"
version = "0.5.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0c9ea0ac24bc397ab3c98583a3c9ba74fa56b09a4449bbe172b9b1ddb016027a"
[[package]]
name = "const-oid"
version = "0.10.2"
@ -96,6 +102,15 @@ dependencies = [
"hybrid-array",
]
[[package]]
name = "ctutils"
version = "0.4.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7d5515a3834141de9eafb9717ad39eea8247b5674e6066c404e8c4b365d2a29e"
dependencies = [
"cmov",
]
[[package]]
name = "curve25519-dalek"
version = "5.0.0"
@ -138,6 +153,7 @@ dependencies = [
"block-buffer",
"const-oid",
"crypto-common",
"ctutils",
]
[[package]]
@ -270,6 +286,15 @@ dependencies = [
"rand_core",
]
[[package]]
name = "hmac"
version = "0.13.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6303bc9732ae41b04cb554b844a762b4115a61bfaa81e3e83050991eeb56863f"
dependencies = [
"digest",
]
[[package]]
name = "http"
version = "1.5.0"
@ -469,6 +494,7 @@ dependencies = [
"ed25519-dalek",
"futures-util",
"getrandom 0.4.3",
"hmac",
"openclaw-gateway-client",
"serde",
"serde_json",
@ -476,6 +502,7 @@ dependencies = [
"thiserror",
"tokio",
"tokio-tungstenite",
"zeroize",
]
[[package]]

View file

@ -12,6 +12,7 @@ futures-util = "0.3.32"
base64 = "0.22.1"
ed25519-dalek = "3.0.0"
getrandom = "0.4.3"
hmac = "0.13.0"
openclaw-gateway-client = { path = "../openclaw-gateway-client" }
serde = { version = "1.0.228", features = ["derive"] }
serde_json = "1.0.150"
@ -26,6 +27,7 @@ tokio = { version = "1.52.3", features = [
"sync",
"time",
] }
zeroize = "1.8.2"
[lints.rust]
unsafe_code = "forbid"

View file

@ -0,0 +1,229 @@
# OpenClaw node sidecar protocol v1
This document describes the byte contract implemented by
`sidecar_protocol.rs`. It is the portable authenticated channel beneath a
future node-runtime message schema. It does not select named pipes, Unix
sockets, inherited anonymous pipes, or another local transport.
## Trust before this protocol
The product supervisor must verify the exact runtime artifact and create a
fresh local-only IPC endpoint before launch. It supplies a random 32-byte
session key, session identifier, and nonzero generation over a protected
bootstrap mechanism. Those values must not be placed in command-line
arguments, broadly inherited environment variables, logs, crash reports, or
world-readable files.
This crate intentionally does not implement platform artifact verification,
secret delivery, process creation, secure storage, or runtime selection. A
peer identity reported inside this protocol is authenticated by the session
key but is not proof that the executable on disk was trusted.
Gateway endpoint selection, challenge signing, node-protocol negotiation, and
device-token persistence stay with the platform-owned `NodeLifecycle`
connection factory. In particular, an `IssuedDeviceToken` is never placed in
sidecar configuration or status traffic. Its attempt number can be correlated
with the secret-free lifecycle `attempt` status, but delivery and durable
acknowledgement require a separate protected product channel outside this
protocol.
## Length prefix and local ceiling
Each frame is preceded by an unsigned four-byte big-endian length. The length
counts the authenticated frame and excludes the prefix. A receiver must apply
its local hard ceiling to the prefix before allocating or reading the frame.
The prefix is not security authority: all internal lengths and fields are
covered by the authentication tag.
The bootstrap exchange always uses protocol minor `0`, the same
pre-negotiation ceiling, and a finite deadline. This lets an older peer read a
newer peer's offer before minor-version negotiation. Negotiated limits are the
minimum of both valid local offers and can never raise a local ceiling. After
the final bootstrap frame, both peers apply the independently verified minor
and frame limit to the active authenticated channel so directional sequence
numbers are not reset.
## Authenticated frame
All integers are unsigned and big-endian.
| Field | Bytes | Meaning |
| ------------------ | -------: | ---------------------------------------------------- |
| Magic | 4 | ASCII `OCSC` |
| Protocol major | 2 | `1` |
| Protocol minor | 2 | `0` during bootstrap; negotiated minor afterward |
| Direction | 1 | `1` supervisor-to-runtime, `2` runtime-to-supervisor |
| Generation | 8 | Nonzero process/session generation |
| Sequence | 8 | Strictly increasing per direction, starting at `1` |
| Session ID length | 2 | UTF-8 byte length |
| Payload length | 4 | JSON payload byte length |
| Session ID | variable | Exact bootstrap session identifier |
| Payload | variable | UTF-8 JSON; the next slice defines typed messages |
| Authentication tag | 32 | HMAC-SHA-256 over every preceding frame byte |
The frame limit includes the authentication tag and excludes the outer length
prefix. The sender and receiver use the same session key; direction is part of
the authenticated header and prevents reflection between peers.
Outbound JSON is serialized directly into the final frame through the local
ceiling. Serialization stops on the first write that would consume the bytes
reserved for the authentication tag; an oversized payload is never fully
materialized or copied into a second plaintext buffer.
A receiver verifies the frame-size ceiling and HMAC before interpreting any
untrusted header or payload field. It then verifies version, direction,
generation, exact next sequence, session identifier, internal lengths, and
payload decoding. Any failure retires the channel; callers must not continue
after a framing, authentication, replay, or generation error. The Rust channel
poisons itself on the first inbound validation failure and rejects every later
send or receive. A transport owner calls `retire()` when length-prefix I/O or
its surrounding IPC transport fails.
A new process gets a new session identifier, key, generation, and sequence
space. A sequence gap or replay is rejected rather than buffered. Rotate the
generation before sequence exhaustion; never reset a sequence in place.
The authenticated channel generation, immutable manifest generation,
`NodeLifecycle` connection attempt, and Gateway pairing generation are separate
scopes. This protocol carries only the first three. Gateway pairing approval,
committed-policy reconciliation, and pairing-generation leases remain Gateway
authority and must not be inferred from a sidecar generation or manifest.
## Negotiation
`SidecarProtocolOffer` carries the peer role and reported identity, protocol
version, additive feature bits, frame/in-flight ceilings, and bootstrap
deadline. Peers must have complementary roles and the same major version.
Features are intersected. The feature mask must be at most `2^53 - 1`, the
largest integer that JSON implementations such as JavaScript can preserve
exactly; larger local or remote offers are invalid. Frame, in-flight, and
deadline values use the lower valid offer. Unknown features remain disabled.
Adapters must perform the intersection with integer arithmetic that preserves
all 53 bits. JavaScript and TypeScript implementations must convert both masks
to `BigInt` before `&` and convert the bounded result back to `Number`; their
native number `&` operator truncates operands to 32 bits and is not conformant.
The supervisor initiates with an authenticated `offer`. The runtime
independently negotiates against its local offer and replies with one
authenticated `accept` containing its offer and the selected parameters. The
supervisor independently recomputes the selection; a mismatch is terminal.
The runtime remains in `AcceptancePending` and retains the bootstrap ceiling
until the acceptance is written successfully. It then commits the selection;
the supervisor commits only after receiving and verifying that frame. This
prevents active traffic from racing ahead of the acceptance and ensures an
acceptance larger than the negotiated ceiling can still be delivered. Both
peers preserve the existing directional sequence state.
The `SidecarHandshake` state machine is bound to one exact authenticated
channel instance and accepts only this two-frame ordering. Substituting another
channel is rejected before frame processing or handshake mutation and retires
the supplied replacement; the original bound handshake can continue.
Wrong roles, incompatible versions, malformed/authentication failures, forged
selection, repeated/out-of-order messages, and unencodable bootstrap messages
retire the handshake and channel. Negotiation must complete before accepting
credentials, configuration, capability registration, or invocation traffic.
## Cross-language vector
[`node-sidecar-protocol-v1.json`](../../test/fixtures/node-sidecar-protocol-v1.json)
contains a test-only session key, payload, and exact encoded data frame.
[`node-sidecar-negotiation-v1.json`](../../test/fixtures/node-sidecar-negotiation-v1.json)
exercises a feature bit above the 32-bit JavaScript bitwise range.
[`node-sidecar-handshake-v1.json`](../../test/fixtures/node-sidecar-handshake-v1.json)
contains both offers, the independently derived selection, and exact offer and
accept frames. Rust tests reproduce and decode every vector. Every non-Rust
adapter must consume the same vectors before it can be selected as a runtime.
[`node-sidecar-runtime-v1.json`](../../test/fixtures/node-sidecar-runtime-v1.json)
contains the typed configuration, configured acknowledgement, admission,
invocation, result, cancellation, and status messages plus their exact compact
JSON encodings. These payloads travel inside the already authenticated,
sequenced frames; the message corpus does not replace the frame vectors.
All integer fields, including integers nested inside invocation parameters or
success payloads, must remain in JSON's exact `-(2^53 - 1)..=2^53 - 1` range.
Serialization and deserialization reject values outside that range; the shared
corpus exercises the positive boundary. Integer-valued decimal or exponent
forms follow the same bound; genuine fractional JSON numbers retain their
normal finite IEEE-754 semantics. Runtime adapter traffic reports the
distinct `SIDECAR_NON_PORTABLE_JSON` failure rather than misclassifying these
values as oversized payloads.
## Runtime bridge
`SidecarRuntimeBridge` can be constructed only from the runtime side after the
validated configuration acknowledgement has been written successfully. The
runtime remains `AcknowledgementPending` until that delivery is committed, so
invocation work cannot race ahead of the supervisor-visible manifest. The
configuration exchange is consumed from its authenticated handshake, and a
successful activation irreversibly moves it to `Activated`; neither phase can
be replayed to multiply concurrency. Beginning configuration locks the
negotiated frame ceiling for the rest of the channel generation, so bridge
preflight budgets cannot become stale. Activation also requires the exact live
authenticated channel. The bridge carries that channel's retirement signal:
retirement blocks new native work and cancels in-flight adapter work. The bridge
accepts one immutable connection manifest and requires its
concurrency/input/output limits to remain within the negotiated sidecar envelope. Command and capability names
use the shared ASCII grammar `[A-Za-z0-9._-]{1,128}`, are duplicate-free, and
are sorted bytewise before acknowledgement; the OpenClaw-owned `system.*`
namespace remains reserved.
The configured output limit must also fit the largest bridge-owned stable
failure envelope. This preserves `SIDECAR_MESSAGE_TOO_LARGE`,
`SIDECAR_NON_PORTABLE_JSON`, and `SIDECAR_CHANNEL_RETIRED` instead of allowing
the generic command-runtime output limiter to rewrite them.
`SidecarConfigurationExchange` permits exactly one supervisor configuration
followed by the runtime's acknowledgement of the independently derived
manifest. It remains bound to the exact channel instance authenticated by the
handshake; a replacement is retired before processing without mutating the
bound exchange. Wrong order, channel role, malformed or unknown fields,
invalid limits/names, and a forged acknowledgement retire the exchange and channel.
It also proves the worst-case secret-free status envelope—including the runtime
version from the authenticated offer—fits the lowered live-channel budget. No
admission or invocation message is accepted by this exchange.
The bridge builds the existing bounded `CommandRuntime`. Every invocation
therefore passes its normal Gateway-manifest, input, concurrency, timeout,
admission, handler, output, and cancellation gates. The product-owned
`SidecarCapabilityAdapter` receives an untrusted typed invocation first for
local admission and only then for native dispatch. A denial or adapter failure
becomes a bounded structured handler failure. The runtime cancellation token is
passed through both waits so product adapter work can stop on Gateway cancel,
timeout, disconnect, or shutdown.
The configured command list is an offered connection manifest, not an approval
or policy grant. The Gateway decides whether a declared command is invocable;
only a Gateway-delivered invocation reaches the bridge's additional
fail-closed local admission. Surface widening requires a newly configured
bridge and a new Gateway connection manifest. Committed policy revocation or a
pairing-generation transition can preserve the physical Gateway connection
while cancelling generation-bound work; that cancellation reaches the adapter
through the existing runtime token rather than a new sidecar wire field.
Logical parameter and result limits are not treated as complete-frame limits.
Before cloning Gateway JSON into an adapter request, and without cloning an
adapter decision/result, the bridge runs a borrowed non-allocating serialization
preflight for the complete admission, invocation, decision, or result message
against the live channel's exact payload budget
(frame ceiling minus fixed header, session identifier, and authentication tag).
Adapter infrastructure errors are first normalized into the same denial or
failure wire shape, so they cannot bypass the complete-message preflight.
A value that fits its logical JSON limit but not its full envelope receives the
stable `SIDECAR_MESSAGE_TOO_LARGE` result without attempting transport.
The bridge also maps secret-free `NodeLifecycle` events into stable status
states and reasons correlated with the immutable manifest generation. A
capability change requires a new bridge and connection/process generation;
registrations cannot be mutated in place after advertisement.
## Not yet implemented
This stack still excludes product IPC selection, protected credential/config
delivery, duplex input/progress transport, a product audit adapter, process
supervision, artifact verification, packaging, Windows integration,
restart/rollback policy, production credential storage, worker/session
hosting, workspace transfer, plugins, host statistics, `system.run`, PTY, MCP,
and skills. The typed runtime messages and adapter boundary do not by
themselves claim production sidecar readiness. Those remaining concerns must
not weaken authentication, bounds, generation, sequencing, cancellation, or
fail-closed behavior.

View file

@ -14,6 +14,9 @@ mod lifecycle;
mod node;
mod reconnect;
mod runtime;
mod sidecar_handshake;
mod sidecar_protocol;
mod sidecar_runtime;
pub use duplex::InvocationIo;
pub use host::{run_host, AuthKind, HostConfig, HostCredentials, HostError};
@ -35,3 +38,21 @@ pub use runtime::{
CancellationToken, CommandRuntime, CommandRuntimeBuilder, HandlerError,
InvocationAdmissionContext, InvocationContext, RuntimeBuildError, RuntimeError,
};
pub use sidecar_handshake::{
SidecarHandshake, SidecarHandshakeError, SidecarHandshakeMessage, SidecarHandshakeState,
SidecarProtocolSelection,
};
pub use sidecar_protocol::{
negotiate_sidecar_protocol, read_sidecar_frame, write_sidecar_frame,
AuthenticatedSidecarChannel, NegotiatedSidecarProtocol, SidecarDirection, SidecarFrameError,
SidecarLimits, SidecarPeerIdentity, SidecarPeerRole, SidecarProtocolError,
SidecarProtocolOffer, SidecarSessionKey, SIDECAR_MAX_FEATURE_BITS, SIDECAR_PROTOCOL_MAJOR,
SIDECAR_PROTOCOL_MINOR,
};
pub use sidecar_runtime::{
SidecarAdapterError, SidecarAdapterFuture, SidecarAdmissionDecision, SidecarCapabilityAdapter,
SidecarCommandRegistration, SidecarConfigurationError, SidecarConfigurationExchange,
SidecarConfigurationState, SidecarInvocation, SidecarInvocationResult, SidecarRuntimeBridge,
SidecarRuntimeBridgeError, SidecarRuntimeConfiguration, SidecarRuntimeManifest,
SidecarRuntimeMessage, SidecarRuntimeReason, SidecarRuntimeState, SidecarRuntimeStatus,
};

View file

@ -0,0 +1,796 @@
//! Authenticated offer/accept state machine for the node sidecar protocol.
use serde::{Deserialize, Serialize};
use thiserror::Error;
use crate::{
negotiate_sidecar_protocol, AuthenticatedSidecarChannel, NegotiatedSidecarProtocol,
SidecarFrameError, SidecarLimits, SidecarPeerRole, SidecarProtocolError, SidecarProtocolOffer,
};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum SidecarHandshakeState {
Starting,
AwaitingAcceptance,
AcceptancePending,
Authenticated,
Failed,
}
#[derive(Clone, Copy, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(rename_all = "camelCase")]
pub struct SidecarProtocolSelection {
pub protocol_major: u16,
pub protocol_minor: u16,
pub feature_bits: u64,
pub limits: SidecarLimits,
}
impl From<&NegotiatedSidecarProtocol> for SidecarProtocolSelection {
fn from(negotiated: &NegotiatedSidecarProtocol) -> Self {
Self {
protocol_major: negotiated.protocol_major,
protocol_minor: negotiated.protocol_minor,
feature_bits: negotiated.feature_bits,
limits: negotiated.limits,
}
}
}
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
#[serde(tag = "type", rename_all = "kebab-case")]
pub enum SidecarHandshakeMessage {
Offer {
offer: SidecarProtocolOffer,
},
Accept {
offer: SidecarProtocolOffer,
selection: SidecarProtocolSelection,
},
}
/// Drives the two-frame authenticated protocol handshake for one peer.
pub struct SidecarHandshake {
local_offer: SidecarProtocolOffer,
state: SidecarHandshakeState,
negotiated: Option<NegotiatedSidecarProtocol>,
pending_negotiated: Option<NegotiatedSidecarProtocol>,
channel_instance_id: Option<u64>,
}
impl SidecarHandshake {
/// Create a handshake for a supervisor or runtime offer.
///
/// # Errors
///
/// Returns an error when the local offer is invalid or uses an unsupported
/// protocol version.
pub fn new(local_offer: SidecarProtocolOffer) -> Result<Self, SidecarHandshakeError> {
validate_local_offer(&local_offer)?;
Ok(Self {
local_offer,
state: SidecarHandshakeState::Starting,
negotiated: None,
pending_negotiated: None,
channel_instance_id: None,
})
}
#[must_use]
pub const fn state(&self) -> SidecarHandshakeState {
self.state
}
#[must_use]
pub const fn local_role(&self) -> SidecarPeerRole {
self.local_offer.peer.role
}
#[must_use]
pub const fn local_peer(&self) -> &crate::SidecarPeerIdentity {
&self.local_offer.peer
}
pub(crate) const fn bound_channel_instance_id(&self) -> Option<u64> {
self.channel_instance_id
}
#[must_use]
pub const fn negotiated(&self) -> Option<&NegotiatedSidecarProtocol> {
self.negotiated.as_ref()
}
/// Encode the supervisor's initial authenticated offer.
///
/// # Errors
///
/// Returns an error when called by a runtime, called out of order, or when
/// the offer cannot be encoded within the bootstrap frame ceiling.
pub fn start(
&mut self,
channel: &mut AuthenticatedSidecarChannel,
) -> Result<Vec<u8>, SidecarHandshakeError> {
if let Err(error) = self.bind_channel(channel) {
channel.retire();
return Err(error);
}
if channel.role() != self.local_offer.peer.role {
return self.fail(channel, SidecarHandshakeError::ChannelRoleMismatch);
}
if self.local_offer.peer.role != SidecarPeerRole::Supervisor {
return self.fail(channel, SidecarHandshakeError::SupervisorMustInitiate);
}
if self.state != SidecarHandshakeState::Starting {
return self.fail(channel, SidecarHandshakeError::UnexpectedMessage);
}
let frame = match channel.seal(&SidecarHandshakeMessage::Offer {
offer: self.local_offer.clone(),
}) {
Ok(frame) => frame,
Err(error) => return self.fail(channel, SidecarHandshakeError::Frame(error)),
};
self.state = SidecarHandshakeState::AwaitingAcceptance;
Ok(frame)
}
/// Consume one authenticated handshake frame and optionally return the
/// runtime's authenticated acceptance frame.
///
/// # Errors
///
/// Returns an error for framing/authentication failure, incompatible or
/// forged negotiation, wrong peer roles, or invalid message ordering.
/// Frame and state-machine errors permanently retire the bound channel and
/// handshake. A different channel instance is rejected before processing:
/// only that supplied replacement is retired, and the bound handshake is
/// unchanged.
pub fn receive(
&mut self,
channel: &mut AuthenticatedSidecarChannel,
frame: &[u8],
) -> Result<Option<Vec<u8>>, SidecarHandshakeError> {
if let Err(error) = self.bind_channel(channel) {
channel.retire();
return Err(error);
}
if channel.role() != self.local_offer.peer.role {
return self.fail(channel, SidecarHandshakeError::ChannelRoleMismatch);
}
let result = self.receive_active(channel, frame);
if result.is_err() {
channel.retire();
self.state = SidecarHandshakeState::Failed;
self.negotiated = None;
self.pending_negotiated = None;
}
result
}
/// Commit the runtime's negotiated channel state after its acceptance
/// frame has been written successfully using the bootstrap ceiling.
///
/// The runtime remains in [`SidecarHandshakeState::AcceptancePending`]
/// until this method succeeds, so active traffic cannot race ahead of the
/// final bootstrap frame. On transport failure, retire the channel before
/// calling this method; completion then fails terminally.
///
/// # Errors
///
/// Returns an error for the wrong channel, role, state, a retired channel,
/// or an invalid negotiated selection. State and negotiation errors retire
/// the bound channel and handshake. A different channel instance is
/// rejected before processing and retires only the supplied replacement.
pub fn complete_acceptance(
&mut self,
channel: &mut AuthenticatedSidecarChannel,
) -> Result<(), SidecarHandshakeError> {
if let Err(error) = self.bind_channel(channel) {
channel.retire();
return Err(error);
}
if channel.role() != SidecarPeerRole::Runtime {
return self.fail(channel, SidecarHandshakeError::ChannelRoleMismatch);
}
if channel.is_retired() {
return self.fail(
channel,
SidecarHandshakeError::Frame(SidecarFrameError::ChannelRetired),
);
}
if self.state != SidecarHandshakeState::AcceptancePending {
return self.fail(channel, SidecarHandshakeError::UnexpectedMessage);
}
let Some(negotiated) = self.pending_negotiated.take() else {
return self.fail(channel, SidecarHandshakeError::UnexpectedMessage);
};
if let Err(error) = channel.apply_negotiated_protocol(&negotiated) {
return self.fail(channel, SidecarHandshakeError::Negotiation(error));
}
self.negotiated = Some(negotiated);
self.state = SidecarHandshakeState::Authenticated;
Ok(())
}
fn receive_active(
&mut self,
channel: &mut AuthenticatedSidecarChannel,
frame: &[u8],
) -> Result<Option<Vec<u8>>, SidecarHandshakeError> {
let message = channel
.open::<SidecarHandshakeMessage>(frame)
.map_err(SidecarHandshakeError::Frame)?;
match (self.local_offer.peer.role, self.state, message) {
(
SidecarPeerRole::Runtime,
SidecarHandshakeState::Starting,
SidecarHandshakeMessage::Offer { offer },
) => self.accept_supervisor(channel, &offer).map(Some),
(
SidecarPeerRole::Supervisor,
SidecarHandshakeState::AwaitingAcceptance,
SidecarHandshakeMessage::Accept { offer, selection },
) => {
self.confirm_runtime(channel, &offer, selection)?;
Ok(None)
}
_ => Err(SidecarHandshakeError::UnexpectedMessage),
}
}
fn accept_supervisor(
&mut self,
channel: &mut AuthenticatedSidecarChannel,
supervisor_offer: &SidecarProtocolOffer,
) -> Result<Vec<u8>, SidecarHandshakeError> {
if supervisor_offer.peer.role != SidecarPeerRole::Supervisor {
return Err(SidecarHandshakeError::WrongPeerRole);
}
validate_peer_identity(supervisor_offer)?;
let negotiated = negotiate_sidecar_protocol(&self.local_offer, supervisor_offer)
.map_err(SidecarHandshakeError::Negotiation)?;
let selection = SidecarProtocolSelection::from(&negotiated);
// The acceptance is the last bootstrap-ceiling frame. Lowering before
// sealing it could make a valid negotiation impossible to acknowledge.
let acceptance = channel
.seal(&SidecarHandshakeMessage::Accept {
offer: self.local_offer.clone(),
selection,
})
.map_err(SidecarHandshakeError::Frame)?;
self.pending_negotiated = Some(negotiated);
self.state = SidecarHandshakeState::AcceptancePending;
Ok(acceptance)
}
fn confirm_runtime(
&mut self,
channel: &mut AuthenticatedSidecarChannel,
runtime_offer: &SidecarProtocolOffer,
claimed: SidecarProtocolSelection,
) -> Result<(), SidecarHandshakeError> {
if runtime_offer.peer.role != SidecarPeerRole::Runtime {
return Err(SidecarHandshakeError::WrongPeerRole);
}
validate_peer_identity(runtime_offer)?;
let negotiated = negotiate_sidecar_protocol(&self.local_offer, runtime_offer)
.map_err(SidecarHandshakeError::Negotiation)?;
if claimed != SidecarProtocolSelection::from(&negotiated) {
return Err(SidecarHandshakeError::SelectionMismatch);
}
channel
.apply_negotiated_protocol(&negotiated)
.map_err(SidecarHandshakeError::Negotiation)?;
self.negotiated = Some(negotiated);
self.state = SidecarHandshakeState::Authenticated;
Ok(())
}
fn fail<T>(
&mut self,
channel: &mut AuthenticatedSidecarChannel,
error: SidecarHandshakeError,
) -> Result<T, SidecarHandshakeError> {
channel.retire();
self.state = SidecarHandshakeState::Failed;
self.negotiated = None;
self.pending_negotiated = None;
Err(error)
}
fn bind_channel(
&mut self,
channel: &AuthenticatedSidecarChannel,
) -> Result<(), SidecarHandshakeError> {
let instance_id = channel.instance_id();
match self.channel_instance_id {
None => {
self.channel_instance_id = Some(instance_id);
Ok(())
}
Some(bound) if bound == instance_id => Ok(()),
Some(_) => Err(SidecarHandshakeError::ChannelInstanceMismatch),
}
}
}
fn validate_local_offer(offer: &SidecarProtocolOffer) -> Result<(), SidecarHandshakeError> {
validate_peer_identity(offer)?;
let counterpart = SidecarProtocolOffer {
protocol_major: offer.protocol_major,
protocol_minor: offer.protocol_minor,
peer: crate::SidecarPeerIdentity {
role: match offer.peer.role {
SidecarPeerRole::Supervisor => SidecarPeerRole::Runtime,
SidecarPeerRole::Runtime => SidecarPeerRole::Supervisor,
},
name: "validation-peer".into(),
version: "0".into(),
artifact_identity: "validation-only".into(),
},
feature_bits: offer.feature_bits,
limits: offer.limits,
};
negotiate_sidecar_protocol(offer, &counterpart)
.map(|_| ())
.map_err(SidecarHandshakeError::Negotiation)
}
fn validate_peer_identity(offer: &SidecarProtocolOffer) -> Result<(), SidecarHandshakeError> {
if offer.peer.name.trim().is_empty()
|| offer.peer.version.trim().is_empty()
|| offer.peer.artifact_identity.trim().is_empty()
{
return Err(SidecarHandshakeError::InvalidPeerIdentity);
}
Ok(())
}
#[derive(Debug, Error)]
pub enum SidecarHandshakeError {
#[error("sidecar handshake frame failed")]
Frame(#[source] SidecarFrameError),
#[error("sidecar protocol negotiation failed")]
Negotiation(#[source] SidecarProtocolError),
#[error("runtime cannot initiate the sidecar handshake")]
SupervisorMustInitiate,
#[error("sidecar acceptance does not match the independently negotiated selection")]
SelectionMismatch,
#[error("unexpected sidecar handshake message")]
UnexpectedMessage,
#[error("sidecar peer identity fields must be nonempty")]
InvalidPeerIdentity,
#[error("sidecar handshake message came from the wrong peer role")]
WrongPeerRole,
#[error("sidecar handshake role does not match the authenticated channel role")]
ChannelRoleMismatch,
#[error("sidecar handshake cannot move between authenticated channel instances")]
ChannelInstanceMismatch,
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{SidecarPeerIdentity, SidecarSessionKey};
use base64::{engine::general_purpose::STANDARD as BASE64, Engine as _};
const KEY: [u8; 32] = [0x3c; 32];
fn offer(
role: SidecarPeerRole,
max_frame_bytes: u32,
feature_bits: u64,
) -> SidecarProtocolOffer {
SidecarProtocolOffer {
protocol_major: crate::SIDECAR_PROTOCOL_MAJOR,
protocol_minor: crate::SIDECAR_PROTOCOL_MINOR,
peer: SidecarPeerIdentity {
role,
name: match role {
SidecarPeerRole::Supervisor => "test-product",
SidecarPeerRole::Runtime => "openclaw-node",
}
.into(),
version: "1.0.0".into(),
artifact_identity: "sha256:test-only".into(),
},
feature_bits,
limits: SidecarLimits {
max_frame_bytes,
max_in_flight: 8,
bootstrap_timeout_ms: 1_000,
},
}
}
fn channel(role: SidecarPeerRole, max_frame_bytes: u32) -> AuthenticatedSidecarChannel {
AuthenticatedSidecarChannel::new(
role,
"handshake-session".into(),
9,
SidecarSessionKey::from_bytes(KEY),
max_frame_bytes,
)
.unwrap()
}
#[test]
fn supervisor_and_runtime_authenticate_the_same_selection() {
let mut supervisor =
SidecarHandshake::new(offer(SidecarPeerRole::Supervisor, 4096, 0b0111)).unwrap();
let mut runtime =
SidecarHandshake::new(offer(SidecarPeerRole::Runtime, 2048, 0b1011)).unwrap();
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
let offer_frame = supervisor.start(&mut supervisor_channel).unwrap();
let acceptance = runtime
.receive(&mut runtime_channel, &offer_frame)
.unwrap()
.unwrap();
assert_eq!(runtime_channel.max_frame_bytes(), 4096);
assert_eq!(runtime.state(), SidecarHandshakeState::AcceptancePending);
assert!(runtime.negotiated().is_none());
runtime.complete_acceptance(&mut runtime_channel).unwrap();
assert_eq!(runtime_channel.max_frame_bytes(), 2048);
assert_eq!(runtime.state(), SidecarHandshakeState::Authenticated);
assert!(supervisor
.receive(&mut supervisor_channel, &acceptance)
.unwrap()
.is_none());
assert_eq!(supervisor_channel.max_frame_bytes(), 2048);
assert_eq!(supervisor.state(), SidecarHandshakeState::Authenticated);
assert_eq!(
SidecarProtocolSelection::from(supervisor.negotiated().unwrap()),
SidecarProtocolSelection::from(runtime.negotiated().unwrap())
);
assert_eq!(supervisor.negotiated().unwrap().feature_bits, 0b0011);
let active = supervisor_channel.seal(&"active").unwrap();
assert_eq!(
u16::from_be_bytes(active[6..8].try_into().unwrap()),
supervisor.negotiated().unwrap().protocol_minor
);
assert_eq!(u64::from_be_bytes(active[17..25].try_into().unwrap()), 2);
assert_eq!(runtime_channel.open::<String>(&active).unwrap(), "active");
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct HandshakeFixture {
schema_version: u8,
session: FixtureSession,
supervisor_offer: SidecarProtocolOffer,
runtime_offer: SidecarProtocolOffer,
selection: SidecarProtocolSelection,
offer_frame_base64: String,
accept_frame_base64: String,
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase")]
struct FixtureSession {
id: String,
generation: u64,
key_base64: String,
}
#[test]
fn cross_language_handshake_vector_is_exact() {
let fixture: HandshakeFixture = serde_json::from_str(include_str!(
"../../../test/fixtures/node-sidecar-handshake-v1.json"
))
.unwrap();
assert_eq!(fixture.schema_version, 1);
let key: [u8; 32] = BASE64
.decode(&fixture.session.key_base64)
.unwrap()
.try_into()
.unwrap();
let bootstrap_frame_bytes = fixture.supervisor_offer.limits.max_frame_bytes;
let new_channel = |role| {
AuthenticatedSidecarChannel::new(
role,
fixture.session.id.clone(),
fixture.session.generation,
SidecarSessionKey::from_bytes(key),
bootstrap_frame_bytes,
)
.unwrap()
};
let mut supervisor = SidecarHandshake::new(fixture.supervisor_offer).unwrap();
let mut runtime = SidecarHandshake::new(fixture.runtime_offer).unwrap();
let mut supervisor_channel = new_channel(SidecarPeerRole::Supervisor);
let mut runtime_channel = new_channel(SidecarPeerRole::Runtime);
let offer_frame = supervisor.start(&mut supervisor_channel).unwrap();
assert_eq!(
offer_frame,
BASE64.decode(fixture.offer_frame_base64).unwrap()
);
let accept_frame = runtime
.receive(&mut runtime_channel, &offer_frame)
.unwrap()
.unwrap();
assert_eq!(
accept_frame,
BASE64.decode(fixture.accept_frame_base64).unwrap()
);
runtime.complete_acceptance(&mut runtime_channel).unwrap();
assert!(supervisor
.receive(&mut supervisor_channel, &accept_frame)
.unwrap()
.is_none());
assert_eq!(
SidecarProtocolSelection::from(supervisor.negotiated().unwrap()),
fixture.selection
);
assert_eq!(
SidecarProtocolSelection::from(runtime.negotiated().unwrap()),
fixture.selection
);
}
#[test]
fn forged_selection_fails_and_retires_supervisor() {
let mut supervisor =
SidecarHandshake::new(offer(SidecarPeerRole::Supervisor, 4096, 0b0011)).unwrap();
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
let mut malicious_runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
let offer_frame = supervisor.start(&mut supervisor_channel).unwrap();
let _: SidecarHandshakeMessage = malicious_runtime_channel.open(&offer_frame).unwrap();
let runtime_offer = offer(SidecarPeerRole::Runtime, 2048, 0b0011);
let forged = malicious_runtime_channel
.seal(&SidecarHandshakeMessage::Accept {
offer: runtime_offer,
selection: SidecarProtocolSelection {
protocol_major: crate::SIDECAR_PROTOCOL_MAJOR,
protocol_minor: crate::SIDECAR_PROTOCOL_MINOR,
feature_bits: u64::MAX,
limits: SidecarLimits {
max_frame_bytes: 4096,
max_in_flight: u16::MAX,
bootstrap_timeout_ms: u32::MAX,
},
},
})
.unwrap();
assert!(matches!(
supervisor.receive(&mut supervisor_channel, &forged),
Err(SidecarHandshakeError::SelectionMismatch)
));
assert_eq!(supervisor.state(), SidecarHandshakeState::Failed);
assert!(supervisor_channel.is_retired());
}
#[test]
fn runtime_commits_selection_only_after_acceptance_delivery() {
let mut supervisor =
SidecarHandshake::new(offer(SidecarPeerRole::Supervisor, 4096, 0)).unwrap();
let mut runtime = SidecarHandshake::new(offer(SidecarPeerRole::Runtime, 128, 0)).unwrap();
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
let offer_frame = supervisor.start(&mut supervisor_channel).unwrap();
let acceptance = runtime
.receive(&mut runtime_channel, &offer_frame)
.unwrap()
.unwrap();
assert!(acceptance.len() > 128);
assert_eq!(runtime.state(), SidecarHandshakeState::AcceptancePending);
assert_eq!(runtime_channel.max_frame_bytes(), 4096);
assert!(runtime.negotiated().is_none());
runtime.complete_acceptance(&mut runtime_channel).unwrap();
assert_eq!(runtime.state(), SidecarHandshakeState::Authenticated);
assert_eq!(runtime_channel.max_frame_bytes(), 128);
}
#[test]
fn failed_acceptance_delivery_is_terminal() {
let mut supervisor =
SidecarHandshake::new(offer(SidecarPeerRole::Supervisor, 4096, 0)).unwrap();
let mut runtime = SidecarHandshake::new(offer(SidecarPeerRole::Runtime, 2048, 0)).unwrap();
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
let offer_frame = supervisor.start(&mut supervisor_channel).unwrap();
runtime
.receive(&mut runtime_channel, &offer_frame)
.unwrap()
.unwrap();
runtime_channel.retire();
assert!(matches!(
runtime.complete_acceptance(&mut runtime_channel),
Err(SidecarHandshakeError::Frame(
SidecarFrameError::ChannelRetired
))
));
assert_eq!(runtime.state(), SidecarHandshakeState::Failed);
assert!(runtime.negotiated().is_none());
}
#[test]
fn handshake_cannot_move_to_another_channel_instance() {
let mut supervisor =
SidecarHandshake::new(offer(SidecarPeerRole::Supervisor, 4096, 0b0011)).unwrap();
let mut original_channel = channel(SidecarPeerRole::Supervisor, 4096);
let _offer_frame = supervisor.start(&mut original_channel).unwrap();
let runtime_offer = offer(SidecarPeerRole::Runtime, 2048, 0b0011);
let selection = SidecarProtocolSelection::from(
&negotiate_sidecar_protocol(&supervisor.local_offer, &runtime_offer).unwrap(),
);
let mut other_runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
let acceptance = other_runtime_channel
.seal(&SidecarHandshakeMessage::Accept {
offer: runtime_offer,
selection,
})
.unwrap();
let mut other_supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
assert!(matches!(
supervisor.receive(&mut other_supervisor_channel, &acceptance),
Err(SidecarHandshakeError::ChannelInstanceMismatch)
));
assert_eq!(
supervisor.state(),
SidecarHandshakeState::AwaitingAcceptance
);
assert!(other_supervisor_channel.is_retired());
assert!(!original_channel.is_retired());
assert!(supervisor
.receive(&mut original_channel, &acceptance)
.unwrap()
.is_none());
assert_eq!(supervisor.state(), SidecarHandshakeState::Authenticated);
}
#[test]
fn incompatible_offer_fails_and_retires_runtime() {
let mut supervisor_offer = offer(SidecarPeerRole::Supervisor, 4096, 0);
supervisor_offer.protocol_major += 1;
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
let incompatible = supervisor_channel
.seal(&SidecarHandshakeMessage::Offer {
offer: supervisor_offer,
})
.unwrap();
let mut runtime = SidecarHandshake::new(offer(SidecarPeerRole::Runtime, 4096, 0)).unwrap();
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
assert!(matches!(
runtime.receive(&mut runtime_channel, &incompatible),
Err(SidecarHandshakeError::Negotiation(
SidecarProtocolError::UnsupportedMajor { .. }
))
));
assert_eq!(runtime.state(), SidecarHandshakeState::Failed);
assert!(runtime_channel.is_retired());
}
#[test]
fn wrong_order_is_terminal() {
let mut runtime = SidecarHandshake::new(offer(SidecarPeerRole::Runtime, 4096, 0)).unwrap();
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
assert!(matches!(
runtime.start(&mut runtime_channel),
Err(SidecarHandshakeError::SupervisorMustInitiate)
));
assert_eq!(runtime.state(), SidecarHandshakeState::Failed);
assert!(runtime_channel.is_retired());
}
#[test]
fn swapped_channel_roles_are_terminal_before_frame_processing() {
let mut supervisor =
SidecarHandshake::new(offer(SidecarPeerRole::Supervisor, 4096, 0)).unwrap();
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
assert!(matches!(
supervisor.start(&mut runtime_channel),
Err(SidecarHandshakeError::ChannelRoleMismatch)
));
assert_eq!(supervisor.state(), SidecarHandshakeState::Failed);
assert!(runtime_channel.is_retired());
let mut runtime = SidecarHandshake::new(offer(SidecarPeerRole::Runtime, 4096, 0)).unwrap();
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
assert!(matches!(
runtime.receive(&mut supervisor_channel, b"ignored"),
Err(SidecarHandshakeError::ChannelRoleMismatch)
));
assert_eq!(runtime.state(), SidecarHandshakeState::Failed);
assert!(supervisor_channel.is_retired());
}
#[test]
fn malformed_frame_retires_handshake_and_channel() {
let mut runtime = SidecarHandshake::new(offer(SidecarPeerRole::Runtime, 4096, 0)).unwrap();
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
assert!(matches!(
runtime.receive(&mut runtime_channel, b"not-authenticated"),
Err(SidecarHandshakeError::Frame(_))
));
assert_eq!(runtime.state(), SidecarHandshakeState::Failed);
assert!(runtime_channel.is_retired());
}
#[test]
fn invalid_local_identity_is_rejected_before_channel_use() {
let mut invalid = offer(SidecarPeerRole::Runtime, 4096, 0);
invalid.peer.artifact_identity.clear();
assert!(matches!(
SidecarHandshake::new(invalid),
Err(SidecarHandshakeError::InvalidPeerIdentity)
));
}
#[test]
fn invalid_supervisor_identity_is_terminal_for_runtime() {
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
let mut invalid = offer(SidecarPeerRole::Supervisor, 4096, 0);
invalid.peer.name = " ".into();
let frame = supervisor_channel
.seal(&SidecarHandshakeMessage::Offer { offer: invalid })
.unwrap();
let mut runtime = SidecarHandshake::new(offer(SidecarPeerRole::Runtime, 4096, 0)).unwrap();
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
assert!(matches!(
runtime.receive(&mut runtime_channel, &frame),
Err(SidecarHandshakeError::InvalidPeerIdentity)
));
assert_eq!(runtime.state(), SidecarHandshakeState::Failed);
assert!(runtime_channel.is_retired());
}
#[test]
fn invalid_runtime_identity_is_terminal_for_supervisor() {
let mut supervisor =
SidecarHandshake::new(offer(SidecarPeerRole::Supervisor, 4096, 0)).unwrap();
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 4096);
let offer_frame = supervisor.start(&mut supervisor_channel).unwrap();
let mut runtime_channel = channel(SidecarPeerRole::Runtime, 4096);
let _: SidecarHandshakeMessage = runtime_channel.open(&offer_frame).unwrap();
let mut invalid = offer(SidecarPeerRole::Runtime, 4096, 0);
invalid.peer.version.clear();
let selection = SidecarProtocolSelection::from(
&negotiate_sidecar_protocol(&supervisor.local_offer, &invalid).unwrap(),
);
let frame = runtime_channel
.seal(&SidecarHandshakeMessage::Accept {
offer: invalid,
selection,
})
.unwrap();
assert!(matches!(
supervisor.receive(&mut supervisor_channel, &frame),
Err(SidecarHandshakeError::InvalidPeerIdentity)
));
assert_eq!(supervisor.state(), SidecarHandshakeState::Failed);
assert!(supervisor_channel.is_retired());
}
#[test]
fn unencodable_initial_offer_is_terminal() {
let mut supervisor =
SidecarHandshake::new(offer(SidecarPeerRole::Supervisor, 128, 0)).unwrap();
let mut supervisor_channel = channel(SidecarPeerRole::Supervisor, 128);
assert!(matches!(
supervisor.start(&mut supervisor_channel),
Err(SidecarHandshakeError::Frame(
SidecarFrameError::FrameTooLarge { .. }
))
));
assert_eq!(supervisor.state(), SidecarHandshakeState::Failed);
assert!(supervisor_channel.is_retired());
}
}

File diff suppressed because it is too large Load diff

File diff suppressed because it is too large Load diff

View file

@ -1,8 +1,13 @@
use ed25519_dalek::{Signer, SigningKey};
use futures_util::{SinkExt, StreamExt};
use openclaw_node_host::{
CommandRuntime, ConnectAuth, HandlerError, NodeClient, NodeClientConfig, NodeConnectOptions,
NodeIdentity, NodeProtocolVersion, NodeSession,
AuthenticatedSidecarChannel, CancellationToken, CommandRuntime, ConnectAuth, HandlerError,
NodeClient, NodeClientConfig, NodeConnectOptions, NodeIdentity, NodeProtocolVersion,
NodeSession, SidecarAdapterError, SidecarAdapterFuture, SidecarAdmissionDecision,
SidecarCapabilityAdapter, SidecarCommandRegistration, SidecarConfigurationExchange,
SidecarHandshake, SidecarInvocation, SidecarInvocationResult, SidecarLimits,
SidecarPeerIdentity, SidecarPeerRole, SidecarProtocolOffer, SidecarRuntimeBridge,
SidecarRuntimeConfiguration, SidecarSessionKey, SIDECAR_PROTOCOL_MAJOR, SIDECAR_PROTOCOL_MINOR,
};
use serde_json::{json, Value};
use std::{
@ -311,6 +316,257 @@ async fn wire_cancellation_during_admission_prevents_handler_construction() {
assert!(!handler_constructed.load(Ordering::SeqCst));
}
#[derive(Default)]
struct AuthorityAdapter {
admissions: AtomicUsize,
invocations: AtomicUsize,
native_effects: AtomicUsize,
retiring_invocation_started: Notify,
}
impl SidecarCapabilityAdapter for AuthorityAdapter {
fn admit(
&self,
invocation: SidecarInvocation,
_cancellation: CancellationToken,
) -> SidecarAdapterFuture<Result<SidecarAdmissionDecision, SidecarAdapterError>> {
self.admissions.fetch_add(1, Ordering::SeqCst);
Box::pin(async move {
if invocation.command == "product.settings" {
Ok(SidecarAdmissionDecision::Deny {
code: "LOCAL_DENY".into(),
message: "denied by product policy".into(),
})
} else {
Ok(SidecarAdmissionDecision::Allow)
}
})
}
fn invoke(
&self,
invocation: SidecarInvocation,
cancellation: CancellationToken,
) -> SidecarAdapterFuture<Result<SidecarInvocationResult, SidecarAdapterError>> {
self.invocations.fetch_add(1, Ordering::SeqCst);
if invocation.params == json!({"retire": true}) {
self.retiring_invocation_started.notify_one();
return Box::pin(async move {
cancellation.cancelled().await;
Err(SidecarAdapterError::new(
"CANCELLED",
"retired before native effects",
))
});
}
self.native_effects.fetch_add(1, Ordering::SeqCst);
Box::pin(async move {
Ok(SidecarInvocationResult::Success {
payload: json!({"handledBy": "sidecar-bridge"}),
})
})
}
}
#[tokio::test]
async fn sidecar_bridge_preserves_authority_through_the_public_runtime() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let server = tokio::spawn(async move {
let (tcp, _) = listener.accept().await.unwrap();
let mut socket = accept_async(tcp).await.unwrap();
send_json(
&mut socket,
json!({"type":"event","event":"connect.challenge",
"payload":{"nonce":"node-nonce","ts":1_700_000_000_123_u64}}),
)
.await;
let connect = receive_json(&mut socket).await;
send_json(
&mut socket,
json!({"type":"res","id":connect["id"],"ok":true,
"payload":{"type":"hello-ok","protocol":4}}),
)
.await;
for (id, command, params, expected) in [
(
"allowed",
"product.status",
json!({}),
("", json!({"handledBy": "sidecar-bridge"})),
),
(
"denied",
"product.settings",
json!({}),
("LOCAL_DENY", Value::Null),
),
(
"retired",
"product.status",
json!({"retire": true}),
("SIDECAR_CHANNEL_RETIRED", Value::Null),
),
] {
send_json(
&mut socket,
json!({"type":"event","event":"node.invoke.request","payload":{
"id":id,"nodeId":"node-1","command":command,
"paramsJSON":params.to_string()
}}),
)
.await;
let result = receive_json(&mut socket).await;
assert_eq!(result["method"], "node.invoke.result");
assert_eq!(result["params"]["id"], id);
if expected.0.is_empty() {
assert_eq!(result["params"]["payload"], expected.1);
} else {
assert_eq!(result["params"]["ok"], false);
assert_eq!(result["params"]["error"]["code"], expected.0);
}
acknowledge(&mut socket, &result).await;
}
socket.close(None).await.unwrap();
});
let adapter = Arc::new(AuthorityAdapter::default());
let (bridge, mut channel) = activated_sidecar_bridge(&adapter);
let runtime = bridge.into_runtime();
let connect_runtime = runtime.clone();
let session = NodeClient::connect(
NodeClientConfig::new(format!("ws://{address}")),
move |_challenge| {
let connect_runtime = connect_runtime.clone();
async move {
Ok::<_, io::Error>(
connect_runtime.activate(
NodeConnectOptions::new("test", "linux")
.auth(ConnectAuth::token("test-token"))
.identity(NodeIdentity::from_secret_bytes([7; 32])),
),
)
}
},
)
.await
.unwrap();
let run = tokio::spawn(async move { runtime.run(session).await });
tokio::time::timeout(
Duration::from_secs(1),
adapter.retiring_invocation_started.notified(),
)
.await
.expect("retiring sidecar invocation did not reach the adapter boundary");
channel.retire();
assert!(run.await.unwrap().is_err());
server.await.unwrap();
assert_eq!(adapter.admissions.load(Ordering::SeqCst), 3);
assert_eq!(adapter.invocations.load(Ordering::SeqCst), 2);
assert_eq!(adapter.native_effects.load(Ordering::SeqCst), 1);
}
fn activated_sidecar_bridge(
adapter: &Arc<AuthorityAdapter>,
) -> (SidecarRuntimeBridge, AuthenticatedSidecarChannel) {
let mut supervisor_handshake =
SidecarHandshake::new(sidecar_offer(SidecarPeerRole::Supervisor)).unwrap();
let mut runtime_handshake =
SidecarHandshake::new(sidecar_offer(SidecarPeerRole::Runtime)).unwrap();
let mut supervisor_channel = sidecar_channel(SidecarPeerRole::Supervisor);
let mut runtime_channel = sidecar_channel(SidecarPeerRole::Runtime);
let offer = supervisor_handshake.start(&mut supervisor_channel).unwrap();
let acceptance = runtime_handshake
.receive(&mut runtime_channel, &offer)
.unwrap()
.unwrap();
runtime_handshake
.complete_acceptance(&mut runtime_channel)
.unwrap();
supervisor_handshake
.receive(&mut supervisor_channel, &acceptance)
.unwrap();
let mut supervisor_exchange = SidecarConfigurationExchange::new(supervisor_handshake).unwrap();
let mut runtime_exchange = SidecarConfigurationExchange::new(runtime_handshake).unwrap();
let configuration = SidecarRuntimeConfiguration {
manifest_generation: 1,
capabilities: vec!["native.status".into()],
commands: vec![
SidecarCommandRegistration {
name: "product.settings".into(),
},
SidecarCommandRegistration {
name: "product.status".into(),
},
],
max_concurrency: 2,
max_input_bytes: 1_024,
max_output_bytes: 1_024,
default_timeout_ms: 1_000,
max_timeout_ms: 5_000,
result_grace_ms: 50,
};
let configure = supervisor_exchange
.start(&mut supervisor_channel, &configuration)
.unwrap();
runtime_exchange
.receive(&mut runtime_channel, &configure)
.unwrap()
.unwrap();
let manifest = runtime_exchange.validated_manifest().unwrap().clone();
let configured = runtime_exchange
.acknowledge(&mut runtime_channel, &manifest)
.unwrap();
runtime_exchange
.complete_acknowledgement(&mut runtime_channel)
.unwrap();
supervisor_exchange
.receive(&mut supervisor_channel, &configured)
.unwrap();
let bridge =
SidecarRuntimeBridge::activate(&mut runtime_exchange, &mut runtime_channel, adapter)
.unwrap();
(bridge, runtime_channel)
}
fn sidecar_channel(role: SidecarPeerRole) -> AuthenticatedSidecarChannel {
AuthenticatedSidecarChannel::new(
role,
"authority-session".into(),
1,
SidecarSessionKey::from_bytes([0x55; 32]),
4_096,
)
.unwrap()
}
fn sidecar_offer(role: SidecarPeerRole) -> SidecarProtocolOffer {
SidecarProtocolOffer {
protocol_major: SIDECAR_PROTOCOL_MAJOR,
protocol_minor: SIDECAR_PROTOCOL_MINOR,
peer: SidecarPeerIdentity {
role,
name: match role {
SidecarPeerRole::Supervisor => "test-supervisor",
SidecarPeerRole::Runtime => "test-runtime",
}
.into(),
version: "test".into(),
artifact_identity: "sha256:test-only".into(),
},
feature_bits: 0,
limits: SidecarLimits {
max_frame_bytes: 4_096,
max_in_flight: 4,
bootstrap_timeout_ms: 1_000,
},
}
}
#[tokio::test]
async fn node_protocol_fallback_uses_fresh_legacy_connect_material_and_recovers_to_v4() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();

View file

@ -0,0 +1,295 @@
use std::{
env,
process::{Command, Stdio},
time::Duration,
};
use openclaw_node_host::{
read_sidecar_frame, write_sidecar_frame, AuthenticatedSidecarChannel, SidecarAdmissionDecision,
SidecarCommandRegistration, SidecarConfigurationExchange, SidecarHandshake,
SidecarHandshakeState, SidecarInvocation, SidecarInvocationResult, SidecarLimits,
SidecarPeerIdentity, SidecarPeerRole, SidecarProtocolOffer, SidecarRuntimeConfiguration,
SidecarRuntimeMessage, SidecarSessionKey, SIDECAR_PROTOCOL_MAJOR, SIDECAR_PROTOCOL_MINOR,
};
use serde_json::json;
use tokio::net::{TcpListener, TcpStream};
const CHILD_ENV: &str = "OPENCLAW_SIDECAR_PROCESS_CHILD";
const CHILD_ADDRESS_ENV: &str = "OPENCLAW_SIDECAR_PROCESS_ADDRESS";
const FRAME_LIMIT: u32 = 4_096;
const SESSION_KEY: [u8; 32] = [0x5a; 32];
const IO_TIMEOUT: Duration = Duration::from_secs(5);
#[tokio::test]
#[expect(
clippy::too_many_lines,
reason = "the parent transcript intentionally keeps the full ordered process exchange visible"
)]
async fn authenticated_sidecar_crosses_a_real_process_boundary() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let child = Command::new(env::current_exe().unwrap())
.arg("--exact")
.arg("sidecar_process_child")
.arg("--nocapture")
.env(CHILD_ENV, "1")
.env(CHILD_ADDRESS_ENV, address.to_string())
.stdin(Stdio::null())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn()
.unwrap();
let (stream, _) = tokio::time::timeout(IO_TIMEOUT, listener.accept())
.await
.expect("sidecar child did not connect")
.unwrap();
let (mut input, mut output) = stream.into_split();
let mut channel = channel(SidecarPeerRole::Supervisor);
let mut handshake = SidecarHandshake::new(offer(SidecarPeerRole::Supervisor)).unwrap();
let offered = handshake.start(&mut channel).unwrap();
write_sidecar_frame(&mut output, &offered, FRAME_LIMIT, IO_TIMEOUT)
.await
.unwrap();
let accepted = read_sidecar_frame(&mut input, FRAME_LIMIT, IO_TIMEOUT)
.await
.unwrap();
assert!(handshake
.receive(&mut channel, &accepted)
.unwrap()
.is_none());
assert_eq!(handshake.state(), SidecarHandshakeState::Authenticated);
let mut exchange = SidecarConfigurationExchange::new(handshake).unwrap();
let configuration = configuration();
let configure = exchange.start(&mut channel, &configuration).unwrap();
write_sidecar_frame(
&mut output,
&configure,
channel.max_frame_bytes(),
IO_TIMEOUT,
)
.await
.unwrap();
let configured = read_sidecar_frame(&mut input, channel.max_frame_bytes(), IO_TIMEOUT)
.await
.unwrap();
assert!(exchange
.receive(&mut channel, &configured)
.unwrap()
.is_none());
let invocation = invocation();
let admission = channel
.seal(&SidecarRuntimeMessage::AdmissionRequest {
invocation: invocation.clone(),
})
.unwrap();
write_sidecar_frame(
&mut output,
&admission,
channel.max_frame_bytes(),
IO_TIMEOUT,
)
.await
.unwrap();
let decision = read_sidecar_frame(&mut input, channel.max_frame_bytes(), IO_TIMEOUT)
.await
.unwrap();
assert_eq!(
channel.open::<SidecarRuntimeMessage>(&decision).unwrap(),
SidecarRuntimeMessage::AdmissionDecision {
invocation_id: invocation.id.clone(),
decision: SidecarAdmissionDecision::Allow,
}
);
let invoke = channel
.seal(&SidecarRuntimeMessage::Invoke {
invocation: invocation.clone(),
})
.unwrap();
write_sidecar_frame(&mut output, &invoke, channel.max_frame_bytes(), IO_TIMEOUT)
.await
.unwrap();
let result = read_sidecar_frame(&mut input, channel.max_frame_bytes(), IO_TIMEOUT)
.await
.unwrap();
assert_eq!(
channel.open::<SidecarRuntimeMessage>(&result).unwrap(),
SidecarRuntimeMessage::Result {
invocation_id: invocation.id,
result: SidecarInvocationResult::Success {
payload: json!({"handledBy": "process-runtime", "value": 42}),
},
}
);
let output = tokio::task::spawn_blocking(move || child.wait_with_output())
.await
.unwrap()
.unwrap();
assert!(
output.status.success(),
"child failed\nstdout:\n{}\nstderr:\n{}",
String::from_utf8_lossy(&output.stdout),
String::from_utf8_lossy(&output.stderr)
);
}
#[tokio::test]
async fn sidecar_process_child() {
if env::var_os(CHILD_ENV).is_none() {
return;
}
let address = env::var(CHILD_ADDRESS_ENV).expect("child address");
let stream = tokio::time::timeout(IO_TIMEOUT, TcpStream::connect(address))
.await
.expect("supervisor was unavailable")
.unwrap();
let (mut input, mut output) = stream.into_split();
let mut channel = channel(SidecarPeerRole::Runtime);
let mut handshake = SidecarHandshake::new(offer(SidecarPeerRole::Runtime)).unwrap();
let offered = read_sidecar_frame(&mut input, FRAME_LIMIT, IO_TIMEOUT)
.await
.unwrap();
let accepted = handshake
.receive(&mut channel, &offered)
.unwrap()
.expect("supervisor offer must produce acceptance");
write_sidecar_frame(&mut output, &accepted, FRAME_LIMIT, IO_TIMEOUT)
.await
.unwrap();
handshake.complete_acceptance(&mut channel).unwrap();
assert_eq!(handshake.state(), SidecarHandshakeState::Authenticated);
let mut exchange = SidecarConfigurationExchange::new(handshake).unwrap();
let configure = read_sidecar_frame(&mut input, channel.max_frame_bytes(), IO_TIMEOUT)
.await
.unwrap();
let received = exchange
.receive(&mut channel, &configure)
.unwrap()
.expect("runtime must receive configuration");
assert_eq!(received, configuration());
let manifest = exchange.validated_manifest().unwrap().clone();
let configured = exchange.acknowledge(&mut channel, &manifest).unwrap();
write_sidecar_frame(
&mut output,
&configured,
channel.max_frame_bytes(),
IO_TIMEOUT,
)
.await
.unwrap();
exchange.complete_acknowledgement(&mut channel).unwrap();
let admission = read_sidecar_frame(&mut input, channel.max_frame_bytes(), IO_TIMEOUT)
.await
.unwrap();
let SidecarRuntimeMessage::AdmissionRequest { invocation } =
channel.open::<SidecarRuntimeMessage>(&admission).unwrap()
else {
panic!("expected admission request");
};
assert_eq!(invocation.command, "product.status");
let decision = channel
.seal(&SidecarRuntimeMessage::AdmissionDecision {
invocation_id: invocation.id,
decision: SidecarAdmissionDecision::Allow,
})
.unwrap();
write_sidecar_frame(
&mut output,
&decision,
channel.max_frame_bytes(),
IO_TIMEOUT,
)
.await
.unwrap();
let invoke = read_sidecar_frame(&mut input, channel.max_frame_bytes(), IO_TIMEOUT)
.await
.unwrap();
let SidecarRuntimeMessage::Invoke { invocation } =
channel.open::<SidecarRuntimeMessage>(&invoke).unwrap()
else {
panic!("expected invocation");
};
assert_eq!(invocation.params, json!({"value": 42}));
let result = channel
.seal(&SidecarRuntimeMessage::Result {
invocation_id: invocation.id,
result: SidecarInvocationResult::Success {
payload: json!({"handledBy": "process-runtime", "value": 42}),
},
})
.unwrap();
write_sidecar_frame(&mut output, &result, channel.max_frame_bytes(), IO_TIMEOUT)
.await
.unwrap();
}
fn configuration() -> SidecarRuntimeConfiguration {
SidecarRuntimeConfiguration {
manifest_generation: 3,
capabilities: vec!["native.status".into()],
commands: vec![SidecarCommandRegistration {
name: "product.status".into(),
}],
max_concurrency: 1,
max_input_bytes: 1_024,
max_output_bytes: 1_024,
default_timeout_ms: 1_000,
max_timeout_ms: 5_000,
result_grace_ms: 50,
}
}
fn invocation() -> SidecarInvocation {
SidecarInvocation {
id: "invoke-process-1".into(),
node_id: "node-process-1".into(),
command: "product.status".into(),
params: json!({"value": 42}),
timeout_ms: Some(1_000),
idempotency_key: Some("process-idempotency-1".into()),
session_key: Some("agent:main:process".into()),
}
}
fn channel(role: SidecarPeerRole) -> AuthenticatedSidecarChannel {
AuthenticatedSidecarChannel::new(
role,
"process-session".into(),
7,
SidecarSessionKey::from_bytes(SESSION_KEY),
FRAME_LIMIT,
)
.unwrap()
}
fn offer(role: SidecarPeerRole) -> SidecarProtocolOffer {
SidecarProtocolOffer {
protocol_major: SIDECAR_PROTOCOL_MAJOR,
protocol_minor: SIDECAR_PROTOCOL_MINOR,
peer: SidecarPeerIdentity {
role,
name: match role {
SidecarPeerRole::Supervisor => "process-supervisor",
SidecarPeerRole::Runtime => "process-runtime",
}
.into(),
version: "test".into(),
artifact_identity: "sha256:process-fixture".into(),
},
feature_bits: 0b0011,
limits: SidecarLimits {
max_frame_bytes: 2_048,
max_in_flight: 8,
bootstrap_timeout_ms: 1_000,
},
}
}

View file

@ -0,0 +1,52 @@
{
"schemaVersion": 1,
"session": {
"id": "handshake-session",
"generation": 9,
"keyBase64": "PDw8PDw8PDw8PDw8PDw8PDw8PDw8PDw8PDw8PDw8PDw="
},
"supervisorOffer": {
"protocolMajor": 1,
"protocolMinor": 0,
"peer": {
"role": "supervisor",
"name": "test-product",
"version": "1.0.0",
"artifactIdentity": "sha256:test-only"
},
"featureBits": 7,
"limits": {
"maxFrameBytes": 4096,
"maxInFlight": 8,
"bootstrapTimeoutMs": 1000
}
},
"runtimeOffer": {
"protocolMajor": 1,
"protocolMinor": 0,
"peer": {
"role": "runtime",
"name": "openclaw-node",
"version": "1.0.0",
"artifactIdentity": "sha256:test-only"
},
"featureBits": 11,
"limits": {
"maxFrameBytes": 2048,
"maxInFlight": 8,
"bootstrapTimeoutMs": 1000
}
},
"selection": {
"protocolMajor": 1,
"protocolMinor": 0,
"featureBits": 3,
"limits": {
"maxFrameBytes": 2048,
"maxInFlight": 8,
"bootstrapTimeoutMs": 1000
}
},
"offerFrameBase64": "T0NTQwABAAABAAAAAAAAAAkAAAAAAAAAAQARAAABA2hhbmRzaGFrZS1zZXNzaW9ueyJ0eXBlIjoib2ZmZXIiLCJvZmZlciI6eyJwcm90b2NvbE1ham9yIjoxLCJwcm90b2NvbE1pbm9yIjowLCJwZWVyIjp7InJvbGUiOiJzdXBlcnZpc29yIiwibmFtZSI6InRlc3QtcHJvZHVjdCIsInZlcnNpb24iOiIxLjAuMCIsImFydGlmYWN0SWRlbnRpdHkiOiJzaGEyNTY6dGVzdC1vbmx5In0sImZlYXR1cmVCaXRzIjo3LCJsaW1pdHMiOnsibWF4RnJhbWVCeXRlcyI6NDA5NiwibWF4SW5GbGlnaHQiOjgsImJvb3RzdHJhcFRpbWVvdXRNcyI6MTAwMH19fTSxSPVri786qBe92/qj1NGhn0efEyqVJfiZrCOdLbr1",
"acceptFrameBase64": "T0NTQwABAAACAAAAAAAAAAkAAAAAAAAAAQARAAABj2hhbmRzaGFrZS1zZXNzaW9ueyJ0eXBlIjoiYWNjZXB0Iiwib2ZmZXIiOnsicHJvdG9jb2xNYWpvciI6MSwicHJvdG9jb2xNaW5vciI6MCwicGVlciI6eyJyb2xlIjoicnVudGltZSIsIm5hbWUiOiJvcGVuY2xhdy1ub2RlIiwidmVyc2lvbiI6IjEuMC4wIiwiYXJ0aWZhY3RJZGVudGl0eSI6InNoYTI1Njp0ZXN0LW9ubHkifSwiZmVhdHVyZUJpdHMiOjExLCJsaW1pdHMiOnsibWF4RnJhbWVCeXRlcyI6MjA0OCwibWF4SW5GbGlnaHQiOjgsImJvb3RzdHJhcFRpbWVvdXRNcyI6MTAwMH19LCJzZWxlY3Rpb24iOnsicHJvdG9jb2xNYWpvciI6MSwicHJvdG9jb2xNaW5vciI6MCwiZmVhdHVyZUJpdHMiOjMsImxpbWl0cyI6eyJtYXhGcmFtZUJ5dGVzIjoyMDQ4LCJtYXhJbkZsaWdodCI6OCwiYm9vdHN0cmFwVGltZW91dE1zIjoxMDAwfX19DY3CiBg6igkouQFlG07FuMhFJuNEFCONtdVA6glc5tU="
}

View file

@ -0,0 +1,36 @@
{
"schemaVersion": 1,
"localOffer": {
"protocolMajor": 1,
"protocolMinor": 0,
"peer": {
"role": "supervisor",
"name": "fixture-supervisor",
"version": "1.0.0",
"artifactIdentity": "sha256:fixture-supervisor"
},
"featureBits": 4503599627370499,
"limits": {
"maxFrameBytes": 4096,
"maxInFlight": 8,
"bootstrapTimeoutMs": 1000
}
},
"remoteOffer": {
"protocolMajor": 1,
"protocolMinor": 0,
"peer": {
"role": "runtime",
"name": "fixture-runtime",
"version": "1.0.0",
"artifactIdentity": "sha256:fixture-runtime"
},
"featureBits": 4503599627370501,
"limits": {
"maxFrameBytes": 2048,
"maxInFlight": 4,
"bootstrapTimeoutMs": 500
}
},
"selectedFeatureBits": 4503599627370497
}

View file

@ -0,0 +1,15 @@
{
"schemaVersion": 1,
"session": {
"id": "session-7",
"generation": 7,
"sessionKeyBase64": "WlpaWlpaWlpaWlpaWlpaWlpaWlpaWlpaWlpaWlpaWlo="
},
"supervisorProbe": {
"payload": {
"requestId": "abc",
"type": "probe"
},
"frameBase64": "T0NTQwABAAABAAAAAAAAAAcAAAAAAAAAAQAJAAAAInNlc3Npb24tN3sicmVxdWVzdElkIjoiYWJjIiwidHlwZSI6InByb2JlIn04my4v5goF1qHj7BfwsC4o3oTzJCbGaE5jWu5WdqOVlQ=="
}
}

View file

@ -0,0 +1,99 @@
{
"schemaVersion": 1,
"messages": [
{
"type": "configure",
"configuration": {
"manifestGeneration": 3,
"capabilities": ["native.settings", "native.status"],
"commands": [
{ "name": "product.settings" },
{ "name": "product.status" }
],
"maxConcurrency": 2,
"maxInputBytes": 1024,
"maxOutputBytes": 1024,
"defaultTimeoutMs": 1000,
"maxTimeoutMs": 5000,
"resultGraceMs": 50
}
},
{
"type": "configured",
"manifest": {
"manifestGeneration": 3,
"capabilities": ["native.settings", "native.status"],
"commands": ["product.settings", "product.status"]
}
},
{
"type": "admission-request",
"invocation": {
"id": "invoke-1",
"nodeId": "node-1",
"command": "product.status",
"params": { "verbose": true },
"timeoutMs": 1000,
"idempotencyKey": "idem-1",
"sessionKey": "agent:main:main"
}
},
{
"type": "admission-decision",
"invocationId": "invoke-1",
"decision": { "outcome": "allow" }
},
{
"type": "invoke",
"invocation": {
"id": "invoke-1",
"nodeId": "node-1",
"command": "product.status",
"params": { "verbose": true },
"timeoutMs": 1000,
"idempotencyKey": "idem-1",
"sessionKey": "agent:main:main"
}
},
{
"type": "result",
"invocationId": "invoke-1",
"result": {
"outcome": "success",
"payload": { "ready": true }
}
},
{ "type": "cancel", "invocationId": "invoke-1" },
{
"type": "status",
"status": {
"state": "ready",
"manifestGeneration": 3,
"runtimeVersion": "1.0.0",
"attempt": 1,
"reason": null
}
},
{
"type": "status",
"status": {
"state": "ready",
"manifestGeneration": 9007199254740991,
"runtimeVersion": "1.0.0",
"attempt": 9007199254740991,
"reason": null
}
}
],
"canonicalJson": [
"{\"type\":\"configure\",\"configuration\":{\"manifestGeneration\":3,\"capabilities\":[\"native.settings\",\"native.status\"],\"commands\":[{\"name\":\"product.settings\"},{\"name\":\"product.status\"}],\"maxConcurrency\":2,\"maxInputBytes\":1024,\"maxOutputBytes\":1024,\"defaultTimeoutMs\":1000,\"maxTimeoutMs\":5000,\"resultGraceMs\":50}}",
"{\"type\":\"configured\",\"manifest\":{\"manifestGeneration\":3,\"capabilities\":[\"native.settings\",\"native.status\"],\"commands\":[\"product.settings\",\"product.status\"]}}",
"{\"type\":\"admission-request\",\"invocation\":{\"id\":\"invoke-1\",\"nodeId\":\"node-1\",\"command\":\"product.status\",\"params\":{\"verbose\":true},\"timeoutMs\":1000,\"idempotencyKey\":\"idem-1\",\"sessionKey\":\"agent:main:main\"}}",
"{\"type\":\"admission-decision\",\"invocationId\":\"invoke-1\",\"decision\":{\"outcome\":\"allow\"}}",
"{\"type\":\"invoke\",\"invocation\":{\"id\":\"invoke-1\",\"nodeId\":\"node-1\",\"command\":\"product.status\",\"params\":{\"verbose\":true},\"timeoutMs\":1000,\"idempotencyKey\":\"idem-1\",\"sessionKey\":\"agent:main:main\"}}",
"{\"type\":\"result\",\"invocationId\":\"invoke-1\",\"result\":{\"outcome\":\"success\",\"payload\":{\"ready\":true}}}",
"{\"type\":\"cancel\",\"invocationId\":\"invoke-1\"}",
"{\"type\":\"status\",\"status\":{\"state\":\"ready\",\"manifestGeneration\":3,\"runtimeVersion\":\"1.0.0\",\"attempt\":1,\"reason\":null}}",
"{\"type\":\"status\",\"status\":{\"state\":\"ready\",\"manifestGeneration\":9007199254740991,\"runtimeVersion\":\"1.0.0\",\"attempt\":9007199254740991,\"reason\":null}}"
]
}