use std::sync::Arc;
use serde_json::{json, Value};
use crate::core::types::Capabilities;
use crate::protocol::connection::CallConnection;
use crate::protocol::wire::{CallError, ResponseEnvelope};
use crate::registry::registration::{
make_handler, Handler, HandlerKind, HandlerRegistration, OperationProvenance, OperationRegistry,
};
use crate::registry::spec::{AccessControl, OperationSpec, OperationType, Visibility};
pub const OP_REGISTER_NAME: &str = "op/register";
#[derive(Debug, Clone)]
pub struct OpRegisterRequest {
pub spec: OperationSpec,
pub replace: bool,
}
impl OpRegisterRequest {
pub fn to_json(&self) -> Value {
json!({
"spec": crate::registry::discovery::spec_to_json_pub(&self.spec),
"replace": self.replace,
})
}
pub fn from_json(value: &Value) -> Result<Self, CallError> {
let spec_json = value
.get("spec")
.ok_or_else(|| CallError::invalid_input("op/register payload missing `spec`"))?;
let name = spec_json
.get("name")
.and_then(|v| v.as_str())
.ok_or_else(|| CallError::invalid_input("op/register spec missing `name`"))?
.to_string();
let spec = crate::client::rebuild_spec_for(spec_json, &name, &None).map_err(|e| {
CallError::invalid_input(format!("op/register spec rebuild failed: {e:?}"))
})?;
let replace = value
.get("replace")
.and_then(|v| v.as_bool())
.unwrap_or(false);
Ok(Self { spec, replace })
}
}
pub fn op_register_spec(access_control: AccessControl) -> OperationSpec {
OperationSpec::new(
OP_REGISTER_NAME,
OperationType::Mutation,
Visibility::External,
json!({
"type": "object",
"properties": {
"spec": { "type": "object" },
"replace": { "type": "boolean" }
},
"required": ["spec"]
}),
json!({
"type": "object",
"properties": {
"name": { "type": "string" },
"registered": { "type": "boolean" }
},
"required": ["name", "registered"]
}),
vec![],
access_control,
None,
)
}
pub fn op_register_handler(
connection: Arc<CallConnection>,
serving_registry: Arc<OperationRegistry>,
) -> Handler {
make_handler(move |input, context| {
let connection = Arc::clone(&connection);
let serving_registry = Arc::clone(&serving_registry);
async move {
let request = match OpRegisterRequest::from_json(&input) {
Ok(r) => r,
Err(e) => return ResponseEnvelope::error(context.request_id, e),
};
if serving_registry.registration(&request.spec.name).is_some() {
return ResponseEnvelope::error(
context.request_id,
CallError::already_exists(format!(
"op/register: `{}` is registered by this side's own serving \
registry; peer-announced ops may not shadow it",
request.spec.name
)),
);
}
if connection.overlay_contains(&request.spec.name) && !request.replace {
return ResponseEnvelope::error(
context.request_id,
CallError::already_exists(format!(
"op/register: `{}` is already registered on this connection; \
set `replace: true` to replace it",
request.spec.name
)),
);
}
let mut spec = request.spec;
spec.visibility = Visibility::Internal;
let remote_name = spec.name.clone();
let handler = forwarding_stub_for_announced_op(
Arc::clone(&connection),
remote_name,
spec.op_type,
);
connection.register_imported(HandlerRegistration::new(
spec.clone(),
handler,
OperationProvenance::FromCall,
None,
None,
Capabilities::new(),
));
ResponseEnvelope::ok(
context.request_id,
json!({ "name": spec.name, "registered": true }),
)
}
})
}
fn forwarding_stub_for_announced_op(
connection: Arc<CallConnection>,
remote_name: String,
op_type: OperationType,
) -> HandlerKind {
match op_type {
OperationType::Query | OperationType::Mutation => HandlerKind::Once(
crate::client::make_forwarding_handler(connection, remote_name),
),
OperationType::Sub | OperationType::Pub => {
HandlerKind::Once(make_handler(|_input, context| async move {
ResponseEnvelope::error(
context.request_id,
CallError::invalid_operation_type(
"peer-announced Sub/Pub ops are not invocable over nested \
composition (request/response only)",
),
)
}))
}
}
}
pub async fn announce_op(
connection: &CallConnection,
spec: OperationSpec,
replace: bool,
) -> ResponseEnvelope {
let request = OpRegisterRequest { spec, replace };
connection
.call_with_payload(serde_json::json!({
"operationId": OP_REGISTER_NAME,
"input": request.to_json(),
}))
.await
}
#[cfg(test)]
mod tests {
use super::*;
use crate::protocol::connection::CallConnection;
use crate::registry::context::OperationContext;
use crate::registry::discovery::install_bootstrap_discovery;
use crate::registry::registration::OperationRegistry;
use crate::registry::spec::Visibility;
use std::collections::HashMap;
fn announced_spec(name: &str) -> OperationSpec {
OperationSpec::new(
name,
OperationType::Query,
Visibility::External,
json!({}),
json!({}),
vec![],
AccessControl::default(),
None,
)
}
fn stub_connection() -> crate::core::types::Connection {
crate::protocol::sink_empty_connection()
}
fn test_context(request_id: &str) -> OperationContext {
OperationContext {
request_id: request_id.to_string(),
parent_request_id: None,
identity: None,
handler_identity: None,
forwarded_for: None,
capabilities: Capabilities::new(),
metadata: HashMap::new(),
scoped_env: crate::registry::context::ScopedPeerEnv::empty(),
env: Arc::new(crate::registry::env::LocalOperationEnv::new(Arc::new(
OperationRegistry::new(),
))),
abort_policy: crate::registry::context::AbortPolicy::default(),
deadline: None,
internal: false,
ownership: None,
}
}
#[test]
fn request_round_trips_through_json() {
let request = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: true,
};
let json = request.to_json();
let parsed = OpRegisterRequest::from_json(&json).expect("parse");
assert_eq!(parsed.spec.name, "worker/exec");
assert_eq!(parsed.spec.op_type, OperationType::Query);
assert!(parsed.replace);
}
#[test]
fn request_missing_spec_is_invalid_input() {
let err = OpRegisterRequest::from_json(&json!({})).unwrap_err();
assert_eq!(err.code, "INVALID_INPUT");
}
#[test]
fn request_missing_name_is_invalid_input() {
let err = OpRegisterRequest::from_json(&json!({ "spec": {} })).unwrap_err();
assert_eq!(err.code, "INVALID_INPUT");
}
#[tokio::test]
async fn handler_registers_announced_op_in_overlay() {
let conn = Arc::new(CallConnection::new(stub_connection()));
let handler = op_register_handler(Arc::clone(&conn), Arc::new(OperationRegistry::new()));
let input = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: false,
}
.to_json();
let response = handler(input, test_context("req-or-1")).await;
assert!(
response.result.is_ok(),
"register succeeded, got {:?}",
response.result
);
let registered = conn.overlay_env().contains("worker/exec");
assert!(registered, "announced op landed in the connection overlay");
assert!(conn.overlay_contains("worker/exec"));
}
#[tokio::test]
async fn handler_rejects_collision_without_replace() {
let conn = Arc::new(CallConnection::new(stub_connection()));
let handler = op_register_handler(Arc::clone(&conn), Arc::new(OperationRegistry::new()));
let input = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: false,
}
.to_json();
let first = handler(input.clone(), test_context("req-or-2a")).await;
assert!(first.result.is_ok());
let second = handler(input, test_context("req-or-2b")).await;
let err = second.result.expect_err("collision rejected");
assert_eq!(err.code, "ALREADY_EXISTS");
}
#[tokio::test]
async fn handler_replaces_with_replace_flag() {
let conn = Arc::new(CallConnection::new(stub_connection()));
let handler = op_register_handler(Arc::clone(&conn), Arc::new(OperationRegistry::new()));
let original = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: false,
}
.to_json();
let first = handler(original, test_context("req-or-3a")).await;
assert!(first.result.is_ok());
let replacement = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: true,
}
.to_json();
let second = handler(replacement, test_context("req-or-3b")).await;
assert!(
second.result.is_ok(),
"replace flag permits re-registration, got {:?}",
second.result
);
}
#[tokio::test]
async fn registered_spec_forced_internal_with_fromcall_provenance() {
let conn = Arc::new(CallConnection::new(stub_connection()));
let handler = op_register_handler(Arc::clone(&conn), Arc::new(OperationRegistry::new()));
let input = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: false,
}
.to_json();
let response = handler(input, test_context("req-or-4")).await;
assert!(response.result.is_ok());
let registration = conn
.overlay_registration("worker/exec")
.expect("registered");
assert_eq!(registration.spec.visibility, Visibility::Internal);
assert_eq!(
registration.provenance,
crate::registry::registration::OperationProvenance::FromCall
);
}
#[tokio::test]
async fn handler_rejects_base_registry_collision_even_with_replace() {
let serving = Arc::new(OperationRegistry::new());
serving
.register(HandlerRegistration::new(
crate::registry::spec::OperationSpec::new(
"fs/readFile",
OperationType::Query,
Visibility::External,
json!({}),
json!({}),
vec![],
AccessControl::default(),
None,
),
HandlerKind::Once(make_handler(|input, ctx| async move {
ResponseEnvelope::ok(ctx.request_id, input)
})),
crate::registry::registration::OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
let conn = Arc::new(CallConnection::new(stub_connection()));
let handler = op_register_handler(Arc::clone(&conn), Arc::clone(&serving));
let input = OpRegisterRequest {
spec: announced_spec("fs/readFile"),
replace: true,
}
.to_json();
let response = handler(input, test_context("req-or-5")).await;
let err = response.result.expect_err("base collision rejected");
assert_eq!(err.code, "ALREADY_EXISTS");
assert!(
!conn.overlay_contains("fs/readFile"),
"the rejected announce never lands in the overlay"
);
assert!(serving.registration("fs/readFile").is_some());
}
#[tokio::test]
async fn handler_rejects_collision_with_internal_serving_op() {
let serving = Arc::new(OperationRegistry::new());
serving
.register(HandlerRegistration::new(
crate::registry::spec::OperationSpec::new(
"internal/vault",
OperationType::Query,
Visibility::Internal,
json!({}),
json!({}),
vec![],
AccessControl::default(),
None,
),
HandlerKind::Once(make_handler(|input, ctx| async move {
ResponseEnvelope::ok(ctx.request_id, input)
})),
crate::registry::registration::OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
let conn = Arc::new(CallConnection::new(stub_connection()));
let handler = op_register_handler(Arc::clone(&conn), Arc::clone(&serving));
let input = OpRegisterRequest {
spec: announced_spec("internal/vault"),
replace: false,
}
.to_json();
let response = handler(input, test_context("req-or-6")).await;
let err = response
.result
.expect_err("internal base collision rejected");
assert_eq!(err.code, "ALREADY_EXISTS");
}
#[tokio::test]
async fn handler_overlay_collision_still_governed_by_replace() {
let serving = Arc::new(OperationRegistry::new());
let conn = Arc::new(CallConnection::new(stub_connection()));
let handler = op_register_handler(Arc::clone(&conn), Arc::clone(&serving));
let first = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: false,
}
.to_json();
assert!(handler(first, test_context("req-or-7a"))
.await
.result
.is_ok());
let collision = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: false,
}
.to_json();
let err = handler(collision, test_context("req-or-7b"))
.await
.result
.expect_err("overlay collision without replace rejected");
assert_eq!(err.code, "ALREADY_EXISTS");
let replacement = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: true,
}
.to_json();
assert!(
handler(replacement, test_context("req-or-7c"))
.await
.result
.is_ok(),
"overlay replace still permitted; base gate is scoped to the serving registry"
);
}
#[tokio::test]
async fn nested_composition_of_base_op_unaffected_by_unrelated_announce() {
let serving = Arc::new(OperationRegistry::new());
serving
.register(HandlerRegistration::new(
crate::registry::spec::OperationSpec::new(
"fs/readFile",
OperationType::Query,
Visibility::External,
json!({}),
json!({}),
vec![],
AccessControl::default(),
None,
),
HandlerKind::Once(make_handler(|_input, ctx| async move {
ResponseEnvelope::ok(ctx.request_id, json!({ "from": "base" }))
})),
crate::registry::registration::OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
let conn = Arc::new(CallConnection::new(stub_connection()));
let handler = op_register_handler(Arc::clone(&conn), Arc::clone(&serving));
let input = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: false,
}
.to_json();
assert!(handler(input, test_context("req-or-8a"))
.await
.result
.is_ok());
assert!(conn.overlay_contains("worker/exec"));
let base = Arc::new(crate::registry::env::LocalOperationEnv::new(Arc::clone(
&serving,
)));
let mut composite = crate::registry::env::PeerCompositeEnv::new(base);
composite.attach_peer("consumer-peer".to_string(), conn.overlay_env());
let env: Arc<dyn crate::registry::env::OperationEnv + Send + Sync> = Arc::new(composite);
let mut ctx = test_context("req-or-8b");
ctx.env = env;
ctx.scoped_env = crate::registry::context::ScopedPeerEnv::new(["fs/readFile"]);
let response = ctx.env.invoke("fs", "readFile", json!({}), &ctx).await;
let out = response.result.expect("base op composes");
assert_eq!(
out,
json!({ "from": "base" }),
"the serving side's own op resolves through composition, not a peer stub"
);
}
#[tokio::test]
async fn announced_op_is_discoverable_via_services_list_peers() {
use crate::registry::env::PeerCompositeEnv;
use crate::registry::{
context::ScopedPeerEnv, discovery::services_list_peers_handler, env::LocalOperationEnv,
};
let serving = Arc::new(OperationRegistry::new());
install_bootstrap_discovery(&serving).expect("bootstrap discovery install");
let peer_identity = crate::core::auth::Identity {
id: "consumer-peer".to_string(),
scopes: vec![],
resources: HashMap::new(),
};
let conn = Arc::new(CallConnection::new_overlay_only(peer_identity));
let handler = op_register_handler(Arc::clone(&conn), Arc::clone(&serving));
let input = OpRegisterRequest {
spec: announced_spec("worker/exec"),
replace: false,
}
.to_json();
assert!(
handler(input, test_context("req-or-9a"))
.await
.result
.is_ok(),
"announce lands in the connection overlay"
);
assert!(conn.overlay_contains("worker/exec"));
let base = Arc::new(LocalOperationEnv::new(Arc::clone(&serving)));
let mut composite = PeerCompositeEnv::new(base);
composite.attach_peer("consumer-peer".to_string(), conn.overlay_env());
let env: Arc<dyn crate::registry::env::OperationEnv + Send + Sync> = Arc::new(composite);
assert!(env.contains("worker/exec"));
let ops = env.peer_operations(&"consumer-peer".to_string());
assert_eq!(
ops,
vec!["worker/exec".to_string()],
"PeerCompositeEnv::peer_operations must surface the peer overlay's announced ops"
);
assert!(
env.peer_operations(&"unknown-peer".to_string()).is_empty(),
"an unattached peer has no operations"
);
let peers_registry = Arc::new(OperationRegistry::new());
install_bootstrap_discovery(&peers_registry).expect("bootstrap install");
let list_handler = services_list_peers_handler(Arc::clone(&peers_registry));
let mut ctx = test_context("req-or-9b");
ctx.env = env;
ctx.scoped_env = ScopedPeerEnv::empty();
let response = list_handler(json!({}), ctx).await;
let out = response.result.expect("list-peers ok");
let peers_arr = out
.get("peers")
.and_then(|v| v.as_array())
.expect("peers array");
let consumer = peers_arr
.iter()
.find(|p| p.get("peer_id").and_then(|v| v.as_str()) == Some("consumer-peer"))
.expect("consumer-peer present in list-peers output");
let names: Vec<&str> = consumer
.get("operations")
.and_then(|v| v.as_array())
.expect("consumer operations array")
.iter()
.filter_map(|o| o.get("name").and_then(|n| n.as_str()))
.collect();
assert!(
names.contains(&"worker/exec"),
"the announced op must be discoverable via services/list-peers (UP-03)"
);
let local = peers_arr
.iter()
.find(|p| p.get("peer_id").and_then(|v| v.as_str()) == Some("local"))
.expect("local peer present");
let local_names: Vec<&str> = local
.get("operations")
.and_then(|v| v.as_array())
.expect("consumer operations array")
.iter()
.filter_map(|o| o.get("name").and_then(|n| n.as_str()))
.collect();
assert!(local_names.contains(&"services/list-peers"));
}
#[test]
fn request_round_trips_flavor_form_channel_open_marker() {
use crate::registry::spec::ChannelOpenSpec;
let request = OpRegisterRequest {
spec: announced_spec("channels/tunnel/direct")
.with_channel_open(ChannelOpenSpec::new("alk/tunnel")),
replace: false,
};
let json = request.to_json();
assert_eq!(json["spec"]["channel_open"], json!(true));
assert_eq!(
json["spec"]["channel_open_alpn"],
json!("alk/tunnel"),
"the flavor form rides the explicit ALPN string"
);
let parsed = OpRegisterRequest::from_json(&json).expect("parse");
let marker = parsed
.spec
.channel_open
.expect("marker survives the announce");
assert_eq!(marker.alpn, "alk/tunnel");
}
}