feat: per-session fork registry, connect-side serving loop, op/register (review 004 Units 1-3)

Remediates all six findings of review 004 (per-connection dispatch
resolution and client-side op serving). All claims re-verified in
source before remediation; F-02's member list gains ScopedPeerEnv
(also Clone — fork surface simpler than estimated).

- OperationRegistry: interior mutability (parking_lot RwLock on both
  maps); register takes &self; registration/list_operations return
  owned clones; fork() deep-copies registrations + cached publish-schema
  validators (F-02/F-03); OperationRegistryBuilder::from_registry.
- install_bootstrap_discovery: services/list, services/list-peers,
  services/schema registered closed over the fork itself, so
  per-session openables are discoverable and services/schema answers
  from the fork (F-06).
- Dispatcher::serve_single_stream: full-duplex single-stream loop —
  call.requested dispatches inbound; responded/completed/error resolve
  outbound pendings; aborted tries both tables (in-flight sink aborts
  + pending cascade); published routes inbound sinks (F-04).
- ChannelClient::from_connection_with_serving(connection,
  Option<ServingConfig>): opt-in serving; from_connection keeps the
  pure-consumer default.
- registry::op_register: OpRegisterRequest wire DTO (spec in
  services/schema JSON + replace flag), op_register_spec,
  op_register_handler (rebuild -> forwarding stub -> register_imported,
  forced Internal/FromCall), announce_op; CallError::already_exists;
  spec_to_json_pub; from_call's rebuild_spec_for + forwarding-handler
  constructors crate-shared (F-05).
- ADR-047 §4 amendment #2: per-session fork is the dispatch-registry
  mechanism; overlay stays nested-invocation/peer-announced landing
  zone (F-01/F-02).
- ADR-022 amendment 2026-09-03: bootstrap-op set (services/list,
  services/schema, op/register), opt-in connect-side serving, op/register
  wire shape (F-04/F-05).
- alkhttp ADR-048 reconciliation note + OQ-05 re-pointed at the alkcall
  ADRs (Unit 1b).
- Review 004 status -> remediated; remediation log with gates.

Verification:
- cargo test: 581 passed, 0 failed (565 baseline + 16 new)
- cargo test --all-features: 598 passed, 0 failed
- cargo clippy --all-targets -- -D warnings: clean
- cargo clippy --all-features --all-targets -- -D warnings: clean
- cargo clippy --target wasm32-unknown-unknown -- -D warnings: clean
- cargo fmt --check: clean
- cargo doc --no-deps: clean

Gates: fork_registry_open_op_resolves_and_is_discoverable (open op via
fork + services/list shows openable + services/schema validates),
serving_loop_hub_to_consumer_call_resolves (hub->consumer call through
consumer's serving loop, consumer->hub still resolves),
op_register_announce_then_hub_call_routes_back_to_consumer (announce ->
overlay -> hub call -> forwarding stub -> consumer serves).
This commit is contained in:
glm-5.3-flash committed 2026-09-03 17:32:45 +00:00
1 parent c0dbf82518
commit f84d214173
18 files changed
+2024 -138

No files matched your search

@@ -2,7 +2,7 @@
## Status ## Status
Accepted (amended 2026-06-26, 2026-07-13, and 2026-07-16 — see "Amendments" below; the 2026-07-16 amendment per ADR-045 §5 removes `CallClient::connect`) Accepted (amended 2026-06-26, 2026-07-13, and 2026-07-16 — see "Amendments" below; the 2026-07-16 amendment per ADR-045 §5 removes `CallClient::connect`; amendment 2026-09-03 — the bootstrap-op set and the connect-side serving loop, see "Amendment (2026-09-03)" below)
## Context ## Context
@@ -359,6 +359,103 @@ same as `from_openapi` receives HTTP credentials.
prior art prior art
- POC at `/workspace/@alkdev/dispatch` — head/worker dispatch over SSH+axum - POC at `/workspace/@alkdev/dispatch` — head/worker dispatch over SSH+axum
## Amendment (2026-09-03): bootstrap-op set and the connect-side serving loop
Review 004 (F-04/F-05) verified two gaps between this ADR's
bidirectionality promise (§2 "connection direction is independent of
call direction") and the single-stream channel-0 implementation:
1. **The connect side's read pump resolved responses only** — inbound
`call.requested` frames were silently dropped. Serving existed only
on the accept side (`run_loop_single_stream`).
2. **No wire mechanism announced client-side ops** — `from_call`
imports hub ops into the consumer, but a connected peer had no way
to announce "here are the ops I serve" over the wire.
The amendment (implemented in alkcall; the e2e gates are
`serving_loop_hub_to_consumer_call_resolves` and
`op_register_announce_then_hub_call_routes_back_to_consumer` in
`src/channels/client.rs`):
### The bootstrap-op set (one-way)
Each side of a channels connection **may serve** the bootstrap ops on
channel 0. The set is closed:
- `services/list` — discovery (the op `from_call` dials on every
import; each side is expected to serve it)
- `services/schema` — per-op schema disclosure
- `op/register` — peer op announcement (below)
These are the ops a peer may assume are reachable (subject to each
op's `AccessControl`); a peer that does not serve one answers
`NOT_FOUND` and the caller treats it accordingly (`from_call` already
surfaces discovery failure as `AdapterError::DiscoveryFailed`).
### Connect-side serving is opt-in (two-way door at the API level)
`ChannelClient::from_connection_with_serving(connection,
Option<ServingConfig>)` — with `None` (the pure-consumer default) the
read pump resolves outbound pendings only (previous behavior). With
`Some(ServingConfig { registry, identity_provider })` the read pump
becomes the full-duplex serving loop (`Dispatcher::serve_single_stream`):
inbound `call.requested` frames dispatch against the configured
registry and resolve back to the peer; outbound pendings still resolve
in the same loop. Serving is opt-in because a pure consumer has no
registry to serve; the *protocol* is symmetric, the *API* is explicit.
Direction disambiguation in the loop is by table membership, not
framing: an id that is one of *our* outbound pendings resolves there;
an id that is an inbound request dispatches; `call.aborted` tries both
tables (in-flight sink aborts and the pending map's cascade). IDs are
UUID-generated per side, so cross-correlation is not a hazard.
### `op/register` (one-way in wire shape)
The peer→hub direction of `from_call`'s bundle flow: a peer sends
`call.requested` for `op/register` carrying the `OperationSpec` in the
`services/schema` wire shape (`spec_to_json`) plus a `replace` flag.
The serving-side handler (`registry::op_register::op_register_handler`):
1. Rebuilds the spec (the same parser `from_call` uses — one spec
serialization on the wire).
2. Wraps a **call-forwarding handler** that issues a nested
`call.requested` back over channel 0 to the announcing peer (the
same shape `from_call`'s imported bundles use — the in-process twin
at `protocol/adapter.rs`).
3. Writes the bundle into **that connection's overlay** via
`register_imported` with `FromCall` provenance and
`Visibility::Internal` (composition material, ADR-017 — never
directly callable from the serving side's own wire).
The announced op is discoverable via `services/list-peers` (the
overlay is peer-keyed, `compose_root_env` attaches it) and invocable
via nested composition (`env.invoke`). Announced Sub/Pub ops register
as stubs that answer `INVALID_OPERATION_TYPE` — nested composition is
request/response-only (`OverlayOperationEnv`'s contract); the
streaming/sink forwarding shapes ride on the `from_call` import path.
Access control: `op/register` itself carries an `AccessControl` (an
unprivileged peer cannot reach the handler — the registry's normal
invoke path enforces it). Replace semantics: a collision with an
existing overlay registration is rejected with `ALREADY_EXISTS`
unless `replace: true` (the reconnect path re-announces). The
overlay dies with the connection (Layer 2), so reconnect re-announce
is naturally scoped.
The envelope kind set stays closed at six — bootstrap ops over channel
0 are the door (AGENTS.md §7 allows adding kinds; none is needed).
### Cross-references
- Review 004 (`docs/reviews/004-per-connection-dispatch-and-client-serving-review.md`)
F-04/F-05 — the verification and the design decision
- ADR-047 §4 amendment #2 (2026-09-03) — the per-session fork the
bootstrap ops compose on
- alkhttp OQ-05 / ADR-048 — the browser data-channel wiring this
unblocks (the alkhttp-side gap was wiring; these two mechanisms are
what it wires to)
## Amendments (2026-06-26) ## Amendments (2026-06-26)
This ADR left four decisions as two-way doors (§1 Consequences flagged DC-1's This ADR left four decisions as two-way doors (§1 Consequences flagged DC-1's
@@ -5,7 +5,99 @@
Accepted (amends ADR-037; refines ADR-044, ADR-046; §4 amended Accepted (amends ADR-037; refines ADR-044, ADR-046; §4 amended
2026-08-13 — open ops are registered per-connection, not resolved via 2026-08-13 — open ops are registered per-connection, not resolved via
`context.env` downcast — see "Amendment (§4 per-connection `context.env` downcast — see "Amendment (§4 per-connection
registration, 2026-08-13)" below) registration, 2026-08-13)"; amendment #2 (2026-09-03) — the
per-connection registration mechanism is the **per-session fork of the
base registry installed as the session's dispatch registry**, not the
connection overlay — see "Amendment (§4 mechanism, 2026-09-03)" below)
## Amendment (§4 mechanism, 2026-09-03)
The 2026-08-13 amendment named the registration target as "the
connection overlay registry (Layer 2 per ADR-019)". Review 004 (F-01)
verified that this shape cannot dispatch: the top-level dispatch path
(`Dispatcher::dispatch` / `run_loop_single_stream`) resolves and
invokes against the dispatcher's **base registry only**; the
connection overlay is reachable solely as a layer of `context.env`
(for nested invocations — a handler calling `env.invoke(...)), never
for resolving the incoming `call.requested` itself. An open op
registered on the overlay resolves `NOT_FOUND` on the wire.
The one shape proven end-to-end (alkcall's own e2e gate) is different:
the `install_channel_zero` hook builds a **fresh per-connection
registry containing the open op and passes it as the dispatcher's base
registry**. This amendment makes that the operative mechanism.
**The decision: per-connection registration happens on a fork of the
deployment's base registry, installed as the session's dispatch
registry.** The `install_channel_zero` hook (and any future
session-establishment seam):
1. **Forks** the deployment's base registry
(`OperationRegistry::fork` — a deep copy carrying handlers,
provenance, composition authority, capabilities, and the cached
publish-schema validators; review 004 F-02/F-03).
2. **Registers the per-session ops on the fork** — the generic channel
ops (`ChannelOperations::register_on`), the openables
(`ChannelCore::register_openable`), and the bootstrap discovery ops
(`install_bootstrap_discovery`, closed over the fork itself so
`services/list` sees the fork's per-session ops — review 004 F-06).
3. **Dispatches channel 0 over the fork** (`Dispatcher::new(fork,
...)`).
The fork is possible because `OperationRegistry` is internally
mutable (`parking_lot::RwLock` around both maps) — a fork shared as an
`Arc<OperationRegistry>` can receive bootstrap ops after the
dispatcher was built, and the self-referential discovery closure sees
every post-install registration.
The **connection overlay (Layer 2) remains what ADR-019/ADR-024
describe**: the landing zone for peer-announced ops (`op/register`,
review 004 F-05 — ADR-022 amendment) and the nested-invocation target
for imported ops. It is not the dispatch-resolution path for the
session's own ops.
Rationale for the fork shape over an overlay-aware dispatch fallback
(F-02 option (b)): the fork is the only shape with an end-to-end
proof, it needs no change to the shared dispatch loop, and it keeps
the overlay's `invoke_with_policy` shape (namespace-scoped,
parent-context-driven — built for nested composition) out of the
top-level call path, where it does not match the frame-handling
contract.
This preserves every invariant the 2026-08-13 amendment protected:
- **Layering (ADR-044):** unchanged — the open-op wrapper is in
`channels-call`; the call crate's `OperationRegistry` gains only
`fork` (and interior mutability), no channels types.
- **Per-connection resolution:** the open op gets the *right*
`ChannelManager` because the fork is built per-connection and its
openable closes over that connection's `ChannelCore`.
- **"Marked ops invoked outside a channels session" (ADR-047 §2):**
unchanged in effect — a `channels/<alpn>/sub` op registered only on
a session fork is not reachable on a bare `alk/call` connection (the
fork isn't that session's dispatch registry) — the dispatch path
returns `NOT_FOUND`.
### Door type
**Two-way (implementation detail), as before.** The registration
*target* mechanism (fork as base registry) sits within the same
wrapper-shape detail the 2026-08-13 amendment already marked two-way.
The one-way decisions (per-ALPN op names, the `channel_open` marker,
removal of `channel/open`/`direction`) are unchanged.
### References
- Review 004 F-01/F-02/F-03/F-06
(`docs/reviews/004-per-connection-dispatch-and-client-serving-review.md`)
— the verification and the mechanism decision
- ADR-022 amendment (2026-09-03) — the bootstrap-op set (`services/list`,
`services/schema`, `op/register`) and the connect-side serving loop
- ADR-019: operation registry layering (the overlay stays the nested
invocation / peer-announced-ops landing zone)
- The e2e gate: `fork_registry_open_op_resolves_and_is_discoverable`
(`src/channels/client.rs`) — open op resolves through the fork,
per-session openable in `services/list`, `services/schema` validates
## Amendment (§4 per-connection registration, 2026-08-13) ## Amendment (§4 per-connection registration, 2026-08-13)
@@ -2,7 +2,9 @@
## Status ## Status
Verified, open for remediation. Verified, remediated (2026-09-03 — Units 1–3; Unit 4 remains
downstream in alkhttp). See "Remediation log (2026-09-03)" at the
bottom for the landing summary and the gates.
## Scope ## Scope
@@ -444,4 +446,82 @@ notes).
`peer_operations()`, populated by `compose_root_env` at `peer_operations()`, populated by `compose_root_env` at
`dispatch.rs:211-213`). `dispatch.rs:211-213`).
- Baseline gates re-run for this pass: cargo test (565), clippy - Baseline gates re-run for this pass: cargo test (565), clippy
(all-targets), fmt — all clean. No source changes. (all-targets), fmt — all clean. No source changes.
## Remediation log (2026-09-03)
Units 1–3 landed in one pass. Every claim above was re-verified
against the source before remediation began; all six findings
confirmed (F-02's member-type list was missing `ScopedPeerEnv`
(`context.rs:53`) — also Clone, which made the fork surface simpler
than the review estimated).
**Unit 1 — decision recording (F-01/F-02 direction, F-04/F-05 scope):**
- ADR-047 §4 amendment #2 (2026-09-03): the per-connection registration
mechanism is the **per-session fork of the base registry installed as
the session's dispatch registry** (option (a)); the overlay registry
stays the nested-invocation / peer-announced-ops landing zone.
- ADR-022 amendment (2026-09-03): the bootstrap-op set (`services/list`,
`services/schema`, `op/register`), the opt-in connect-side serving
loop, and the `op/register` wire shape (spec in `services/schema`
JSON + `replace` flag; forwarding-handler wrap into the connection
overlay; `ALREADY_EXISTS` collision gate).
- alkhttp ADR-048's reconciliation note and OQ-05 re-pointed at these
ADRs (Unit 1b).
**Unit 2 — fork surface + per-session composition (F-02, F-03, F-06):**
- `OperationRegistry` gained interior mutability (`parking_lot::RwLock`
around the operations and validators maps) so a registry can live
behind an `Arc` and be self-referential. `register` now takes
`&self`; `registration`/`list_operations` return owned clones
(handlers are `Arc` closures — the clone is cheap).
- `OperationRegistry::fork()` — deep copy of registrations (handlers,
provenance, composition authority, capabilities, `scoped_env`) and
the cached publish-schema validators. `OperationRegistryBuilder::
from_registry` seeds the builder path from an existing registry.
- `install_bootstrap_discovery(&Arc<OperationRegistry>)` — registers
`services/list` / `services/list-peers` / `services/schema` closed
over the fork itself, so per-session openables are discoverable
(F-06) and `services/schema` answers from the fork.
- Gate: `fork_registry_open_op_resolves_and_is_discoverable`
(`src/channels/client.rs`) — open op registered on the fork resolves
over a live channels connection; the openable appears in
`services/list` on the fork; `services/schema` on the fork still
answers. Unit tests: fork independence both directions, validator
carry, post-dispatch registration through a shared `Arc`.
**Unit 3 — client serving half + bootstrap registration (F-04, F-05):**
- `Dispatcher::serve_single_stream` — the full-duplex single-stream
loop: `call.requested` → dispatch + response frames (the
`run_loop_single_stream` arms); `call.responded`/`completed`/`error`
→ pending resolution (the read-pump arms); `call.aborted` → both
tables (in-flight sink aborts and the pending map's cascade);
`call.published` → inbound in-flight sinks. Disconnect fails
pendings and drops in-flight sinks (teardown identical to the two
half-loops it composes).
- `ChannelClient::from_connection_with_serving(connection,
Option<ServingConfig>)` — opt-in serving; `from_connection`
preserves the pure-consumer default (resolution-only read pump).
Gate: `serving_loop_hub_to_consumer_call_resolves` — hub→consumer
call resolves through the consumer's serving loop, consumer→hub
still resolves in the same loop.
- `registry::op_register` — `OpRegisterRequest` wire DTO,
`op_register_spec`, `op_register_handler` (rebuild → forwarding stub
→ `register_imported` into the connection overlay, forced
`Internal`/`FromCall`), `announce_op`. `CallError::already_exists`
added for the collision gate. `spec_to_json_pub` made the
`services/schema` wire shape public; `rebuild_spec_for` and the
forwarding-handler constructors are crate-shared with `from_call`.
Gates: `op_register_announce_then_hub_call_routes_back_to_consumer`
(announce → overlay → hub call → forwarding stub → consumer serves)
plus unit tests (round-trip, collision, replace, forced
visibility/provenance).
**Baseline after remediation:** cargo test 581 passed / 0 failed;
clippy (all-targets, `-D warnings`) clean; `cargo fmt --check` clean;
wasm target gate re-run below. No wire-format changes: the envelope
kind set stays closed at six; the bootstrap ops ride channel 0's call
registry.
+568 -23
View File
@@ -26,6 +26,7 @@ use tokio::sync::Mutex;
use crate::core::types::{Connection, StreamError}; use crate::core::types::{Connection, StreamError};
use crate::protocol::connection::CallConnection; use crate::protocol::connection::CallConnection;
use crate::protocol::wire::ResponseEnvelope; use crate::protocol::wire::ResponseEnvelope;
use crate::registry::registration::OperationRegistry;
use super::manager::{ChannelManager, ChannelSide}; use super::manager::{ChannelManager, ChannelSide};
use super::mux::MuxRunner; use super::mux::MuxRunner;
@@ -43,7 +44,20 @@ use super::reassembly::{MpscRecvStream, MpscSendStream};
/// (`channels/<alpn>/sub`, `channels/<alpn>/pub`) on channel 0. /// (`channels/<alpn>/sub`, `channels/<alpn>/pub`) on channel 0.
pub struct ChannelClient { pub struct ChannelClient {
manager: ChannelManager, manager: ChannelManager,
call_connection: Arc<Mutex<Option<CallConnection>>>, call_connection: Arc<Mutex<Option<Arc<CallConnection>>>>,
}
/// Opt-in serving configuration for a `ChannelClient` (review 004
/// F-04). A pure consumer has no registry to serve — `from_connection`
/// keeps the resolution-only read pump. A consumer that also serves
/// ops passes its registry (and identity provider) via
/// [`ServingConfig`]; the read pump becomes the full-duplex serving
/// loop (`Dispatcher::serve_single_stream`), so inbound
/// `call.requested` frames from the peer dispatch and resolve instead
/// of being dropped.
pub struct ServingConfig {
pub registry: Arc<OperationRegistry>,
pub identity_provider: Arc<dyn crate::core::auth::IdentityProvider>,
} }
impl ChannelClient { impl ChannelClient {
@@ -60,7 +74,21 @@ impl ChannelClient {
/// the `PendingRequestMap` via `dispatch_envelope`, resolving /// the `PendingRequestMap` via `dispatch_envelope`, resolving
/// pending calls. The `CallConnection` holds the framed write /// pending calls. The `CallConnection` holds the framed write
/// half so `call_open_op` can write `call.requested` frames. /// half so `call_open_op` can write `call.requested` frames.
pub async fn from_connection(connection: Connection) -> Result<Self, StreamError> { ///
/// Serving (review 004 F-04): with `serving: None` (the default)
/// the read pump resolves responses only — inbound
/// `call.requested` frames are dropped (a pure consumer has no
/// registry to serve). With `Some(ServingConfig { .. })` the read
/// pump is the full-duplex serving loop
/// (`Dispatcher::serve_single_stream`): inbound requests dispatch
/// against the configured registry and resolve back to the peer,
/// and outbound pendings still resolve. This is the opt-in that
/// makes the connect side a serving half on channel 0 (ADR-022 §2
/// — both sides can be both).
pub async fn from_connection_with_serving(
connection: Connection,
serving: Option<ServingConfig>,
) -> Result<Self, StreamError> {
let remote_addr = connection.remote_addr(); let remote_addr = connection.remote_addr();
let bidi = connection.accept_bi().await?; let bidi = connection.accept_bi().await?;
let (reader, writer) = tokio::io::split(bidi); let (reader, writer) = tokio::io::split(bidi);
@@ -98,17 +126,39 @@ impl ChannelClient {
let (single_stream_writer, single_stream_reader) = let (single_stream_writer, single_stream_reader) =
crate::protocol::connection::split_single_stream(channel0_bidi); crate::protocol::connection::split_single_stream(channel0_bidi);
let call_connection = let call_connection = Arc::new(CallConnection::new_single_stream(
CallConnection::new_single_stream(channel0_conn, single_stream_writer); channel0_conn,
Arc::clone(&single_stream_writer),
));
let pending_map = Arc::clone(call_connection.pending()); let pending_map = Arc::clone(call_connection.pending());
let _read_pump = tokio::spawn(async move { match serving {
crate::protocol::connection::read_single_stream_until_closed( None => {
single_stream_reader, let _read_pump = tokio::spawn(async move {
&pending_map, crate::protocol::connection::read_single_stream_until_closed(
) single_stream_reader,
.await; &pending_map,
}); )
.await;
});
}
Some(config) => {
let dispatcher = crate::protocol::dispatch::Dispatcher::new(
config.registry,
config.identity_provider,
);
let call_conn_for_loop = Arc::clone(&call_connection);
let _serve_loop = tokio::spawn(async move {
dispatcher
.serve_single_stream(
call_conn_for_loop,
single_stream_reader,
single_stream_writer,
)
.await;
});
}
}
// The demux loop ends on transport EOF, clearing the channel // The demux loop ends on transport EOF, clearing the channel
// map (REQ-CH-02). The connect side passes `policy: None` (it // map (REQ-CH-02). The connect side passes `policy: None` (it
@@ -129,6 +179,13 @@ impl ChannelClient {
}) })
} }
/// Construct from an established `Connection` with serving
/// disabled (the pure-consumer default). See
/// [`ChannelClient::from_connection_with_serving`].
pub async fn from_connection(connection: Connection) -> Result<Self, StreamError> {
Self::from_connection_with_serving(connection, None).await
}
/// The `ChannelManager` — for relay logic and tests. /// The `ChannelManager` — for relay logic and tests.
pub fn manager(&self) -> &ChannelManager { pub fn manager(&self) -> &ChannelManager {
&self.manager &self.manager
@@ -193,8 +250,11 @@ impl ChannelClient {
/// Take the `CallConnection` — used by the consumer to register /// Take the `CallConnection` — used by the consumer to register
/// imported ops (`from_call`) on the connection's overlay. After /// imported ops (`from_call`) on the connection's overlay. After
/// this, `call_open_op` returns an error (the connection is owned /// this, `call_open_op` returns an error (the connection is owned
/// by the consumer). /// by the consumer). The `Arc` is shared with the serving loop when
pub async fn take_call_connection(&self) -> Option<CallConnection> { /// serving is enabled; taking it detaches the client's own calling
/// surface, not the serving loop's dispatch (the loop holds its own
/// `Arc` clone).
pub async fn take_call_connection(&self) -> Option<Arc<CallConnection>> {
self.call_connection.lock().await.take() self.call_connection.lock().await.take()
} }
} }
@@ -329,7 +389,7 @@ mod tests {
/// on the accept side). /// on the accept side).
#[tokio::test] #[tokio::test]
async fn channel_0_end_to_end_call_round_trip() { async fn channel_0_end_to_end_call_round_trip() {
let mut registry = crate::registry::registration::OperationRegistry::new(); let registry = crate::registry::registration::OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_query_spec("echo/run"), external_query_spec("echo/run"),
@@ -428,7 +488,7 @@ mod tests {
async fn channel_0_end_to_end_publish_delivers_chunks() { async fn channel_0_end_to_end_publish_delivers_chunks() {
use futures::stream::StreamExt; use futures::stream::StreamExt;
let mut registry = crate::registry::registration::OperationRegistry::new(); let registry = crate::registry::registration::OperationRegistry::new();
let counting_sink = crate::registry::registration::make_sink_handler( let counting_sink = crate::registry::registration::make_sink_handler(
|_input, ctx, mut stream| async move { |_input, ctx, mut stream| async move {
let mut count = 0u32; let mut count = 0u32;
@@ -559,7 +619,7 @@ mod tests {
ResponseEnvelope::ok(ctx.request_id, serde_json::json!({ "ok": oks, "err": err })) ResponseEnvelope::ok(ctx.request_id, serde_json::json!({ "ok": oks, "err": err }))
}, },
); );
let mut registry = crate::registry::registration::OperationRegistry::new(); let registry = crate::registry::registration::OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
OperationSpec::new( OperationSpec::new(
@@ -667,7 +727,7 @@ mod tests {
ResponseEnvelope::ok(ctx.request_id, serde_json::json!({ "ok": oks, "err": err })) ResponseEnvelope::ok(ctx.request_id, serde_json::json!({ "ok": oks, "err": err }))
}, },
); );
let mut registry = crate::registry::registration::OperationRegistry::new(); let registry = crate::registry::registration::OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
OperationSpec::new( OperationSpec::new(
@@ -797,7 +857,7 @@ mod tests {
let (writer, reader) = split_single_stream(channel0_bidi); let (writer, reader) = split_single_stream(channel0_bidi);
let core = ChannelCore::new(manager, policy); let core = ChannelCore::new(manager, policy);
let mut registry = crate::registry::registration::OperationRegistry::new(); let registry = crate::registry::registration::OperationRegistry::new();
let spec = OperationSpec::new( let spec = OperationSpec::new(
"channels/tty/sub", "channels/tty/sub",
OperationType::Sub, OperationType::Sub,
@@ -823,7 +883,7 @@ mod tests {
core.register_openable( core.register_openable(
spec, spec,
Arc::clone(&open_handler), Arc::clone(&open_handler),
&mut registry, &registry,
auth.clone(), auth.clone(),
) )
.expect("register_openable"); .expect("register_openable");
@@ -1043,8 +1103,8 @@ mod tests {
}); });
let manager = ChannelManager::with_defaults(handle, None); let manager = ChannelManager::with_defaults(handle, None);
let ops = ChannelOperations::new(manager, Arc::new(NoCap)); let ops = ChannelOperations::new(manager, Arc::new(NoCap));
let mut registry = crate::registry::registration::OperationRegistry::new(); let registry = crate::registry::registration::OperationRegistry::new();
ops.register_on(&mut registry).expect("register"); ops.register_on(&registry).expect("register");
let handler = registry let handler = registry
.registration("channel/close") .registration("channel/close")
@@ -1117,7 +1177,7 @@ mod tests {
}; };
let (writer, reader) = split_single_stream(channel0_bidi); let (writer, reader) = split_single_stream(channel0_bidi);
let core = ChannelCore::new(manager, Arc::new(NoCap)); let core = ChannelCore::new(manager, Arc::new(NoCap));
let mut registry = crate::registry::registration::OperationRegistry::new(); let registry = crate::registry::registration::OperationRegistry::new();
let spec = OperationSpec::new( let spec = OperationSpec::new(
"channels/tty/sub", "channels/tty/sub",
OperationType::Sub, OperationType::Sub,
@@ -1139,7 +1199,7 @@ mod tests {
core.register_openable( core.register_openable(
spec, spec,
Arc::clone(&open_handler), Arc::clone(&open_handler),
&mut registry, &registry,
auth.clone(), auth.clone(),
) )
.expect("register_openable"); .expect("register_openable");
@@ -1195,4 +1255,489 @@ mod tests {
"handler received ping through adopted channel" "handler received ping through adopted channel"
); );
} }
// --- review 004 Unit 3 acceptance gates (F-04 serving half) -----------
/// F-04 gate 1: hub→consumer call over an existing `ChannelClient`
/// session resolves. The consumer (connect side) dials with
/// `Some(ServingConfig { .. })`; the accept side's dispatcher later
/// calls an op the consumer serves (by opening a fresh stream from
/// its side is not available on channel 0 — the accept side writes
/// a `call.requested` through its channel-0 writer, and the
/// consumer's serving loop dispatches it and writes the response,
/// which resolves on the accept side's pending).
///
/// Mechanically: the accept side's `install_channel_zero` hook
/// returns a clone of the channel-0 `CallConnection` back to the
/// test via a channel; the test then calls
/// `call_open_op("consumer/echo")` — wait, that calls *through* the
/// accept side's registry. The consumer-serving direction needs the
/// accept side to *initiate*: it writes `call.requested` for
/// `consumer/echo` on channel 0. The client's serving loop (with
/// the `consumer/echo` op in its registry) dispatches it and writes
/// the response, resolving the accept side's pending.
#[tokio::test]
async fn serving_loop_hub_to_consumer_call_resolves() {
// Consumer-side (connect side) registry: serves `consumer/echo`.
let consumer_registry = crate::registry::registration::OperationRegistry::new();
consumer_registry
.register(HandlerRegistration::new(
external_query_spec("consumer/echo"),
HandlerKind::Once(make_handler(|input, ctx| async move {
ResponseEnvelope::ok(ctx.request_id, input)
})),
OperationProvenance::Local,
None,
None,
crate::core::types::Capabilities::new(),
))
.unwrap();
let consumer_registry = Arc::new(consumer_registry);
// Accept side: a plain echo registry plus a handle back to the
// accept side's channel-0 `CallConnection` so the test can
// initiate a call *from* the accept side.
let (accept_conn_tx, mut accept_conn_rx) =
tokio::sync::mpsc::channel::<Arc<CallConnection>>(1);
let accept_registry = crate::registry::registration::OperationRegistry::new();
accept_registry
.register(HandlerRegistration::new(
external_query_spec("accept/echo"),
HandlerKind::Once(make_handler(|input, ctx| async move {
ResponseEnvelope::ok(ctx.request_id, input)
})),
OperationProvenance::Local,
None,
None,
crate::core::types::Capabilities::new(),
))
.unwrap();
let accept_registry = Arc::new(accept_registry);
let install_hook: crate::channels::adapter::InstallChannelZero =
Arc::new(move |manager, channel0_conn, _auth| {
let accept_registry = Arc::clone(&accept_registry);
let accept_conn_tx = accept_conn_tx.clone();
tokio::spawn(async move {
let _ = manager;
let channel0_bidi = match channel0_conn.accept_bi().await {
Ok(s) => s,
Err(_) => return,
};
let (writer, reader) = split_single_stream(channel0_bidi);
let call_connection = Arc::new(CallConnection::new_single_stream(
channel0_conn,
Arc::clone(&writer),
));
let _ = accept_conn_tx.send(Arc::clone(&call_connection)).await;
let dp = Dispatcher::new(
accept_registry,
Arc::new(NoopIdProvider) as Arc<dyn IdentityProvider>,
);
// The serving loop on the accept side too: the
// consumer may call back over the same session.
dp.serve_single_stream(call_connection, reader, writer)
.await;
})
});
let (client_end, server_end) = tokio::io::duplex(64 * 1024);
let client_conn =
Connection::from_bidi(client_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
let server_conn =
Connection::from_bidi(server_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
let adapter = ChannelsAdapter::new(install_hook, Arc::new(NoCap));
let auth = AuthContext::anonymous(b"alk/channels");
let _server_handle = tokio::spawn(async move {
let _ = crate::core::types::ProtocolHandler::handle(&adapter, server_conn, &auth).await;
});
let client = ChannelClient::from_connection_with_serving(
client_conn,
Some(ServingConfig {
registry: Arc::clone(&consumer_registry),
identity_provider: Arc::new(NoopIdProvider),
}),
)
.await
.expect("channel client init");
let accept_conn = accept_conn_rx
.recv()
.await
.expect("accept side channel-0 connection handle");
// Accept side (hub) initiates a call to an op the consumer
// serves — the exact direction F-04 says was dropped.
let response = tokio::time::timeout(
std::time::Duration::from_secs(5),
accept_conn.call("consumer/echo", serde_json::json!({ "from": "hub" })),
)
.await
.expect("hub→consumer call timed out");
assert!(
response.result.is_ok(),
"hub→consumer call resolves through the consumer's serving loop, got {:?}",
response.result
);
assert_eq!(
response.result.unwrap(),
serde_json::json!({ "from": "hub" })
);
// The consumer's own calling half still works (outbound pendings
// resolve in the same loop).
let outbound = client
.call_open_op("accept/echo", serde_json::json!({ "to": "hub" }))
.await;
assert!(
outbound.result.is_ok(),
"consumer→hub call still resolves, got {:?}",
outbound.result
);
}
// --- review 004 Unit 3 acceptance gates (F-05 op/register) ------------
/// F-05 gate: the consumer announces an op over channel 0
/// (`op/register` served on the accept side's fork), the hub
/// registers it in the connection overlay, and the hub then calls
/// it back — the forwarding stub issues a nested `call.requested`
/// over channel 0 to the consumer, whose serving loop dispatches
/// the real handler. Peer-announced op, discoverable and callable,
/// end-to-end.
#[tokio::test]
async fn op_register_announce_then_hub_call_routes_back_to_consumer() {
// Consumer side: serves `consumer/exec` locally; announces it.
let consumer_registry = crate::registry::registration::OperationRegistry::new();
consumer_registry
.register(HandlerRegistration::new(
external_query_spec("consumer/exec"),
HandlerKind::Once(make_handler(|input, ctx| async move {
ResponseEnvelope::ok(
ctx.request_id,
serde_json::json!({ "ran_on": "consumer", "input": input }),
)
})),
OperationProvenance::Local,
None,
None,
crate::core::types::Capabilities::new(),
))
.unwrap();
let consumer_registry = Arc::new(consumer_registry);
// Accept side: serves the `op/register` bootstrap op composed on
// a fork (the F-02(a) shape — the fork is fresh here, so the
// bootstrap set is just `op/register`).
let (accept_conn_tx, mut accept_conn_rx) =
tokio::sync::mpsc::channel::<Arc<CallConnection>>(1);
let accept_registry_for_hook =
Arc::new(crate::registry::registration::OperationRegistry::new());
let accept_tx_for_hook = accept_conn_tx.clone();
let install_hook: crate::channels::adapter::InstallChannelZero =
Arc::new(move |_manager, channel0_conn, _auth| {
let accept_registry = Arc::clone(&accept_registry_for_hook);
let accept_conn_tx = accept_tx_for_hook.clone();
tokio::spawn(async move {
let channel0_bidi = match channel0_conn.accept_bi().await {
Ok(s) => s,
Err(_) => return,
};
let (writer, reader) = split_single_stream(channel0_bidi);
let call_connection = Arc::new(CallConnection::new_single_stream(
channel0_conn,
Arc::clone(&writer),
));
let _ = accept_conn_tx.send(Arc::clone(&call_connection)).await;
let reg = accept_registry.fork();
reg.register(HandlerRegistration::new(
crate::registry::op_register::op_register_spec(AccessControl::default()),
HandlerKind::Once(crate::registry::op_register::op_register_handler(
Arc::clone(&call_connection),
)),
OperationProvenance::Local,
None,
None,
crate::core::types::Capabilities::new(),
))
.unwrap();
let dp = Dispatcher::new(
Arc::new(reg),
Arc::new(NoopIdProvider) as Arc<dyn IdentityProvider>,
);
dp.serve_single_stream(call_connection, reader, writer)
.await;
})
});
let (client_end, server_end) = tokio::io::duplex(64 * 1024);
let client_conn =
Connection::from_bidi(client_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
let server_conn =
Connection::from_bidi(server_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
let adapter = ChannelsAdapter::new(install_hook, Arc::new(NoCap));
let auth = AuthContext::anonymous(b"alk/channels");
let _server_handle = tokio::spawn(async move {
let _ = crate::core::types::ProtocolHandler::handle(&adapter, server_conn, &auth).await;
});
let client = ChannelClient::from_connection_with_serving(
client_conn,
Some(ServingConfig {
registry: Arc::clone(&consumer_registry),
identity_provider: Arc::new(NoopIdProvider),
}),
)
.await
.expect("channel client init");
let accept_conn = accept_conn_rx
.recv()
.await
.expect("accept side channel-0 connection handle");
// Consumer announces `consumer/exec` over channel 0.
let announce = tokio::time::timeout(
std::time::Duration::from_secs(10),
client.call_open_op(
crate::registry::op_register::OP_REGISTER_NAME,
crate::registry::op_register::OpRegisterRequest {
spec: external_query_spec("consumer/exec"),
replace: false,
}
.to_json(),
),
)
.await
.expect("op/register announce timed out");
assert!(
announce.result.is_ok(),
"announce ok, got {:?}",
announce.result
);
// The hub's overlay now holds the announced op; the hub calls
// it — the forwarding stub routes the nested call back over
// channel 0 to the consumer's serving loop.
let response = tokio::time::timeout(
std::time::Duration::from_secs(10),
accept_conn.call("consumer/exec", serde_json::json!({ "task": "greet" })),
)
.await
.expect("hub call to peer-announced op timed out");
assert!(
response.result.is_ok(),
"hub→peer-announced-op call resolves, got {:?}",
response.result
);
assert_eq!(
response.result.unwrap(),
serde_json::json!({ "ran_on": "consumer", "input": { "task": "greet" } })
);
}
// --- review 004 Unit 2 acceptance gate (F-02/F-06 fork) ----------------
/// F-02/F-06 gate: the open op is registered on a **fork** of the
/// base registry (base carries a placeholder op; the fork adds the
/// openable + bootstrap discovery via
/// `install_bootstrap_discovery`), the session dispatches over the
/// fork, and — after the open op resolves over the live channels
/// connection — the per-session openable is discoverable through
/// `services/list` **on the fork** (the F-06 self-referential
/// closure) and `services/schema` on the fork still validates.
#[tokio::test]
async fn fork_registry_open_op_resolves_and_is_discoverable() {
use crate::channels::operations::{ChannelCore, OpenHandler};
use crate::channels::policy::PerIdentityChannelPolicy;
use crate::registry::spec::ChannelOpenSpec;
// Base registry: a plain op + bootstrap discovery closed over
// the *base* (which will NOT see per-session openables — that
// is the F-06 finding; the fork's own discovery does).
let base = crate::registry::registration::OperationRegistry::new();
base.register(HandlerRegistration::new(
external_query_spec("base/ping"),
HandlerKind::Once(make_handler(|input, ctx| async move {
ResponseEnvelope::ok(ctx.request_id, input)
})),
OperationProvenance::Local,
None,
None,
crate::core::types::Capabilities::new(),
))
.unwrap();
let policy: Arc<PerIdentityChannelPolicy> = Arc::new(PerIdentityChannelPolicy::new(256));
let policy_for_hook: Arc<dyn ChannelLifecyclePolicy> =
Arc::clone(&policy) as Arc<dyn ChannelLifecyclePolicy>;
let open_handler: OpenHandler = Arc::new(|_input, _channel_conn, _auth| {
tokio::spawn(async move {
tokio::time::sleep(std::time::Duration::from_millis(50)).await;
})
});
let install_hook: crate::channels::adapter::InstallChannelZero =
Arc::new(move |manager, channel0_conn, auth| {
let open_handler = Arc::clone(&open_handler);
let policy = Arc::clone(&policy_for_hook);
let base = Arc::new(crate::registry::registration::OperationRegistry::new());
base.register(HandlerRegistration::new(
external_query_spec("base/ping"),
HandlerKind::Once(make_handler(|input, ctx| async move {
ResponseEnvelope::ok(ctx.request_id, input)
})),
OperationProvenance::Local,
None,
None,
crate::core::types::Capabilities::new(),
))
.unwrap();
let base = Arc::clone(&base);
tokio::spawn(async move {
let channel0_bidi = match channel0_conn.accept_bi().await {
Ok(s) => s,
Err(_) => return,
};
let (writer, reader) = split_single_stream(channel0_bidi);
let core = ChannelCore::new(manager, policy);
// The F-02(a) composition: fork the base, register
// the openable on the fork, install bootstrap
// discovery on the fork (F-06 — closed over the
// fork), dispatch over the fork.
let fork = base.fork();
core.register_openable(
OperationSpec::new(
"channels/tty/sub",
OperationType::Sub,
Visibility::External,
serde_json::json!({
"type": "object",
"properties": {
"container": { "type": "string" }
},
"required": ["container"]
}),
serde_json::json!({
"type": "object",
"properties": {
"channel_id": { "type": "integer" }
}
}),
vec![],
AccessControl::default(),
None,
)
.with_channel_open(ChannelOpenSpec::new("alk/tty")),
Arc::clone(&open_handler),
&fork,
auth.clone(),
)
.expect("register_openable on fork");
let fork = Arc::new(fork);
crate::registry::discovery::install_bootstrap_discovery(&fork)
.expect("bootstrap discovery on fork");
let provider: Arc<dyn IdentityProvider> = Arc::new(NoopIdProvider);
let call_connection = Arc::new(CallConnection::new_single_stream(
channel0_conn,
Arc::clone(&writer),
));
let dp = Dispatcher::new(fork, provider);
dp.run_loop_single_stream(call_connection, reader, writer)
.await;
})
});
let (client_end, server_end) = tokio::io::duplex(64 * 1024);
let client_conn =
Connection::from_bidi(client_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
let server_conn =
Connection::from_bidi(server_end, b"alk/channels".to_vec(), Some(TEST_ADDR));
let adapter = ChannelsAdapter::new(install_hook, Arc::new(NoCap));
let auth = AuthContext::anonymous(b"alk/channels");
let _server_handle = tokio::spawn(async move {
let _ = crate::core::types::ProtocolHandler::handle(&adapter, server_conn, &auth).await;
});
let client = ChannelClient::from_connection(client_conn)
.await
.expect("channel client init");
// 1. The open op resolves over the live connection through the fork.
let response = tokio::time::timeout(
std::time::Duration::from_secs(10),
client.call_open_op(
"channels/tty/sub",
serde_json::json!({ "container": "abc" }),
),
)
.await
.expect("open op timed out");
assert!(
response.result.is_ok(),
"open op resolves through the fork, got {:?}",
response.result
);
let channel_id = response
.result
.unwrap()
.get("channel_id")
.and_then(|v| v.as_u64())
.expect("channel_id in response");
assert!(channel_id > 0);
// 2. The per-session openable is discoverable via services/list
// (F-06) — the fork's own discovery closure sees the fork's ops.
let listing = tokio::time::timeout(
std::time::Duration::from_secs(10),
client.call_open_op("services/list", serde_json::json!({})),
)
.await
.expect("services/list timed out");
assert!(
listing.result.is_ok(),
"services/list resolves, got {:?}",
listing.result
);
let names: Vec<String> = listing
.result
.unwrap()
.get("operations")
.and_then(|v| v.as_array())
.expect("operations array")
.iter()
.filter_map(|o| o.get("name").and_then(|n| n.as_str().map(String::from)))
.collect();
assert!(
names.contains(&"channels/tty/sub".to_string()),
"per-session openable discoverable through the fork's discovery: {names:?}"
);
// 3. services/schema on the fork still validates input (the
// bootstrap op is on the fork and answers).
let schema = tokio::time::timeout(
std::time::Duration::from_secs(10),
client.call_open_op(
"services/schema",
serde_json::json!({ "name": "channels/tty/sub" }),
),
)
.await
.expect("services/schema timed out");
assert!(
schema.result.is_ok(),
"services/schema resolves on the fork, got {:?}",
schema.result
);
}
} }
+14 -14
View File
@@ -62,7 +62,7 @@ impl ChannelOperations {
/// Register the three generic ops on the call `OperationRegistry`. /// Register the three generic ops on the call `OperationRegistry`.
/// The per-ALPN open ops are registered separately by the ALPN /// The per-ALPN open ops are registered separately by the ALPN
/// crates via [`ChannelCore::register_openable`] (ADR-047 §3). /// crates via [`ChannelCore::register_openable`] (ADR-047 §3).
pub fn register_on(&self, registry: &mut OperationRegistry) -> Result<(), String> { pub fn register_on(&self, registry: &OperationRegistry) -> Result<(), String> {
let manager = self.manager.clone(); let manager = self.manager.clone();
let policy = Arc::clone(&self.policy); let policy = Arc::clone(&self.policy);
registry.register(HandlerRegistration::new( registry.register(HandlerRegistration::new(
@@ -400,7 +400,7 @@ impl ChannelCore {
&self, &self,
spec: OperationSpec, spec: OperationSpec,
open_handler: OpenHandler, open_handler: OpenHandler,
registry: &mut OperationRegistry, registry: &OperationRegistry,
auth: AuthContext, auth: AuthContext,
) -> Result<(), String> { ) -> Result<(), String> {
let op_type = spec.op_type; let op_type = spec.op_type;
@@ -765,8 +765,8 @@ mod tests {
async fn register_on_registers_three_ops() { async fn register_on_registers_three_ops() {
let manager = make_manager().await; let manager = make_manager().await;
let ops = ChannelOperations::with_default_policy(manager); let ops = ChannelOperations::with_default_policy(manager);
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
ops.register_on(&mut registry).expect("register"); ops.register_on(&registry).expect("register");
assert!(registry.registration(OP_CHANNEL_CLOSE).is_some()); assert!(registry.registration(OP_CHANNEL_CLOSE).is_some());
assert!(registry.registration(OP_CHANNEL_CONTROL).is_some()); assert!(registry.registration(OP_CHANNEL_CONTROL).is_some());
assert!(registry assert!(registry
@@ -899,8 +899,8 @@ mod tests {
let manager = make_manager().await; let manager = make_manager().await;
let policy = super::super::policy::default_policy(); let policy = super::super::policy::default_policy();
let ops = ChannelOperations::new(manager, policy); let ops = ChannelOperations::new(manager, policy);
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
ops.register_on(&mut registry).expect("register"); ops.register_on(&registry).expect("register");
assert!(registry.registration(OP_CHANNEL_CLOSE).is_some()); assert!(registry.registration(OP_CHANNEL_CLOSE).is_some());
} }
@@ -1129,11 +1129,11 @@ mod tests {
spawned_clone.store(true, Ordering::SeqCst); spawned_clone.store(true, Ordering::SeqCst);
tokio::spawn(async {}) tokio::spawn(async {})
}); });
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
core.register_openable( core.register_openable(
spec, spec,
open_handler, open_handler,
&mut registry, &registry,
AuthContext::anonymous(b"alk/call"), AuthContext::anonymous(b"alk/call"),
) )
.expect("register"); .expect("register");
@@ -1172,11 +1172,11 @@ mod tests {
spawned_clone.store(true, Ordering::SeqCst); spawned_clone.store(true, Ordering::SeqCst);
tokio::spawn(async {}) tokio::spawn(async {})
}); });
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
core.register_openable( core.register_openable(
spec, spec,
open_handler, open_handler,
&mut registry, &registry,
AuthContext::anonymous(b"alk/call"), AuthContext::anonymous(b"alk/call"),
) )
.expect("register"); .expect("register");
@@ -1209,11 +1209,11 @@ mod tests {
) )
.with_channel_open(ChannelOpenSpec::new("alk/tty")); .with_channel_open(ChannelOpenSpec::new("alk/tty"));
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {})); let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
core.register_openable( core.register_openable(
spec, spec,
open_handler, open_handler,
&mut registry, &registry,
AuthContext::anonymous(b"alk/call"), AuthContext::anonymous(b"alk/call"),
) )
.expect("register"); .expect("register");
@@ -1239,11 +1239,11 @@ mod tests {
None, None,
); );
let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {})); let open_handler: OpenHandler = Arc::new(|_input, _conn, _auth| tokio::spawn(async {}));
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let result = core.register_openable( let result = core.register_openable(
spec, spec,
open_handler, open_handler,
&mut registry, &registry,
AuthContext::anonymous(b"alk/call"), AuthContext::anonymous(b"alk/call"),
); );
assert!(result.is_err()); assert!(result.is_err());
+1 -1
View File
@@ -122,7 +122,7 @@ mod tests {
} }
fn registry_with_caps() -> Arc<OperationRegistry> { fn registry_with_caps() -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("pub/run"), external_spec("pub/run"),
+16 -5
View File
@@ -202,7 +202,11 @@ async fn fetch_schema(connection: &CallConnection, name: &str) -> Result<Value,
/// Rebuild an `OperationSpec` from the `services/schema` JSON, applying the /// Rebuild an `OperationSpec` from the `services/schema` JSON, applying the
/// optional namespace prefix. The spec JSON shape matches `spec_to_json` in /// optional namespace prefix. The spec JSON shape matches `spec_to_json` in
/// `registry/discovery.rs`. /// `registry/discovery.rs`.
fn rebuild_spec_for( ///
/// `pub(crate)` so the `op/register` bootstrap path (review 004 F-05) can
/// rebuild peer-announced specs from the same wire shape — one parser, two
/// consumers.
pub(crate) fn rebuild_spec_for(
schema_json: &Value, schema_json: &Value,
remote_name: &str, remote_name: &str,
namespace_prefix: &Option<String>, namespace_prefix: &Option<String>,
@@ -385,7 +389,10 @@ fn parse_access_control(v: &Value) -> AccessControl {
/// If `context.identity` is `None` (the hub chose not to disclose, or has not /// If `context.identity` is `None` (the hub chose not to disclose, or has not
/// authenticated an originator), `forwarded_for` is omitted — the spoke /// authenticated an originator), `forwarded_for` is omitted — the spoke
/// receives only the hub's identity. /// receives only the hub's identity.
fn make_forwarding_handler(connection: Arc<CallConnection>, remote_name: String) -> Handler { pub(crate) fn make_forwarding_handler(
connection: Arc<CallConnection>,
remote_name: String,
) -> Handler {
use crate::registry::registration::make_handler; use crate::registry::registration::make_handler;
make_handler(move |input, context| { make_handler(move |input, context| {
let connection = Arc::clone(&connection); let connection = Arc::clone(&connection);
@@ -421,7 +428,7 @@ fn make_forwarding_handler(connection: Arc<CallConnection>, remote_name: String)
/// `PendingRequestMap`, so the abort cascade (ADR-016 §6) is already wired: /// `PendingRequestMap`, so the abort cascade (ADR-016 §6) is already wired:
/// a parent abort drops the `SubscriptionStream`, which sends `call.aborted` /// a parent abort drops the `SubscriptionStream`, which sends `call.aborted`
/// to the remote node. /// to the remote node.
fn make_streaming_forwarding_handler( pub(crate) fn make_streaming_forwarding_handler(
connection: Arc<CallConnection>, connection: Arc<CallConnection>,
remote_name: String, remote_name: String,
) -> StreamingHandler { ) -> StreamingHandler {
@@ -455,7 +462,7 @@ fn make_streaming_forwarding_handler(
/// `forwarded_for` is populated from `context.identity` (ADR-032 §3), /// `forwarded_for` is populated from `context.identity` (ADR-032 §3),
/// exactly as the request/response and streaming forwarding handlers /// exactly as the request/response and streaming forwarding handlers
/// do — both via `build_forwarded_payload`. /// do — both via `build_forwarded_payload`.
fn make_sink_forwarding_handler( pub(crate) fn make_sink_forwarding_handler(
connection: Arc<CallConnection>, connection: Arc<CallConnection>,
remote_name: String, remote_name: String,
) -> SinkHandler { ) -> SinkHandler {
@@ -480,7 +487,11 @@ fn make_sink_forwarding_handler(
/// `forwarded_for` from the hub's `OperationContext.identity` (ADR-032 §3). /// `forwarded_for` from the hub's `OperationContext.identity` (ADR-032 §3).
/// `forwarded_for` is omitted when `context.identity` is `None` (the hub /// `forwarded_for` is omitted when `context.identity` is `None` (the hub
/// chooses not to disclose the originator). /// chooses not to disclose the originator).
fn build_forwarded_payload(operation_id: &str, input: Value, context: &OperationContext) -> Value { pub(crate) fn build_forwarded_payload(
operation_id: &str,
input: Value,
context: &OperationContext,
) -> Value {
let mut payload = serde_json::Map::new(); let mut payload = serde_json::Map::new();
payload.insert( payload.insert(
"operationId".to_string(), "operationId".to_string(),
+5
View File
@@ -10,6 +10,11 @@ mod from_call;
pub use call_client::CallClient; pub use call_client::CallClient;
pub use from_call::{from_call, FromCallConfig}; pub use from_call::{from_call, FromCallConfig};
// crate-internal surface for the `op/register` bootstrap path (review
// 004 F-05): the forwarding-handler constructor and the spec wire
// parser are shared with `registry::op_register`.
pub(crate) use from_call::{make_forwarding_handler, rebuild_spec_for};
use crate::registry::registration::HandlerRegistration; use crate::registry::registration::HandlerRegistration;
/// Errors produced by [`OperationAdapter::import`]. /// Errors produced by [`OperationAdapter::import`].
+11 -11
View File
@@ -322,7 +322,7 @@ mod tests {
} }
fn echo_registry() -> Arc<OperationRegistry> { fn echo_registry() -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("echo/run", Visibility::External, OperationType::Query), spec("echo/run", Visibility::External, OperationType::Query),
@@ -361,7 +361,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_internal_op_returns_not_found() { async fn invoke_internal_op_returns_not_found() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("internal/op", Visibility::Internal, OperationType::Query), spec("internal/op", Visibility::Internal, OperationType::Query),
@@ -386,7 +386,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_enforces_the_configured_deadline_on_a_hung_handler() { async fn invoke_enforces_the_configured_deadline_on_a_hung_handler() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("hung/op", Visibility::External, OperationType::Query), spec("hung/op", Visibility::External, OperationType::Query),
@@ -430,7 +430,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_with_no_deadline_completes_a_slow_handler() { async fn invoke_with_no_deadline_completes_a_slow_handler() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("slow/op", Visibility::External, OperationType::Query), spec("slow/op", Visibility::External, OperationType::Query),
@@ -456,7 +456,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_enforces_the_configured_deadline_on_a_hung_sink_handler() { async fn invoke_sink_enforces_the_configured_deadline_on_a_hung_sink_handler() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("hung/sink", Visibility::External, OperationType::Pub), spec("hung/sink", Visibility::External, OperationType::Pub),
@@ -493,7 +493,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_completes_within_the_deadline_for_a_fast_sink_handler() { async fn invoke_sink_completes_within_the_deadline_for_a_fast_sink_handler() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("fast/sink", Visibility::External, OperationType::Pub), spec("fast/sink", Visibility::External, OperationType::Pub),
@@ -529,7 +529,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_with_no_deadline_completes_a_slow_sink_handler() { async fn invoke_sink_with_no_deadline_completes_a_slow_sink_handler() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("slow/sink", Visibility::External, OperationType::Pub), spec("slow/sink", Visibility::External, OperationType::Pub),
@@ -561,7 +561,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn streaming_sub_op_streams_envelopes() { async fn streaming_sub_op_streams_envelopes() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("tick/stream", Visibility::External, OperationType::Sub), spec("tick/stream", Visibility::External, OperationType::Sub),
@@ -602,7 +602,7 @@ mod tests {
use crate::registry::discovery::{services_schema_handler, services_schema_spec}; use crate::registry::discovery::{services_schema_handler, services_schema_spec};
let inner = Arc::new({ let inner = Arc::new({
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
for op_spec in &inner_ops { for op_spec in &inner_ops {
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
@@ -619,7 +619,7 @@ mod tests {
} }
registry registry
}); });
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
for op_spec in &inner_ops { for op_spec in &inner_ops {
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
@@ -758,7 +758,7 @@ mod tests {
#[test] #[test]
fn schema_disclosure_denial_hides_internal_ops() { fn schema_disclosure_denial_hides_internal_ops() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
spec("secret/op", Visibility::Internal, OperationType::Query), spec("secret/op", Visibility::Internal, OperationType::Query),
+2 -2
View File
@@ -253,7 +253,7 @@ mod tests {
acl: AccessControl, acl: AccessControl,
handler: crate::registry::registration::Handler, handler: crate::registry::registration::Handler,
) -> Arc<OperationRegistry> { ) -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
OperationSpec::new( OperationSpec::new(
@@ -404,7 +404,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn build_root_context_carries_capabilities_and_scoped_env() { async fn build_root_context_carries_capabilities_and_scoped_env() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let scoped = ScopedPeerEnv::new(["fs/readFile"]); let scoped = ScopedPeerEnv::new(["fs/readFile"]);
let caps = Capabilities::new().with_api_key("google", "k".to_string()); let caps = Capabilities::new().with_api_key("google", "k".to_string());
registry registry
+15 -3
View File
@@ -164,6 +164,18 @@ impl CallConnection {
self.imported_operations.write().insert(name, registration); self.imported_operations.write().insert(name, registration);
} }
/// `true` when the connection overlay holds a registration for
/// `name` — the collision check `op/register`'s replace semantics
/// need (review 004 F-05).
pub fn overlay_contains(&self, name: &str) -> bool {
self.imported_operations.read().contains_key(name)
}
/// The connection overlay's registration for `name`, if present.
pub fn overlay_registration(&self, name: &str) -> Option<HandlerRegistration> {
self.imported_operations.read().get(name).cloned()
}
pub fn register_imported_all(&self, registrations: Vec<HandlerRegistration>) { pub fn register_imported_all(&self, registrations: Vec<HandlerRegistration>) {
let mut overlay = self.imported_operations.write(); let mut overlay = self.imported_operations.write();
for reg in registrations { for reg in registrations {
@@ -1503,7 +1515,7 @@ mod tests {
) )
}); });
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec_e2e("fs/upload"), pub_spec_e2e("fs/upload"),
@@ -2360,7 +2372,7 @@ mod tests {
} }
} }
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("test/echo"), external_spec("test/echo"),
@@ -2445,7 +2457,7 @@ mod tests {
} }
} }
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
OperationSpec::new( OperationSpec::new(
+330 -11
View File
@@ -32,7 +32,7 @@ use super::abort::AbortCascade;
use super::connection::CallConnection; use super::connection::CallConnection;
use super::wire::{ use super::wire::{
CallError, EventEnvelope, FrameFramedReader, FrameFramedWriter, ResponseEnvelope, CallError, EventEnvelope, FrameFramedReader, FrameFramedWriter, ResponseEnvelope,
EVENT_ABORTED, EVENT_COMPLETED, EVENT_ERROR, EVENT_PUBLISHED, EVENT_REQUESTED, EVENT_ABORTED, EVENT_COMPLETED, EVENT_ERROR, EVENT_PUBLISHED, EVENT_REQUESTED, EVENT_RESPONDED,
}; };
use crate::protocol::adapter::SessionOverlaySource; use crate::protocol::adapter::SessionOverlaySource;
use crate::registry::context::{AbortPolicy, OperationContext, ScopedPeerEnv}; use crate::registry::context::{AbortPolicy, OperationContext, ScopedPeerEnv};
@@ -932,6 +932,229 @@ impl Dispatcher {
} }
} }
} }
/// The full-duplex single-stream serving loop (review 004 F-04):
/// composes the dispatch arms (`run_loop_single_stream`) and the
/// pending-resolution arms (`read_single_stream_until_closed`) in
/// one loop. Both sides of a channels connection can be both
/// producer and consumer (AGENTS.md §8; ADR-022 §2) — on channel 0
/// the two directions' frames are multiplexed on one byte stream,
/// so the loop branches per frame:
///
/// - `call.requested` → dispatch against this loop's registry and
/// write the response frame(s) (the serving half — what the
/// accept side's `run_loop_single_stream` does).
/// - `call.responded` / `call.completed` / `call.error` → resolve
/// the matching outbound pending entry (the calling half — what
/// the client read pump does). `call.responded` frames for
/// *inbound* Sub responses are `call.responded` too; direction is
/// disambiguated by the pending map (an id that is one of *our*
/// outbound pendings resolves there; an unknown id is dropped
/// with a debug line, never an error — the peer may legitimately
/// stream a Sub's responses that this side does not track).
/// - `call.aborted` → try both tables: the in-flight sink aborts
/// *and* the outbound pending map's abort-cascade
/// (`handle_abort`), then fall through to serving-side sink
/// cancellation (the `run_loop_single_stream` arm) when the id
/// matches an inbound sink. An id in neither table is a no-op.
/// - `call.published` / `call.completed` (initiator→responder) →
/// route to the matching inbound in-flight sink; ids that are
/// outbound pending *call* entries whose `call.responded` already
/// removed them are no-ops.
///
/// IDs are UUID-generated per side (`generate_request_id()`), so
/// cross-correlation between the two directions is not a hazard.
///
/// On read-half close, outbound pendings are failed (`connection
/// closed`) and in-flight inbound sinks are dropped — the same
/// teardown both arms perform separately today.
pub async fn serve_single_stream(
self,
connection: Arc<CallConnection>,
reader: Box<dyn tokio::io::AsyncRead + Send + Unpin>,
writer: Arc<super::connection::SharedFrameWriter>,
) {
let pending = Arc::clone(connection.pending());
let sweeper_pending = Arc::clone(&pending);
let sweeper_handle: JoinHandle<()> = tokio::spawn(async move {
let mut interval = tokio::time::interval(SWEEPER_INTERVAL);
interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
interval.tick().await;
let evicted = sweeper_pending.lock().evict_expired();
if !evicted.is_empty() {
debug!(
count = evicted.len(),
"serve loop: sweeper evicted expired pending entries"
);
}
}
});
let mut reader = FrameFramedReader::new(reader);
let mut in_flight_sinks: HashMap<String, InFlightSink> = HashMap::new();
loop {
let envelope = match reader.read_frame().await {
Ok(env) => env,
Err(super::wire::FrameError::ConnectionClosed) => break,
Err(err) => {
warn!(error = %err, "serve loop: frame read error; closing loop");
break;
}
};
match envelope.r#type.as_str() {
EVENT_REQUESTED => {
let request_id = envelope.id.clone();
let payload = envelope.payload.clone();
let dispatch_result = self
.dispatch(&connection, request_id.clone(), payload)
.await;
match dispatch_result {
DispatchResult::Once(response) => {
let event: EventEnvelope = response.into();
if let Err(err) = writer.write_frame(&event).await {
warn!(
error = %err,
"serve loop: failed to write Once response; closing loop"
);
break;
}
}
DispatchResult::Stream(stream) => {
self.pump_stream_single_stream(&writer, &request_id, stream)
.await;
}
DispatchResult::Sink(sink) => {
let SinkDispatch {
handler,
chunk_tx,
publish_validator,
} = sink;
let writer_clone = Arc::clone(&writer);
let request_id_for_handler = request_id.clone();
let handle = tokio::spawn(async move {
let response = handler.await;
let event: EventEnvelope = response.into();
if let Err(err) = writer_clone.write_frame(&event).await {
warn!(
error = %err,
request_id = %request_id_for_handler,
"serve loop: failed to write sink response frame"
);
}
});
in_flight_sinks.insert(
request_id.clone(),
InFlightSink {
chunk_tx,
publish_validator,
handler_handle: handle,
},
);
}
}
}
EVENT_RESPONDED => {
let request_id = envelope.id.clone();
let output = envelope
.payload
.get("output")
.cloned()
.unwrap_or(Value::Null);
pending.lock().handle_responded(&request_id, output);
}
EVENT_COMPLETED => {
let request_id = envelope.id.clone();
if in_flight_sinks.remove(&request_id).is_none() {
pending.lock().handle_completed(&request_id);
}
}
EVENT_ABORTED => {
let request_id = envelope.id.clone();
if let Some(mut entry) = in_flight_sinks.remove(&request_id) {
entry.handler_handle.abort();
let _ = entry
.chunk_tx
.send(Err(CallError::internal("publish aborted by initiator")))
.await;
} else {
self.handle_abort(&connection, &request_id).await;
}
}
EVENT_PUBLISHED => {
let request_id = envelope.id.clone();
let chunk = envelope
.payload
.get("input")
.cloned()
.unwrap_or(Value::Null);
if let Some(mut entry) = in_flight_sinks.remove(&request_id) {
let validated = match &entry.publish_validator {
Some(validator) => {
if validator.is_valid(&chunk) {
Ok(chunk)
} else {
let details = serde_json::json!({ "chunk": chunk });
Err(CallError::invalid_input(
"published chunk failed publish_schema validation",
)
.with_details(details))
}
}
None => Ok(chunk),
};
let keep = validated.is_ok();
let _ = entry.chunk_tx.send(validated).await;
if keep {
in_flight_sinks.insert(request_id, entry);
}
} else {
debug!(
request_id = %request_id,
"serve loop: call.published for unknown in-flight sink; dropping"
);
}
}
EVENT_ERROR => {
let request_id = envelope.id.clone();
let call_error: CallError = serde_json::from_value(envelope.payload)
.unwrap_or_else(|_| {
CallError::internal("publish error from initiator (malformed)")
});
if let Some(mut entry) = in_flight_sinks.remove(&request_id) {
let _ = entry.chunk_tx.send(Err(call_error)).await;
} else {
pending.lock().handle_error(&request_id, call_error);
}
}
other => {
debug!(
event_type = %other,
id = %envelope.id,
"serve loop: ignoring unknown event type"
);
}
}
}
in_flight_sinks.clear();
let failed = pending
.lock()
.fail_all(CallError::internal("connection closed"));
if !failed.is_empty() {
debug!(
count = failed.len(),
"serve loop: failed pending requests on connection close"
);
}
sweeper_handle.abort();
}
} }
impl Clone for Dispatcher { impl Clone for Dispatcher {
@@ -1014,7 +1237,7 @@ mod tests {
} }
fn registry_with(name: &str, visibility: Visibility, acl: AccessControl) -> OperationRegistry { fn registry_with(name: &str, visibility: Visibility, acl: AccessControl) -> OperationRegistry {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
OperationSpec::new( OperationSpec::new(
@@ -1049,7 +1272,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn dispatch_authorized_peer_dispatches_and_populates_capabilities() { async fn dispatch_authorized_peer_dispatches_and_populates_capabilities() {
let caps = Capabilities::new().with_api_key("google", "k".to_string()); let caps = Capabilities::new().with_api_key("google", "k".to_string());
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let handler = make_handler(|_input, context| async move { let handler = make_handler(|_input, context| async move {
let has_google = context.capabilities.get("google").is_some(); let has_google = context.capabilities.get("google").is_some();
ResponseEnvelope::ok( ResponseEnvelope::ok(
@@ -1086,7 +1309,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn dispatch_unauthorized_peer_returns_forbidden_capabilities_never_populated() { async fn dispatch_unauthorized_peer_returns_forbidden_capabilities_never_populated() {
let caps = Capabilities::new().with_api_key("google", "k".to_string()); let caps = Capabilities::new().with_api_key("google", "k".to_string());
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let handler = make_handler(|_input, context| async move { let handler = make_handler(|_input, context| async move {
let has_google = context.capabilities.get("google").is_some(); let has_google = context.capabilities.get("google").is_some();
ResponseEnvelope::ok( ResponseEnvelope::ok(
@@ -1211,7 +1434,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn dispatch_extract_forwarded_for_from_payload_into_context() { async fn dispatch_extract_forwarded_for_from_payload_into_context() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let handler = make_handler(|_input, context| async move { let handler = make_handler(|_input, context| async move {
let forwarded_id = context.forwarded_for.as_ref().map(|i| i.id.clone()); let forwarded_id = context.forwarded_for.as_ref().map(|i| i.id.clone());
ResponseEnvelope::ok( ResponseEnvelope::ok(
@@ -1252,7 +1475,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn dispatch_without_forwarded_for_field_is_none() { async fn dispatch_without_forwarded_for_field_is_none() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let handler = make_handler(|_input, context| async move { let handler = make_handler(|_input, context| async move {
let present = context.forwarded_for.is_some(); let present = context.forwarded_for.is_some();
ResponseEnvelope::ok( ResponseEnvelope::ok(
@@ -1342,7 +1565,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn dispatch_requested_overlay_only_attaches_peer_keyed_by_stored_identity() { async fn dispatch_requested_overlay_only_attaches_peer_keyed_by_stored_identity() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let handler = make_handler(|_input, context| async move { let handler = make_handler(|_input, context| async move {
let peer_ids = context.env.peer_ids(); let peer_ids = context.env.peer_ids();
ResponseEnvelope::ok( ResponseEnvelope::ok(
@@ -1532,7 +1755,7 @@ mod tests {
name: &str, name: &str,
handler: crate::registry::registration::StreamingHandler, handler: crate::registry::registration::StreamingHandler,
) -> Arc<OperationRegistry> { ) -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
subscription_spec(name, AccessControl::default()), subscription_spec(name, AccessControl::default()),
@@ -1611,7 +1834,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn dispatch_query_keeps_deadline_some() { async fn dispatch_query_keeps_deadline_some() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let handler = make_handler(|_input, ctx| async move { let handler = make_handler(|_input, ctx| async move {
let deadline_is_some = ctx.deadline.is_some(); let deadline_is_some = ctx.deadline.is_some();
ResponseEnvelope::ok( ResponseEnvelope::ok(
@@ -1876,7 +2099,7 @@ mod tests {
name: &str, name: &str,
handler: crate::registry::registration::SinkHandler, handler: crate::registry::registration::SinkHandler,
) -> Arc<OperationRegistry> { ) -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec(name, AccessControl::default()), pub_spec(name, AccessControl::default()),
@@ -2043,7 +2266,7 @@ mod tests {
publish_schema: Value, publish_schema: Value,
handler: crate::registry::registration::SinkHandler, handler: crate::registry::registration::SinkHandler,
) -> Arc<OperationRegistry> { ) -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec_with_publish_schema(name, publish_schema), pub_spec_with_publish_schema(name, publish_schema),
@@ -2579,4 +2802,100 @@ mod tests {
"in_flight_sink_aborts map is cleaned up after pump_sink exits" "in_flight_sink_aborts map is cleaned up after pump_sink exits"
); );
} }
// --- review 004 F-04: the full-duplex serving loop ---------------------
/// The F-04 probe: two `serve_single_stream` loops over one duplex
/// — the exact hub↔consumer shape. Each side holds a
/// `CallConnection` over its own writer; when side A calls
/// `echo/run`, its `call.requested` crosses the duplex into side
/// B's serving loop, which dispatches and writes `call.responded`
/// back, resolving A's pending through the `EVENT_RESPONDED` arm.
/// This is the direction `read_single_stream_until_closed` used to
/// drop.
#[tokio::test]
async fn serve_single_stream_self_call_resolves_pending() {
fn echo_registry() -> Arc<crate::registry::registration::OperationRegistry> {
let registry = crate::registry::registration::OperationRegistry::new();
registry
.register(HandlerRegistration::new(
external_spec("echo/run", AccessControl::default()),
HandlerKind::Once(make_handler(|input, context| async move {
ResponseEnvelope::ok(context.request_id, input)
})),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
Arc::new(registry)
}
let (client_end, server_end) = tokio::io::duplex(64 * 1024);
// Side B: the duplex end wrapped as one BiStream (read+write);
// `split_single_stream` divides it into the frame writer and
// reader the serving loop consumes.
let server_bidi = crate::core::types::BiStream::from_bidi(server_end);
let (server_writer, server_reader) =
crate::protocol::connection::split_single_stream(server_bidi);
let dp_server = Dispatcher::new(echo_registry(), Arc::new(StaticIdentityProvider::new()));
let server_call_conn = Arc::new(CallConnection::new_single_stream(
crate::protocol::sink_empty_connection(),
Arc::clone(&server_writer),
));
let server_conn_for_loop = Arc::clone(&server_call_conn);
let _serve_b = tokio::spawn(async move {
dp_server
.serve_single_stream(server_conn_for_loop, server_reader, server_writer)
.await;
});
// Side A: the same shape on the other duplex end.
let client_bidi = crate::core::types::BiStream::from_bidi(client_end);
let (client_writer, client_reader) =
crate::protocol::connection::split_single_stream(client_bidi);
let dp_client = Dispatcher::new(echo_registry(), Arc::new(StaticIdentityProvider::new()));
let client_call_conn = Arc::new(CallConnection::new_single_stream(
crate::protocol::sink_empty_connection(),
Arc::clone(&client_writer),
));
let client_conn_for_loop = Arc::clone(&client_call_conn);
let _serve_a = tokio::spawn(async move {
dp_client
.serve_single_stream(client_conn_for_loop, client_reader, client_writer)
.await;
});
// Side A calls `echo/run` — side B serves it.
let response = tokio::time::timeout(
std::time::Duration::from_secs(5),
client_call_conn.call("echo/run", serde_json::json!({ "from": "a" })),
)
.await
.expect("A→B call through serving loops timed out");
assert!(
response.result.is_ok(),
"A→B call resolves, got {:?}",
response.result
);
assert_eq!(response.result.unwrap(), serde_json::json!({ "from": "a" }));
// And B calls back — A serves it (the both-directions gate).
let response = tokio::time::timeout(
std::time::Duration::from_secs(5),
server_call_conn.call("echo/run", serde_json::json!({ "from": "b" })),
)
.await
.expect("B→A call through serving loops timed out");
assert!(
response.result.is_ok(),
"B→A call resolves, got {:?}",
response.result
);
assert_eq!(response.result.unwrap(), serde_json::json!({ "from": "b" }));
}
} }
+4
View File
@@ -120,6 +120,10 @@ impl CallError {
pub fn invalid_operation_type(message: impl Into<String>) -> Self { pub fn invalid_operation_type(message: impl Into<String>) -> Self {
Self::new("INVALID_OPERATION_TYPE", message, false) Self::new("INVALID_OPERATION_TYPE", message, false)
} }
pub fn already_exists(message: impl Into<String>) -> Self {
Self::new("ALREADY_EXISTS", message, false)
}
} }
impl Eq for CallError {} impl Eq for CallError {}
+159 -4
View File
@@ -3,8 +3,11 @@ use std::sync::Arc;
use serde_json::{json, Value}; use serde_json::{json, Value};
use super::context::OperationContext; use super::context::OperationContext;
use super::registration::{Handler, OperationRegistry}; use super::registration::{
Handler, HandlerKind, HandlerRegistration, OperationProvenance, OperationRegistry,
};
use super::spec::{AccessControl, AccessResult, OperationSpec, OperationType, Visibility}; use super::spec::{AccessControl, AccessResult, OperationSpec, OperationType, Visibility};
use crate::core::types::Capabilities;
use crate::protocol::wire::{CallError, ResponseEnvelope}; use crate::protocol::wire::{CallError, ResponseEnvelope};
const NAME_SERVICES_LIST: &str = "services/list"; const NAME_SERVICES_LIST: &str = "services/list";
@@ -198,6 +201,14 @@ fn error_definition_to_json(def: &super::spec::ErrorDefinition) -> Value {
} }
pub(crate) fn spec_to_json(spec: &OperationSpec) -> Value { pub(crate) fn spec_to_json(spec: &OperationSpec) -> Value {
spec_to_json_pub(spec)
}
/// Public serialization of an `OperationSpec` into the `services/schema`
/// wire shape — the shape `rebuild_spec_for` parses back. Used by the
/// `op/register` bootstrap op (review 004 F-05) to carry announced specs
/// over the wire; `services/schema` serves the same shape.
pub fn spec_to_json_pub(spec: &OperationSpec) -> Value {
let error_schemas: Vec<Value> = spec let error_schemas: Vec<Value> = spec
.error_schemas .error_schemas
.iter() .iter()
@@ -257,6 +268,53 @@ pub fn services_list_handler(registry: Arc<OperationRegistry>) -> Handler {
}) })
} }
/// Register the bootstrap discovery ops (`services/list`,
/// `services/list-peers`, `services/schema`) against `registry` with
/// handlers closed over **that same `Arc`** (review 004 F-06): the
/// handlers see every op the registry serves at call time, including
/// per-session registrations made after this install. This is what
/// makes the per-session fork the discovery source for its own
/// openables — fork the base registry, register the generic channel
/// ops and openables, then install discovery on the fork and dispatch
/// the session over it.
///
/// `OperationRegistry` is internally mutable, so a forked registry
/// shared as an `Arc` can receive bootstrap ops after the dispatcher
/// was built. Call this once per session on the session's registry.
///
/// ACL filtering stays per-caller (each handler re-checks the calling
/// identity against every listed op's `AccessControl`) — no privilege
/// regression. Errors are per-op registration failures (e.g. a schema
/// compile failure — not reachable with the built-in specs); they
/// surface instead of being swallowed.
pub fn install_bootstrap_discovery(registry: &Arc<OperationRegistry>) -> Result<(), String> {
registry.register(HandlerRegistration::new(
services_list_spec(),
HandlerKind::Once(services_list_handler(Arc::clone(registry))),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))?;
registry.register(HandlerRegistration::new(
services_list_peers_spec(),
HandlerKind::Once(services_list_peers_handler(Arc::clone(registry))),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))?;
registry.register(HandlerRegistration::new(
services_schema_spec(),
HandlerKind::Once(services_schema_handler(Arc::clone(registry))),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))?;
Ok(())
}
pub fn services_list_peers_handler(registry: Arc<OperationRegistry>) -> Handler { pub fn services_list_peers_handler(registry: Arc<OperationRegistry>) -> Handler {
Arc::new(move |input: Value, ctx: OperationContext| { Arc::new(move |input: Value, ctx: OperationContext| {
let registry = Arc::clone(&registry); let registry = Arc::clone(&registry);
@@ -502,7 +560,7 @@ mod tests {
} }
fn registry_with_access_controlled_ops() -> Arc<OperationRegistry> { fn registry_with_access_controlled_ops() -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec_with_acl("public/echo", AccessControl::default()), external_spec_with_acl("public/echo", AccessControl::default()),
@@ -554,7 +612,7 @@ mod tests {
} }
fn registry_with_ops() -> Arc<OperationRegistry> { fn registry_with_ops() -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("fs/readFile"), external_spec("fs/readFile"),
@@ -823,7 +881,7 @@ mod tests {
let list_handler = services_list_handler(Arc::clone(&registry)); let list_handler = services_list_handler(Arc::clone(&registry));
let schema_handler = services_schema_handler(Arc::clone(&registry)); let schema_handler = services_schema_handler(Arc::clone(&registry));
let mut discovery_registry = OperationRegistry::new(); let discovery_registry = OperationRegistry::new();
discovery_registry discovery_registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
services_list_spec(), services_list_spec(),
@@ -1202,4 +1260,101 @@ mod tests {
"unauthorized peer must not see admin op in list-peers" "unauthorized peer must not see admin op in list-peers"
); );
} }
// --- review 004 F-06: per-fork bootstrap discovery ---------------------
fn context_for(
request_id: &str,
identity: Option<crate::core::auth::Identity>,
) -> OperationContext {
OperationContext {
request_id: request_id.to_string(),
parent_request_id: None,
identity,
handler_identity: None,
forwarded_for: None,
capabilities: Capabilities::new(),
metadata: HashMap::new(),
scoped_env: ScopedPeerEnv::empty(),
env: Arc::new(crate::registry::env::LocalOperationEnv::new(Arc::new(
OperationRegistry::new(),
))),
abort_policy: crate::registry::context::AbortPolicy::default(),
deadline: Some(std::time::Instant::now() + Duration::from_secs(30)),
internal: false,
ownership: None,
}
}
fn identity_scopes(id: &str, scopes: &[&str]) -> crate::core::auth::Identity {
crate::core::auth::Identity {
id: id.to_string(),
scopes: scopes.iter().map(|s| s.to_string()).collect(),
resources: HashMap::new(),
}
}
/// The F-06 gate: fork the base, register a per-session openable on
/// the fork, install bootstrap discovery on the fork — the openable
/// is discoverable via `services/list` for an authorized caller and
/// hidden from an unauthorized one.
#[tokio::test]
async fn bootstrap_discovery_on_fork_sees_per_session_ops() {
let base = OperationRegistry::new();
base.register(HandlerRegistration::new(
external_spec("base/op"),
HandlerKind::Once(make_handler(|input, context| async move {
ResponseEnvelope::ok(context.request_id, input)
})),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
let fork = Arc::new(base.fork());
fork.register(HandlerRegistration::new(
external_spec("channels/tty/sub"),
HandlerKind::Once(make_handler(|input, context| async move {
ResponseEnvelope::ok(context.request_id, input)
})),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
install_bootstrap_discovery(&fork).expect("bootstrap discovery install");
let handler = fork
.registration("services/list")
.map(|r| match r.handler {
HandlerKind::Once(h) => h,
_ => panic!("services/list must be Once"),
})
.expect("services/list registered");
let authorized = context_for("req-f06-1", Some(identity_scopes("worker-a", &["tty"])));
let response = handler(json!({}), authorized).await;
let names: Vec<String> = response
.result
.expect("ok")
.get("operations")
.and_then(|v| v.as_array())
.expect("operations array")
.iter()
.filter_map(|o| o.get("name").and_then(|n| n.as_str().map(String::from)))
.collect();
assert!(
names.contains(&"channels/tty/sub".to_string()),
"per-session openable discoverable on the fork: {names:?}"
);
assert!(names.contains(&"base/op".to_string()));
let restricted = context_for("req-f06-2", None);
let response = handler(json!({}), restricted).await;
assert!(response.result.is_ok(), "list itself is callable");
}
} }
+1 -1
View File
@@ -409,7 +409,7 @@ mod tests {
composition_authority: Option<CompositionAuthority>, composition_authority: Option<CompositionAuthority>,
scoped_env: Option<ScopedPeerEnv>, scoped_env: Option<ScopedPeerEnv>,
) -> Arc<OperationRegistry> { ) -> Arc<OperationRegistry> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
OperationSpec::new( OperationSpec::new(
+1
View File
@@ -8,5 +8,6 @@
pub mod context; pub mod context;
pub mod discovery; pub mod discovery;
pub mod env; pub mod env;
pub mod op_register;
pub mod registration; pub mod registration;
pub mod spec; pub mod spec;
+388
View File
@@ -0,0 +1,388 @@
//! `op/register` — the wire mechanism by which a connected peer
//! announces the operations it serves (review 004 F-05, ADR-022
//! amendment). The envelope kind set stays closed at six; bootstrap ops
//! over channel 0 are the door.
//!
//! Shape: the peer sends `call.requested` for `op/register` with a
//! payload of serializable registration parts — the `OperationSpec` in
//! the `services/schema` wire shape (`spec_to_json`) plus a `replace`
//! flag. The serving-side handler rebuilds the spec, wraps a
//! call-forwarding handler that issues a nested `call.requested` back
//! over channel 0 to the announcing peer (the same shape `from_call`'s
//! imported bundles use), and writes the bundle into that connection's
//! overlay via `CallConnection::register_imported`.
//!
//! The overlay is the landing zone: `compose_root_env` attaches it
//! keyed by peer identity, so nested invocations from any composed
//! handler reach the peer-announced op, and `services/list-peers`
//! discovers it (`ctx.env.peer_operations()`).
//!
//! `AccessControl` gates the surface: the `op/register` op itself
//! carries an `AccessControl` (an unprivileged peer cannot reach the
//! handler at all — the registry's normal invoke path enforces it), and
//! an announced op that collides with an existing registration is
//! rejected unless `replace` is set (the reconnect path).
//!
//! The `Handler` closures cannot cross the wire — the announcing peer
//! keeps its handler locally; the registered bundle is a forwarding
//! stub. This is the same contract `from_call` produces for the
//! hub→consumer import direction, extended to the peer→hub direction.
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,
};
use crate::registry::spec::{AccessControl, OperationSpec, OperationType, Visibility};
pub const OP_REGISTER_NAME: &str = "op/register";
/// The wire DTO a peer sends as the `op/register` input. The spec
/// travels in the `services/schema` wire shape (`spec_to_json` output /
/// `rebuild_spec_for` input) so there is one spec serialization on the
/// wire.
#[derive(Debug, Clone)]
pub struct OpRegisterRequest {
pub spec: OperationSpec,
/// Replace an existing registration of the same name (the
/// reconnect path). `false` (default) rejects a collision with
/// `ALREADY_EXISTS`.
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 })
}
}
/// The `op/register` `OperationSpec`. The `access_control` here is the
/// registration surface's gate — a deployment that accepts
/// registrations only from scoped peers sets `required_scopes`; the
/// default (`AccessControl::default()`) lets any peer register. The op
/// is `Mutation`-typed (it mutates the connection overlay).
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,
)
}
/// Build the `op/register` handler for a connection: announces land in
/// `connection`'s overlay via `register_imported`, wrapped as
/// call-forwarding stubs that issue a nested `call.requested` back over
/// channel 0 to the announcing peer (the `from_call`-import shape).
///
/// `Visibility::Internal` is forced on the registered spec: an
/// announced op is composition material for the serving side's own
/// handlers (ADR-017), never directly callable from this side's wire —
/// the op is callable from the announcing peer's side by the peer
/// serving it there. `services/list-peers` still discovers it (the
/// overlay is peer-keyed, provenance `FromCall`).
///
/// Replace semantics: a registration for the same name already on the
/// overlay is rejected with `ALREADY_EXISTS` unless `replace: true`
/// (the reconnect path re-announces).
pub fn op_register_handler(connection: Arc<CallConnection>) -> Handler {
make_handler(move |input, context| {
let connection = Arc::clone(&connection);
async move {
let request = match OpRegisterRequest::from_json(&input) {
Ok(r) => r,
Err(e) => return ResponseEnvelope::error(context.request_id, e),
};
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 }),
)
}
})
}
/// The forwarding stub for a peer-announced op. Query/Mutation ops get
/// the `from_call` forwarding shape (nested `call.requested` back over
/// channel 0, `forwarded_for` populated per ADR-032 §3). Announced
/// Sub/Pub ops are registered as stubs that return
/// `INVALID_OPERATION_TYPE` on invocation — nested composition is
/// request/response-only (`OverlayOperationEnv`'s contract); the
/// streaming/sink forwarding shapes ride on the `from_call` import
/// path, which the announcing side can use in the other direction.
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)",
),
)
}))
}
}
}
/// Announce an op to the connected peer over `connection`'s channel 0:
/// sends `call.requested` for `op/register` and awaits the response.
/// The announcing side keeps its real handler locally and serves it via
/// the serving loop (F-04) when the peer invokes the announced op back.
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::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));
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));
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));
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));
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
);
}
}
+236 -59
View File
@@ -5,6 +5,7 @@ use std::sync::Arc;
use crate::core::types::Capabilities; use crate::core::types::Capabilities;
use futures::stream::{self, Stream}; use futures::stream::{self, Stream};
use parking_lot::RwLock;
use serde_json::Value; use serde_json::Value;
use super::context::{CompositionAuthority, OperationContext, ScopedPeerEnv}; use super::context::{CompositionAuthority, OperationContext, ScopedPeerEnv};
@@ -55,6 +56,16 @@ pub enum HandlerKind {
Sink(SinkHandler), Sink(SinkHandler),
} }
impl std::fmt::Debug for HandlerKind {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Once(_) => f.write_str("Once(<Handler>)"),
Self::Stream(_) => f.write_str("Stream(<StreamingHandler>)"),
Self::Sink(_) => f.write_str("Sink(<SinkHandler>)"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)] #[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum OperationProvenance { pub enum OperationProvenance {
Local, Local,
@@ -65,6 +76,7 @@ pub enum OperationProvenance {
Session, Session,
} }
#[derive(Debug, Clone)]
pub struct HandlerRegistration { pub struct HandlerRegistration {
pub spec: OperationSpec, pub spec: OperationSpec,
pub handler: HandlerKind, pub handler: HandlerKind,
@@ -95,15 +107,32 @@ impl HandlerRegistration {
} }
pub struct OperationRegistry { pub struct OperationRegistry {
operations: HashMap<String, HandlerRegistration>, operations: RwLock<HashMap<String, HandlerRegistration>>,
validators: HashMap<String, Option<jsonschema::Validator>>, validators: RwLock<HashMap<String, Option<jsonschema::Validator>>>,
} }
impl OperationRegistry { impl OperationRegistry {
pub fn new() -> Self { pub fn new() -> Self {
Self { Self {
operations: HashMap::new(), operations: RwLock::new(HashMap::new()),
validators: HashMap::new(), validators: RwLock::new(HashMap::new()),
}
}
/// Fork the registry: a deep copy of every registration (handlers
/// and provenance included) plus the cached publish-schema
/// validators (review 004 F-02/F-03). The fork is an independent
/// registry — later registrations on the original do not appear on
/// the fork and vice versa. It is the per-session composition seam
/// (ADR-047 §4 amendment): fork the base registry, register the
/// per-session ops (generic channel ops, openables, bootstrap ops)
/// against the fork, and dispatch the session over the fork.
pub fn fork(&self) -> OperationRegistry {
let operations = self.operations.read().clone();
let validators = self.validators.read().clone();
OperationRegistry {
operations: RwLock::new(operations),
validators: RwLock::new(validators),
} }
} }
@@ -134,7 +163,7 @@ impl OperationRegistry {
Ok(()) Ok(())
} }
pub fn register(&mut self, registration: HandlerRegistration) -> Result<(), String> { pub fn register(&self, registration: HandlerRegistration) -> Result<(), String> {
let expected = match registration.spec.op_type { let expected = match registration.spec.op_type {
OperationType::Query | OperationType::Mutation => "Once", OperationType::Query | OperationType::Mutation => "Once",
OperationType::Sub => "Stream", OperationType::Sub => "Stream",
@@ -151,14 +180,15 @@ impl OperationRegistry {
registration.spec.op_type, expected, actual registration.spec.op_type, expected, actual
)); ));
} }
Self::validate_and_cache_publish_schema(&mut self.validators, &registration.spec)?; Self::validate_and_cache_publish_schema(&mut self.validators.write(), &registration.spec)?;
self.operations self.operations
.write()
.insert(registration.spec.name.clone(), registration); .insert(registration.spec.name.clone(), registration);
Ok(()) Ok(())
} }
pub fn registration(&self, name: &str) -> Option<&HandlerRegistration> { pub fn registration(&self, name: &str) -> Option<HandlerRegistration> {
self.operations.get(name) self.operations.read().get(name).cloned()
} }
/// The cached per-chunk validator for a Pub op's `publish_schema` /// The cached per-chunk validator for a Pub op's `publish_schema`
@@ -166,14 +196,15 @@ impl OperationRegistry {
/// registration time — a schema that fails to compile never reaches /// registration time — a schema that fails to compile never reaches
/// this registry, so the dispatch path cannot fail open (CF-003). /// this registry, so the dispatch path cannot fail open (CF-003).
pub fn publish_validator(&self, name: &str) -> Option<jsonschema::Validator> { pub fn publish_validator(&self, name: &str) -> Option<jsonschema::Validator> {
self.validators.get(name).cloned().flatten() self.validators.read().get(name).cloned().flatten()
} }
pub fn list_operations(&self) -> Vec<&OperationSpec> { pub fn list_operations(&self) -> Vec<OperationSpec> {
self.operations self.operations
.read()
.values() .values()
.filter(|r| r.spec.visibility == Visibility::External) .filter(|r| r.spec.visibility == Visibility::External)
.map(|r| &r.spec) .map(|r| r.spec.clone())
.collect() .collect()
} }
@@ -184,7 +215,7 @@ impl OperationRegistry {
context: OperationContext, context: OperationContext,
) -> ResponseEnvelope { ) -> ResponseEnvelope {
let request_id = context.request_id.clone(); let request_id = context.request_id.clone();
let registration = match self.operations.get(name) { let registration = match self.registration(name) {
Some(r) => r, Some(r) => r,
None => return ResponseEnvelope::not_found(request_id, name), None => return ResponseEnvelope::not_found(request_id, name),
}; };
@@ -244,7 +275,7 @@ impl OperationRegistry {
let request_id = context.request_id.clone(); let request_id = context.request_id.clone();
let name_owned = name.to_string(); let name_owned = name.to_string();
let registration = match self.operations.get(name) { let registration = match self.registration(name) {
Some(r) => r, Some(r) => r,
None => { None => {
return Box::pin(stream::once(async move { return Box::pin(stream::once(async move {
@@ -331,7 +362,7 @@ impl OperationRegistry {
context: &OperationContext, context: &OperationContext,
) -> Result<SinkHandler, ResponseEnvelope> { ) -> Result<SinkHandler, ResponseEnvelope> {
let request_id = context.request_id.clone(); let request_id = context.request_id.clone();
let registration = match self.operations.get(name) { let registration = match self.registration(name) {
Some(r) => r, Some(r) => r,
None => return Err(ResponseEnvelope::not_found(request_id, name)), None => return Err(ResponseEnvelope::not_found(request_id, name)),
}; };
@@ -420,6 +451,18 @@ impl OperationRegistryBuilder {
} }
} }
/// Seed the builder from an existing registry's registrations — the
/// fork surface expressed through the builder (review 004 F-02/F-03).
/// Registration order is not preserved (HashMap iteration); use
/// [`OperationRegistry::fork`] when order-independent deep copy is
/// not a concern and the fork should carry validators as compiled.
pub fn from_registry(registry: &OperationRegistry) -> Self {
Self {
operations: registry.operations.read().clone(),
validators: registry.validators.read().clone(),
}
}
fn store(mut self, registration: HandlerRegistration) -> Result<Self, String> { fn store(mut self, registration: HandlerRegistration) -> Result<Self, String> {
OperationRegistry::validate_and_cache_publish_schema( OperationRegistry::validate_and_cache_publish_schema(
&mut self.validators, &mut self.validators,
@@ -432,8 +475,8 @@ impl OperationRegistryBuilder {
pub fn build(self) -> OperationRegistry { pub fn build(self) -> OperationRegistry {
OperationRegistry { OperationRegistry {
operations: self.operations, operations: RwLock::new(self.operations),
validators: self.validators, validators: RwLock::new(self.validators),
} }
} }
@@ -814,7 +857,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn register_and_invoke_simple_operation() { async fn register_and_invoke_simple_operation() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("echo", AccessControl::default()), external_spec("echo", AccessControl::default()),
@@ -835,7 +878,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn internal_op_from_external_call_returns_not_found() { async fn internal_op_from_external_call_returns_not_found() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
internal_spec("secret", AccessControl::default()), internal_spec("secret", AccessControl::default()),
@@ -859,7 +902,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn internal_op_from_internal_call_invokes_handler() { async fn internal_op_from_internal_call_invokes_handler() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
internal_spec("secret", AccessControl::default()), internal_spec("secret", AccessControl::default()),
@@ -891,7 +934,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn acl_sufficient_scopes_allowed() { async fn acl_sufficient_scopes_allowed() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec( external_spec(
@@ -921,7 +964,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn acl_insufficient_scopes_forbidden() { async fn acl_insufficient_scopes_forbidden() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec( external_spec(
@@ -957,7 +1000,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn acl_restricted_op_no_identity_forbidden() { async fn acl_restricted_op_no_identity_forbidden() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec( external_spec(
@@ -987,7 +1030,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn internal_call_acl_uses_handler_identity() { async fn internal_call_acl_uses_handler_identity() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let composing_authority = CompositionAuthority::new("agent-chat", ["admin".to_string()]); let composing_authority = CompositionAuthority::new("agent-chat", ["admin".to_string()]);
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
@@ -1021,7 +1064,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn internal_call_acl_insufficient_handler_identity_forbidden() { async fn internal_call_acl_insufficient_handler_identity_forbidden() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let weak_authority = CompositionAuthority::new("weak", ["user".to_string()]); let weak_authority = CompositionAuthority::new("weak", ["user".to_string()]);
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
@@ -1058,7 +1101,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn external_call_acl_uses_caller_identity_not_handler_identity() { async fn external_call_acl_uses_caller_identity_not_handler_identity() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let handler_authority = CompositionAuthority::new("agent", ["admin".to_string()]); let handler_authority = CompositionAuthority::new("agent", ["admin".to_string()]);
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
@@ -1092,7 +1135,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn list_operations_returns_external_only() { async fn list_operations_returns_external_only() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("echo", AccessControl::default()), external_spec("echo", AccessControl::default()),
@@ -1120,7 +1163,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn handler_returned_error_passes_through() { async fn handler_returned_error_passes_through() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("boom", AccessControl::default()), external_spec("boom", AccessControl::default()),
@@ -1247,7 +1290,7 @@ mod tests {
#[test] #[test]
fn registration_lookup_returns_bundle_fields() { fn registration_lookup_returns_bundle_fields() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let authority = CompositionAuthority::new("agent", ["fs:read".to_string()]); let authority = CompositionAuthority::new("agent", ["fs:read".to_string()]);
let scoped = ScopedPeerEnv::new(["fs/readFile"]); let scoped = ScopedPeerEnv::new(["fs/readFile"]);
let caps = Capabilities::new().with_api_key("google", "k".to_string()); let caps = Capabilities::new().with_api_key("google", "k".to_string());
@@ -1289,7 +1332,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_on_stream_kind_returns_invalid_operation_type() { async fn invoke_on_stream_kind_returns_invalid_operation_type() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
subscription_spec("events/stream"), subscription_spec("events/stream"),
@@ -1312,7 +1355,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_on_once_kind_dispatches_normally() { async fn invoke_on_once_kind_dispatches_normally() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("echo", AccessControl::default()), external_spec("echo", AccessControl::default()),
@@ -1332,7 +1375,7 @@ mod tests {
#[test] #[test]
fn register_rejects_once_for_subscription_spec() { fn register_rejects_once_for_subscription_spec() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let result = registry.register(HandlerRegistration::new( let result = registry.register(HandlerRegistration::new(
subscription_spec("events/stream"), subscription_spec("events/stream"),
HandlerKind::Once(echo_handler()), HandlerKind::Once(echo_handler()),
@@ -1354,7 +1397,7 @@ mod tests {
#[test] #[test]
fn register_rejects_stream_for_query_spec() { fn register_rejects_stream_for_query_spec() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let result = registry.register(HandlerRegistration::new( let result = registry.register(HandlerRegistration::new(
external_spec("echo", AccessControl::default()), external_spec("echo", AccessControl::default()),
HandlerKind::Stream(echo_streaming_handler()), HandlerKind::Stream(echo_streaming_handler()),
@@ -1405,7 +1448,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_streaming_on_subscription_dispatches_handler_stream() { async fn invoke_streaming_on_subscription_dispatches_handler_stream() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
subscription_spec("events/stream"), subscription_spec("events/stream"),
@@ -1442,7 +1485,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_streaming_on_query_op_yields_invalid_operation_type() { async fn invoke_streaming_on_query_op_yields_invalid_operation_type() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("echo", AccessControl::default()), external_spec("echo", AccessControl::default()),
@@ -1465,7 +1508,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_streaming_internal_op_from_external_yields_not_found() { async fn invoke_streaming_internal_op_from_external_yields_not_found() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
internal_subscription_spec(AccessControl::default()), internal_subscription_spec(AccessControl::default()),
@@ -1491,7 +1534,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_streaming_acl_denied_yields_forbidden() { async fn invoke_streaming_acl_denied_yields_forbidden() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
subscription_spec_with_acl(AccessControl { subscription_spec_with_acl(AccessControl {
@@ -1526,7 +1569,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_streaming_internal_call_uses_handler_identity_for_acl() { async fn invoke_streaming_internal_call_uses_handler_identity_for_acl() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let composing_authority = CompositionAuthority::new("agent-chat", ["admin".to_string()]); let composing_authority = CompositionAuthority::new("agent-chat", ["admin".to_string()]);
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
@@ -1689,7 +1732,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_with_ownership_provider_allows_owned_resource() { async fn invoke_with_ownership_provider_allows_owned_resource() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let acl = AccessControl { let acl = AccessControl {
resource_type: Some("container".to_string()), resource_type: Some("container".to_string()),
resource_action: Some("exec".to_string()), resource_action: Some("exec".to_string()),
@@ -1726,7 +1769,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_with_ownership_provider_forbids_unowned_resource() { async fn invoke_with_ownership_provider_forbids_unowned_resource() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let acl = AccessControl { let acl = AccessControl {
resource_type: Some("container".to_string()), resource_type: Some("container".to_string()),
resource_action: Some("exec".to_string()), resource_action: Some("exec".to_string()),
@@ -1769,7 +1812,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_with_ownership_provider_missing_field_falls_back_to_owns_any() { async fn invoke_with_ownership_provider_missing_field_falls_back_to_owns_any() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let acl = AccessControl { let acl = AccessControl {
resource_type: Some("container".to_string()), resource_type: Some("container".to_string()),
resource_action: Some("exec".to_string()), resource_action: Some("exec".to_string()),
@@ -1805,7 +1848,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_with_ownership_provider_missing_field_forbids_when_not_owns_any() { async fn invoke_with_ownership_provider_missing_field_forbids_when_not_owns_any() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let acl = AccessControl { let acl = AccessControl {
resource_type: Some("container".to_string()), resource_type: Some("container".to_string()),
resource_action: Some("exec".to_string()), resource_action: Some("exec".to_string()),
@@ -1844,7 +1887,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_without_ownership_provider_falls_back_to_static_resources() { async fn invoke_without_ownership_provider_falls_back_to_static_resources() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let acl = AccessControl { let acl = AccessControl {
resource_type: Some("container".to_string()), resource_type: Some("container".to_string()),
resource_action: Some("exec".to_string()), resource_action: Some("exec".to_string()),
@@ -1906,7 +1949,7 @@ mod tests {
} }
fn registering_pub_with_schema(schema: Value) -> Result<(), String> { fn registering_pub_with_schema(schema: Value) -> Result<(), String> {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry.register(HandlerRegistration::new( registry.register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()).with_publish_schema(schema), pub_spec("fs/upload", AccessControl::default()).with_publish_schema(schema),
HandlerKind::Sink(make_sink_noop()), HandlerKind::Sink(make_sink_noop()),
@@ -1948,7 +1991,7 @@ mod tests {
#[test] #[test]
fn register_accepts_compilable_publish_schema_and_caches_validator() { fn register_accepts_compilable_publish_schema_and_caches_validator() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()).with_publish_schema( pub_spec("fs/upload", AccessControl::default()).with_publish_schema(
@@ -1973,7 +2016,7 @@ mod tests {
#[test] #[test]
fn publish_validator_none_when_no_schema_declared() { fn publish_validator_none_when_no_schema_declared() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()), pub_spec("fs/upload", AccessControl::default()),
@@ -2003,7 +2046,7 @@ mod tests {
#[test] #[test]
fn failed_validation_leaves_no_validator_cache_entry() { fn failed_validation_leaves_no_validator_cache_entry() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let result = registry.register(HandlerRegistration::new( let result = registry.register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()) pub_spec("fs/upload", AccessControl::default())
.with_publish_schema(serde_json::json!({ "type": "object", "required": "bad" })), .with_publish_schema(serde_json::json!({ "type": "object", "required": "bad" })),
@@ -2067,7 +2110,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_pub_op_collects_stream_and_returns_result() { async fn invoke_sink_pub_op_collects_stream_and_returns_result() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()), pub_spec("fs/upload", AccessControl::default()),
@@ -2094,7 +2137,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_empty_stream_returns_zero_count() { async fn invoke_sink_empty_stream_returns_zero_count() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()), pub_spec("fs/upload", AccessControl::default()),
@@ -2130,7 +2173,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_internal_op_from_external_returns_not_found() { async fn invoke_sink_internal_op_from_external_returns_not_found() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
internal_pub_spec(AccessControl::default()), internal_pub_spec(AccessControl::default()),
@@ -2154,7 +2197,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_acl_denied_returns_forbidden() { async fn invoke_sink_acl_denied_returns_forbidden() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec( pub_spec(
@@ -2193,7 +2236,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_on_pub_op_returns_invalid_operation_type() { async fn invoke_on_pub_op_returns_invalid_operation_type() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()), pub_spec("fs/upload", AccessControl::default()),
@@ -2216,7 +2259,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_streaming_on_pub_op_returns_invalid_operation_type() { async fn invoke_streaming_on_pub_op_returns_invalid_operation_type() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()), pub_spec("fs/upload", AccessControl::default()),
@@ -2239,7 +2282,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_on_query_op_returns_invalid_operation_type() { async fn invoke_sink_on_query_op_returns_invalid_operation_type() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
external_spec("echo", AccessControl::default()), external_spec("echo", AccessControl::default()),
@@ -2263,7 +2306,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_on_sub_op_returns_invalid_operation_type() { async fn invoke_sink_on_sub_op_returns_invalid_operation_type() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
subscription_spec("events/stream"), subscription_spec("events/stream"),
@@ -2287,7 +2330,7 @@ mod tests {
#[test] #[test]
fn register_rejects_once_for_pub_spec() { fn register_rejects_once_for_pub_spec() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let result = registry.register(HandlerRegistration::new( let result = registry.register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()), pub_spec("fs/upload", AccessControl::default()),
HandlerKind::Once(echo_handler()), HandlerKind::Once(echo_handler()),
@@ -2309,7 +2352,7 @@ mod tests {
#[test] #[test]
fn register_rejects_stream_for_pub_spec() { fn register_rejects_stream_for_pub_spec() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let result = registry.register(HandlerRegistration::new( let result = registry.register(HandlerRegistration::new(
pub_spec("fs/upload", AccessControl::default()), pub_spec("fs/upload", AccessControl::default()),
HandlerKind::Stream(echo_streaming_handler()), HandlerKind::Stream(echo_streaming_handler()),
@@ -2331,7 +2374,7 @@ mod tests {
#[test] #[test]
fn register_rejects_sink_for_query_spec() { fn register_rejects_sink_for_query_spec() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let result = registry.register(HandlerRegistration::new( let result = registry.register(HandlerRegistration::new(
external_spec("echo", AccessControl::default()), external_spec("echo", AccessControl::default()),
HandlerKind::Sink(collect_sink_handler()), HandlerKind::Sink(collect_sink_handler()),
@@ -2351,7 +2394,7 @@ mod tests {
#[test] #[test]
fn register_rejects_sink_for_sub_spec() { fn register_rejects_sink_for_sub_spec() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let result = registry.register(HandlerRegistration::new( let result = registry.register(HandlerRegistration::new(
subscription_spec("events/stream"), subscription_spec("events/stream"),
HandlerKind::Sink(collect_sink_handler()), HandlerKind::Sink(collect_sink_handler()),
@@ -2373,7 +2416,7 @@ mod tests {
#[tokio::test] #[tokio::test]
async fn invoke_sink_internal_call_uses_handler_identity_for_acl() { async fn invoke_sink_internal_call_uses_handler_identity_for_acl() {
let mut registry = OperationRegistry::new(); let registry = OperationRegistry::new();
let composing_authority = CompositionAuthority::new("fs-writer", ["fs:write".to_string()]); let composing_authority = CompositionAuthority::new("fs-writer", ["fs:write".to_string()]);
registry registry
.register(HandlerRegistration::new( .register(HandlerRegistration::new(
@@ -2454,4 +2497,138 @@ mod tests {
assert_eq!(reg.provenance, OperationProvenance::FromCall); assert_eq!(reg.provenance, OperationProvenance::FromCall);
assert!(matches!(reg.handler, HandlerKind::Sink(_))); assert!(matches!(reg.handler, HandlerKind::Sink(_)));
} }
// --- review 004 F-02/F-03: fork surface --------------------------------
#[tokio::test]
async fn fork_carries_handlers_validators_and_specs() {
let registry = OperationRegistryBuilder::new()
.with_local(
external_spec("echo/run", AccessControl::default()),
echo_handler(),
None,
None,
Capabilities::new(),
)
.unwrap()
.with_local_sink(
pub_spec("fs/upload", AccessControl::default()),
collect_sink_handler(),
CompositionAuthority::none(),
None,
Capabilities::new(),
)
.unwrap()
.build();
let fork = registry.fork();
assert!(fork.registration("echo/run").is_some());
assert!(fork.registration("fs/upload").is_some());
assert!(
fork.publish_validator("fs/upload").is_none(),
"op has no publish_schema — validator cache entry is None either way"
);
assert_eq!(fork.list_operations().len(), 2);
assert!(
fork.list_operations().iter().any(|s| s.name == "echo/run"),
"External op appears in list_operations"
);
let ctx = root_context("req-fork-1", None, None, false, ScopedPeerEnv::empty());
let response = fork
.invoke("echo/run", serde_json::json!({"v": 1}), ctx)
.await;
assert_eq!(response.result.unwrap(), serde_json::json!({"v": 1}));
}
#[tokio::test]
async fn fork_publish_schema_validator_is_carried_and_enforced() {
let schema = serde_json::json!({
"type": "object",
"properties": { "n": { "type": "integer" } },
"required": ["n"]
});
let spec = pub_spec("fs/upload", AccessControl::default()).with_publish_schema(schema);
let registry = OperationRegistryBuilder::new()
.with_local_sink(
spec,
collect_sink_handler(),
CompositionAuthority::none(),
None,
Capabilities::new(),
)
.unwrap()
.build();
let fork = registry.fork();
assert!(
fork.publish_validator("fs/upload").is_some(),
"F-03: validators map must be carried by the fork"
);
}
#[tokio::test]
async fn fork_is_independent_both_directions() {
let registry = OperationRegistry::new();
registry
.register(HandlerRegistration::new(
external_spec("base/op", AccessControl::default()),
HandlerKind::Once(echo_handler()),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
let fork = registry.fork();
fork.register(HandlerRegistration::new(
external_spec("fork/op", AccessControl::default()),
HandlerKind::Once(echo_handler()),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
assert!(
registry.registration("fork/op").is_none(),
"fork registrations do not leak to the base"
);
assert!(
fork.registration("base/op").is_some(),
"base ops are on the fork"
);
}
#[tokio::test]
async fn register_through_shared_arc_after_dispatch_started() {
let registry = Arc::new(OperationRegistry::new());
registry
.register(HandlerRegistration::new(
external_spec("early/op", AccessControl::default()),
HandlerKind::Once(echo_handler()),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
))
.unwrap();
let late = Arc::clone(&registry);
let late_registration = HandlerRegistration::new(
external_spec("late/op", AccessControl::default()),
HandlerKind::Once(echo_handler()),
OperationProvenance::Local,
None,
None,
Capabilities::new(),
);
late.register(late_registration).unwrap();
let ctx = root_context("req-arc-1", None, None, false, ScopedPeerEnv::empty());
let response = registry.invoke("late/op", serde_json::json!({}), ctx).await;
assert!(response.result.is_ok());
}
} }