docs(arch): hub-wiring cluster — aggregated env, peer_operations, from_call cleanup, alknet-hub crate spec

Three ADRs addressing gaps surfaced by the first hub consumer (alkapi):

- ADR-067: Aggregated peer-env wiring — Dispatcher::with_aggregated_env hook
  so compose_root_env reads a shared PeerCompositeEnv across all calls,
  not a fresh per-call one. The hub-defining gap (alkapi OQ-08 / G.1).

- ADR-068: PeerCompositeEnv::peer_operations override — adds
  list_operation_names() to OperationEnv, overrides on OverlayOperationEnv
  and PeerCompositeEnv. Fixes services/list-peers returning empty operation
  lists for non-local peers (alkapi G.6).

- ADR-069: from_call is a manual free function, not auto-wired — reverses
  the aspirational OQ-27 resolution to match the implementation. Cleans
  the 'v1 default' hedging language in ADR-017 and client-and-adapters.md
  (alkapi G.4).

New crate spec: alknet-hub — reusable hub pattern (aggregated env,
connection lifecycle, worker supervision with backoff, service discovery).

Spec fixes: builder API drift (with_local/with_local_streaming separation),
ScopedOperationEnv → ScopedPeerEnv type name, OQ count 51→54.

Three new OQs: OQ-52 (wait_for_close), OQ-53 (backoff defaults),
OQ-54 (inbound hook placement).
This commit is contained in:
glm-5.2 committed 2026-07-09 16:01:30 +00:00
1 parent 2282647f8b
commit 87b2c2a5ef
15 files changed
+1052 -47

No files matched your search

+5 -1
View File
@@ -97,6 +97,7 @@ adapter location map is now consistent: all HTTP-backed adapters
| [crates/vault/encryption.md](crates/vault/encryption.md) | stable | AES-256-GCM, EncryptedData, key versioning, salt (Phase B reserved) |
| [crates/vault/service.md](crates/vault/service.md) | stable | VaultServiceHandle lifecycle, direct dispatch, cache, error model |
| [crates/vault/protocol.md](crates/vault/protocol.md) | stable | DerivedKey redaction, KeyType, serialization behavior |
| [crates/hub/README.md](crates/hub/README.md) | draft | alknet-hub crate — aggregated peer env, connection lifecycle, worker supervision, service discovery |
## ADR Table
@@ -168,10 +169,13 @@ adapter location map is now consistent: all HTTP-backed adapters
| [064](decisions/064-irpc-never-integrated-hand-rolled-framing.md) | irpc Was Never Integrated — Hand-Rolled EventEnvelope Framing | Accepted (supersedes ADR-005) |
| [065](decisions/065-connection-from-stream-generic-single-stream.md) | `Connection::from_stream` — Generic Single-Stream Connections | Accepted |
| [066](decisions/066-from-jsonschema-as-http-adapter.md) | `from_jsonschema` as HTTP-Backed Single-Endpoint Adapter in alknet-http | Accepted (supersedes the `from_jsonschema` clause of ADR-017 §5 and the `FromJsonSchema` provenance row of ADR-022) |
| [067](decisions/067-aggregated-peer-env-wiring.md) | Aggregated Peer-Environment Wiring for Hub Deployments | Proposed |
| [068](decisions/068-peer-composite-env-peer-operations.md) | PeerCompositeEnv::peer_operations Override | Proposed |
| [069](decisions/069-from-call-manual-free-function.md) | from_call Is a Manual Free Function, Not Auto-Wired | Proposed |
## Open Questions
Open questions are tracked in [open-questions.md](open-questions.md) — an index of theme-grouped tables (51 OQs across 16 themes) with a cross-theme [Deferred / Blocked](open-questions.md#deferred--blocked) section surfacing the safe-exit deferrals. Each OQ lives in its own file under [`questions/`](questions/) (`NNN-slug.md`, mirroring the ADR convention).
Open questions are tracked in [open-questions.md](open-questions.md) — an index of theme-grouped tables (54 OQs across 17 themes) with a cross-theme [Deferred / Blocked](open-questions.md#deferred--blocked) section surfacing the safe-exit deferrals. Each OQ lives in its own file under [`questions/`](questions/) (`NNN-slug.md`, mirroring the ADR convention).
## Document Lifecycle
+1 -1
View File
@@ -58,7 +58,7 @@ Structured RPC: operations, request/response, streaming subscriptions, and servi
| OQ-19 | Session-scoped operation registries | resolved | Agent-written operations overlaid on curated registry via `OperationEnv` trait layering. Protocol doesn't need changes; `OperationEnv` must remain a trait. Generalized by ADR-024 to cover connection-scoped overlays. |
| OQ-25 | ~~Remote-safe marking shape~~ | **dissolved** (ADR-029) | `remote_safe`/`trusted_peer` retired; peer authorization is `AccessControl::check(peer_identity)` |
| OQ-26 | OperationAdapter error type (AdapterError variants) | **resolved** | `DiscoveryFailed`, `SchemaParse`, `Transport`, `Unauthorized`, `SamePeerCollision`; `#[non_exhaustive]` |
| OQ-27 | from_call re-import trigger | **resolved** | Auto-re-import on connection establishment; `refresh()` is a feature addition |
| OQ-27 | from_call re-import trigger | **resolved** | `from_call` is a manual free function; the assembly layer calls it after `connect()`. `refresh()` is a genuine feature addition. See ADR-069. |
| OQ-28 | from_call namespace collision | **resolved** | Same-peer collision = error; cross-peer dissolved by ADR-029 (separate sub-overlays) |
| OQ-29 | CallClient TLS client-auth | **resolved** | Wire quinn client-auth; key-type-aware server cert verification; fingerprint normalization |
| OQ-30 | `PeerRef::Any` routing policy | **resolved** | Insertion-order first-match; richer routing is a feature extension |
@@ -474,7 +474,7 @@ fn build_root_context(
metadata: HashMap::new(), // fresh per request
deadline: Some(Instant::now() + self.default_timeout), // root deadline (W7)
scoped_env: registration.scoped_env.clone()
.unwrap_or_else(ScopedOperationEnv::empty), // from the bundle, empty for leaves
.unwrap_or_else(ScopedPeerEnv::empty), // from the bundle, empty for leaves
// Per-call env composition (ADR-024 + ADR-029): the root env is a
// PeerCompositeEnv — the curated base + this connection's imported-
// ops overlay (peer-keyed in the head's aggregation env, ADR-029 §1)
@@ -590,8 +590,9 @@ See [open-questions.md](../../open-questions.md) for full details.
variants (`DiscoveryFailed`, `SchemaParse`, `Transport`, `Unauthorized`,
`SamePeerCollision`); `#[non_exhaustive]`. See
[client-and-adapters.md](client-and-adapters.md).
- **OQ-27** (resolved): `from_call` re-import trigger — auto-re-import on
connection establishment. See [client-and-adapters.md](client-and-adapters.md).
- **OQ-27** (resolved): `from_call` re-import trigger — `from_call` is a manual
free function; the assembly layer calls it after `connect()`. See
[ADR-069](../../decisions/069-from-call-manual-free-function.md).
- **OQ-28** (resolved): `from_call` namespace collision — same-peer collision
= error; cross-peer dissolved by ADR-029 (separate sub-overlays). See
[client-and-adapters.md](client-and-adapters.md).
@@ -352,11 +352,14 @@ The flow (ADR-017 §3):
4. The caller registers the bundles via
`CallConnection::register_imported_all()`.
**Re-import on reconnection** (DC-2, OQ-27): `from_call` runs automatically on
connection establishment. The overlay is per-connection (Layer 2, ADR-024), so
a stale overlay dies with the connection; re-import on reconnect is naturally
scoped to the new connection. This is the v1 default; explicit re-import via a
future `CallConnection::refresh()` is additive.
**Re-import on reconnection** (DC-2, OQ-27): `from_call` is a free function;
the assembly layer calls it after `connect()`. The overlay is per-connection
(Layer 2, ADR-024), so a stale overlay dies with the connection; re-import on
reconnect is naturally scoped to the new connection. A
`CallConnection::refresh()` method for mid-connection re-discovery is a
genuine feature addition — non-breaking, additive — if a deployment needs
manual re-discovery without drop-and-reconnect. See
[ADR-069](../../decisions/069-from-call-manual-free-function.md).
**Namespace collision** (DC-3, OQ-28): under the peer-graph model (ADR-029),
cross-peer collision dissolves — same name on different peers is fine (they
@@ -649,8 +652,10 @@ Based on the gap analysis and the downstream unblock chain:
holds a `PeerCompositeEnv` with `connections: HashMap<PeerId, Arc<dyn OperationEnv>>`,
not a singular connection overlay. `invoke_peer()` routes to the right peer
via `PeerRef::Specific` / `PeerRef::Any` (ADR-029 §1-2).
- **`from_call` re-import is auto-on-reconnect.** v1 default; the overlay is
per-connection so re-import is naturally scoped (DC-2, OQ-27).
- **`from_call` is a manual free function.** The assembly layer calls it after
`connect()`. The overlay is per-connection so re-import on reconnect is
naturally scoped (DC-2, OQ-27). See
[ADR-069](../../decisions/069-from-call-manual-free-function.md).
- **`from_call` namespace collision is same-peer only.** Cross-peer collision
dissolves (same name on different peers is fine — separate sub-overlays,
ADR-029 §5). Same-peer collision stays an error. `namespace_prefix` is
@@ -713,9 +718,10 @@ See [open-questions.md](../../open-questions.md) for full details.
- **OQ-26** (resolved): `AdapterError` variants — `DiscoveryFailed`,
`SchemaParse`, `Transport`, `Unauthorized`, `SamePeerCollision`
(replaces flat `Conflict`). `#[non_exhaustive]`.
- **OQ-27** (resolved): `from_call` re-import trigger — auto-re-import on
connection establishment. `CallConnection::refresh()` is a feature
addition, not an unmade decision.
- **OQ-27** (resolved): `from_call` re-import trigger — `from_call` is a manual
free function; the assembly layer calls it after `connect()`. A
`CallConnection::refresh()` method is a genuine feature addition —
non-breaking, additive. See [ADR-069](../../decisions/069-from-call-manual-free-function.md).
- **OQ-28** (resolved): `from_call` namespace collision — same-peer
collision = error; cross-peer dissolved by ADR-029 (separate sub-overlays).
`namespace_prefix` is optional local-naming sugar.
@@ -243,7 +243,7 @@ pub struct OperationContext {
/// Populated from the registration bundle's `scoped_env` (ADR-022).
/// The reachability check in `OperationEnv::invoke()` consults
/// `scoped_env.allows(&name)`. This is data, not a dispatch trait.
pub scoped_env: ScopedOperationEnv,
pub scoped_env: ScopedPeerEnv,
/// Composition dispatch trait. A handler calls `env.invoke(...)` to
/// compose child operations. This is `Arc<dyn OperationEnv>` (a trait
/// object), not a concrete struct — the trait-object design is what
@@ -303,7 +303,7 @@ impl OperationContext {
- `forwarded_for`: The original caller when this call was forwarded by a `from_call` handler (ADR-032). **Metadata only** — `AccessControl::check` never reads it; the ACL always authorizes the direct caller's `identity`. Handlers may read it for logging, auditing, per-user rate limiting, or application context. Populated from `call.requested.forwarded_for` by the dispatch path; set to `None` for composed children (wire-ingress only). The forwarder's claim, not a verified identity — a malicious hub can lie (same property as HTTP `X-Forwarded-For`). See ADR-032.
- `capabilities`: Outbound credentials the handler may use (decrypted API keys, scoped vault access) — see [Capability Injection](#capability-injection) below
- `metadata`: Request-scoped context (tracing IDs, connection info). **Must not hold secret material** — see ADR-014. **Does not propagate through `OperationEnv::invoke()`** — nested calls get fresh metadata. The tracing link between parent and child is `parent_request_id`, not metadata propagation. Anything a handler needs to pass to a child goes in the call `input`.
- `scoped_env`: The reachability set — the operations this handler may compose. Populated from the registration bundle's `scoped_env` (ADR-022). The reachability check in `OperationEnv::invoke()` consults `scoped_env.allows(&name)`. This is *data* (a `ScopedOperationEnv` struct), not a dispatch trait. `None`/empty for leaves.
- `scoped_env`: The reachability set — the operations this handler may compose. Populated from the registration bundle's `scoped_env` (ADR-022). The reachability check in `OperationEnv::invoke()` consults `scoped_env.allows(&name)`. This is *data* (a `ScopedPeerEnv` struct), not a dispatch trait. `None`/empty for leaves.
- `env`: The composition dispatch trait (`Arc<dyn OperationEnv + Send + Sync>`). A handler calls `context.env.invoke(...)` to compose child operations. This is a trait object, not a concrete struct — the trait-object design enables registry layering (ADR-024): the CallAdapter composes the root env per call from the active layers (curated base + connection overlay + session overlay), and overlays wrap the base via trait layering. Same pattern as `IdentityProvider` (ADR-004). See ADR-024.
- `internal`: When `true`, this call originated from composition (a handler calling another operation via `OperationEnv`), not from a wire request. This switches the authority context: ACL runs against `handler_identity`, not `identity`. The `internal` field uses module-private construction — handlers construct `OperationContext` through `OperationEnv::invoke()` which sets `internal: true`, or through the `CallAdapter` dispatch path which sets `internal: false`. The field is not `pub` for writes; only `pub fn is_internal(&self) -> bool` is exposed for reads. See ADR-015.
@@ -433,24 +433,37 @@ impl CompositionAuthority {
- `scoped_env`: The set of operations this handler may reach via `env.invoke()`. `None` for leaves (empty env). The reachability control from ADR-015.
- `capabilities`: Outbound credentials (decrypted API keys, signing keys). Populated by the assembly layer from the vault at registration time. See [Capability Injection](#capability-injection).
The `OperationRegistryBuilder` provides a fluent API with convenience methods for common cases. The builder absorbs the `HandlerKind` wrapping internally — `.with_local()` and `.with_leaf()` take the raw `Handler` (or `StreamingHandler`) and wrap it in the right `HandlerKind` based on `spec.op_type` (ADR-049):
The `OperationRegistryBuilder` provides a fluent API with convenience methods for common cases. The builder validates handler kind against `spec.op_type` at registration time — `with_local` / `with_leaf` accept `Handler` (for `Query`/`Mutation` ops), `with_local_streaming` / `with_leaf_streaming` accept `StreamingHandler` (for `Subscription` ops). Passing a `StreamingHandler` to `with_local` or a `Handler` to `with_local_streaming` is a registration-time error:
```rust
// with_local: Local provenance, full bundle — all 5 args required.
// Accepts Handler (for Query/Mutation ops). Validates op_type at registration.
// with_local(spec, handler, composition_authority, scoped_env, capabilities)
// The builder inspects spec.op_type and wraps in HandlerKind::Once
// (Query/Mutation) or HandlerKind::Stream (Subscription) automatically.
// with_local_streaming: Local provenance, full bundle — all 5 args required.
// Accepts StreamingHandler (for Subscription ops). Validates op_type at registration.
// with_local_streaming(spec, streaming_handler, composition_authority, scoped_env, capabilities)
// with_leaf: Leaf provenance (default FromOpenAPI), no composition authority.
// Accepts Handler (for Query/Mutation ops).
// with_leaf(spec, handler, capabilities)
// with_leaf_streaming: Leaf provenance (default FromOpenAPI), no composition authority.
// Accepts StreamingHandler (for Subscription ops).
// with_leaf_streaming(spec, streaming_handler, capabilities)
// with_leaf_provenance / with_leaf_streaming_provenance: explicit provenance variant.
let registry = OperationRegistryBuilder::new()
// Built-in service discovery (Local, no composition — empty authority, empty env, empty caps)
.with_local(services_list_spec(), Arc::new(services_list_handler),
CompositionAuthority::none(), ScopedOperationEnv::empty(), Capabilities::new())
CompositionAuthority::none(), ScopedPeerEnv::empty(), Capabilities::new())
.with_local(services_schema_spec(), Arc::new(schema_handler),
CompositionAuthority::none(), ScopedOperationEnv::empty(), Capabilities::new())
CompositionAuthority::none(), ScopedPeerEnv::empty(), Capabilities::new())
// Agent handler (Local, Subscription — streams call.responded as the
// LLM generates tokens; builder wraps in HandlerKind::Stream)
.with_local(agent_chat_spec(), Arc::new(agent_chat_streaming_handler),
// LLM generates tokens; uses with_local_streaming for the StreamingHandler)
.with_local_streaming(agent_chat_spec(), Arc::new(agent_chat_streaming_handler),
CompositionAuthority::new("agent-chat", ["llm:call", "fs:read", "vastai:query"]),
ScopedOperationEnv::new(["fs/readFile", "vastai/listMachines", "llm/generate"]),
ScopedPeerEnv::new(["fs/readFile", "vastai/listMachines", "llm/generate"]),
Capabilities::new().with_api_key("google", google_api_key))
// Imported ops (leaves — no authority, no scoped env; capabilities for outbound HTTP)
.with_leaf(vastai_listMachines_spec(), Arc::new(vastai_handler), vastai_credentials)
@@ -621,7 +634,7 @@ impl OperationEnv for LocalOperationEnv {
abort_policy: policy, // Explicit policy (from invoke() default or invoke_with_policy)
deadline: parent.deadline, // Inherit parent's deadline (children don't get a fresh 30s)
scoped_env: registration.scoped_env.clone()
.unwrap_or_else(ScopedOperationEnv::empty), // Child's own scoped env (empty for leaves)
.unwrap_or_else(ScopedPeerEnv::empty), // Child's own scoped env (empty for leaves)
// Dispatch trait: the child inherits the parent's env (the same
// composite of curated base + active overlays). See ADR-024.
env: parent.env.clone(),
@@ -820,9 +833,9 @@ let vastai_credentials = Capabilities::new().with_http_token("vastai", vastai_to
let registry = OperationRegistryBuilder::new()
// Built-in service discovery (Local, no composition — empty caps)
.with_local(services_list_spec(), Arc::new(services_list_handler),
CompositionAuthority::none(), ScopedOperationEnv::empty(), Capabilities::new())
CompositionAuthority::none(), ScopedPeerEnv::empty(), Capabilities::new())
.with_local(services_schema_spec(), Arc::new(schema_handler),
CompositionAuthority::none(), ScopedOperationEnv::empty(), Capabilities::new())
CompositionAuthority::none(), ScopedPeerEnv::empty(), Capabilities::new())
// Agent handler (Local, Subscription — composes; streaming handler
// wrapped in HandlerKind::Stream by the builder per ADR-049)
.with(HandlerRegistration {
@@ -831,7 +844,7 @@ let registry = OperationRegistryBuilder::new()
provenance: OperationProvenance::Local,
composition_authority: Some(CompositionAuthority::new(
"agent-chat", ["llm:call", "fs:read", "vastai:query"])),
scoped_env: Some(ScopedOperationEnv::new(
scoped_env: Some(ScopedPeerEnv::new(
["fs/readFile", "vastai/listMachines", "llm/generate"])),
capabilities: Capabilities::new().with_api_key("google", google_api_key),
})
@@ -940,9 +953,10 @@ See [open-questions.md](../../open-questions.md) for full details.
variants: `DiscoveryFailed`, `SchemaParse`, `Transport`, `Unauthorized`,
`SamePeerCollision` (replaces flat `Conflict`). `#[non_exhaustive]`. See
[client-and-adapters.md](client-and-adapters.md).
- **OQ-27** (resolved): `from_call` re-import trigger — auto-re-import on
connection establishment. `CallConnection::refresh()` is a feature
addition, not an unmade decision. See [client-and-adapters.md](client-and-adapters.md).
- **OQ-27** (resolved): `from_call` re-import trigger — `from_call` is a manual
free function; the assembly layer calls it after `connect()`. A
`CallConnection::refresh()` method is a genuine feature addition —
non-breaking, additive. See [ADR-069](../../decisions/069-from-call-manual-free-function.md).
- **OQ-28** (resolved): `from_call` namespace collision — same-peer
collision = error; cross-peer dissolved by ADR-029 (separate sub-overlays).
`namespace_prefix` is optional local-naming sugar. See
+408
View File
@@ -0,0 +1,408 @@
---
status: draft
last_updated: 2026-07-09
---
# alknet-hub
The hub pattern as a reusable crate: aggregated peer environment, connection
lifecycle hooks, worker supervision, and service discovery across connected
peers. Consumes `alknet-call` types; provides the assembly-layer wiring every
hub needs.
## What
`alknet-hub` is a thin crate that provides the head/worker (hub-spoke) pattern
as a reusable library. It depends on `alknet-call` and `alknet-core`. It does
not introduce new types — it wires the existing `PeerCompositeEnv`,
`CallClient`, `CallAdapter`, `from_call`, and `Dispatcher` types into a
coherent hub runtime.
A hub is a call-protocol node that:
1. **Aggregates multiple worker connections** into a single
`PeerCompositeEnv` shared across all calls.
2. **Discovers worker operations** via `from_call` on connection establishment.
3. **Supervises worker connections** — reconnects on drop with backoff,
re-discovers ops on reconnect, detaches on disconnect.
4. **Exposes service discovery** — `services/list-peers` returns each
connected worker's operation list.
The crate provides these four capabilities. A downstream consumer (alkapi,
any future hub) constructs a `Hub`, registers its curated operations, and
starts accepting/dialing workers. The hub handles the rest.
## Why
The alkapi project identified that the hub pattern requires wiring that
alknet-call does not provide out of the box:
- `Dispatcher::compose_root_env` builds a fresh `PeerCompositeEnv` per call
with only the current connection — multi-worker aggregation is not wired
(ADR-067).
- `PeerCompositeEnv::peer_operations` is not overridden — `services/list-peers`
returns empty operation lists for non-local peers (ADR-068).
- `from_call` is a free function, not wired into `CallClient::connect` — the
assembly layer must call it manually after every connect (ADR-069).
- There is no worker supervision loop — reconnection, backoff, and
re-discovery are assembly-layer concerns.
These are not design flaws; they are the correct separation of concerns.
`alknet-call` provides the types and the routing logic; the hub wiring is a
consumer concern. But it is a concern *every* hub consumer shares. Rather than
each downstream project (alkapi, future hubs) building the same wiring
independently, `alknet-hub` provides it once, as a reusable crate.
## Architecture
### Hub struct
The `Hub` is the central type. It owns the aggregated `PeerCompositeEnv`, the
`OperationRegistry`, and the `Dispatcher`:
```rust
pub struct Hub {
registry: Arc<OperationRegistry>,
aggregated_env: Arc<RwLock<PeerCompositeEnv>>,
dispatcher: Dispatcher,
identity_provider: Arc<dyn IdentityProvider>,
}
```
Construction:
```rust
impl Hub {
pub fn new(
registry: Arc<OperationRegistry>,
identity_provider: Arc<dyn IdentityProvider>,
) -> Self {
let base: Arc<dyn OperationEnv + Send + Sync> =
Arc::new(LocalOperationEnv::new(Arc::clone(&registry)));
let aggregated_env = Arc::new(RwLock::new(PeerCompositeEnv::new(base)));
let dispatcher = Dispatcher::new(Arc::clone(&registry), Arc::clone(&identity_provider))
.with_aggregated_env(Arc::clone(&aggregated_env));
Self {
registry,
aggregated_env,
dispatcher,
identity_provider,
}
}
/// The shared aggregated PeerCompositeEnv. The assembly layer wires
/// this into CallAdapter::with_aggregated_env so every call's
/// compose_root_env sees all connected workers.
pub fn aggregated_env(&self) -> &Arc<RwLock<PeerCompositeEnv>> {
&self.aggregated_env
}
}
```
The `Hub` exposes builder methods for optional hooks (`with_session_source`,
`with_ownership_provider`, `with_timeout`) that delegate to the `Dispatcher`.
### Worker connection lifecycle
The hub provides two paths for worker connections:
#### Inbound workers (workers dial the hub)
The hub's `CallAdapter` accepts inbound connections. The hub provides an
`on_worker_connected` hook that the assembly layer wires into the accept path:
```rust
impl Hub {
/// Called by the assembly layer when a worker connects inbound (via
/// CallAdapter::handle). Runs from_call, registers the bundles in the
/// connection's overlay, and attaches the peer to the aggregated env.
pub async fn on_worker_connected(
&self,
connection: &CallConnection,
config: FromCallConfig,
) -> Result<PeerId, HubError> {
let peer_id = connection.identity()
.map(|id| id.id.clone())
.ok_or(HubError::NoPeerIdentity)?;
let bundles = from_call(connection, config).await?;
connection.register_imported_all(bundles);
self.aggregated_env
.write()
.unwrap()
.attach_peer(peer_id.clone(), connection.overlay_env());
Ok(peer_id)
}
/// Called by the assembly layer when a worker disconnects (run_loop exits).
pub fn on_worker_disconnected(&self, peer_id: &PeerId) {
self.aggregated_env.write().unwrap().detach_peer(peer_id);
}
}
```
#### Outbound workers (hub dials workers)
The hub provides a `dial_worker` method that combines `CallClient::connect`,
`from_call`, and `attach_peer` into a single operation. It also provides a
`supervise_worker` method that wraps `dial_worker` in a reconnect loop with
configurable backoff:
```rust
impl Hub {
/// Dial a worker, discover its operations, and attach it to the
/// aggregated env. Returns the worker's PeerId and the live
/// CallConnection.
pub async fn dial_worker(
&self,
addr: SocketAddr,
credentials: CallCredentials,
config: FromCallConfig,
) -> Result<(PeerId, CallConnection), HubError> {
let client = CallClient::new(
Arc::clone(&self.registry),
Arc::clone(&self.identity_provider),
);
let connection = client.connect(addr, credentials).await?;
let peer_id = connection.identity()
.map(|id| id.id.clone())
.ok_or(HubError::NoPeerIdentity)?;
let bundles = from_call(&connection, config).await?;
connection.register_imported_all(bundles);
self.aggregated_env
.write()
.unwrap()
.attach_peer(peer_id.clone(), connection.overlay_env());
Ok((peer_id, connection))
}
/// Supervise an outbound worker connection: dial, discover, attach.
/// On disconnect, detach and retry with backoff. Runs until the
/// Hub is dropped (the returned JoinHandle can be aborted).
pub fn supervise_worker(
self: &Arc<Self>,
addr: SocketAddr,
credentials: CallCredentials,
config: FromCallConfig,
backoff: BackoffConfig,
) -> tokio::task::JoinHandle<()> {
let hub = Arc::clone(self);
tokio::spawn(async move {
let mut retries = 0;
loop {
match hub.dial_worker(addr, credentials.clone(), config.clone()).await {
Ok((peer_id, connection)) => {
retries = 0;
// Wait for the connection to drop (run_loop exits).
// Committed interim (OQ-52): poll the underlying
// Connection's accept_bi() until it returns
// ConnectionClosed. A CallConnection::closed()
// method is the target resolution.
loop {
match connection.connection() {
Some(conn) => {
match conn.accept_bi().await {
Err(StreamError::ConnectionClosed)
| Err(StreamError::StreamClosed)
| Err(StreamError::Timeout) => break,
_ => continue,
}
}
None => break,
}
}
hub.on_worker_disconnected(&peer_id);
}
Err(e) => {
tracing::warn!(?e, retries, "worker dial failed; retrying");
}
}
let delay = backoff.delay_for(retries);
tokio::time::sleep(delay).await;
retries += 1;
}
})
}
}
```
### Backoff configuration
```rust
pub struct BackoffConfig {
pub initial: Duration,
pub max: Duration,
pub multiplier: f64,
}
impl BackoffConfig {
pub fn delay_for(&self, retries: u32) -> Duration {
let delay = self.initial.as_millis() as f64
* self.multiplier.powi(retries as i32);
let delay_ms = delay.min(self.max.as_millis() as f64) as u64;
Duration::from_millis(delay_ms)
}
}
impl Default for BackoffConfig {
fn default() -> Self {
Self {
initial: Duration::from_secs(1),
max: Duration::from_secs(60),
multiplier: 2.0,
}
}
}
```
### Service discovery
The hub registers the built-in service discovery operations (`services/list`,
`services/schema`, `services/list-peers`) automatically. The
`services/list-peers` handler returns each connected worker's operation list
via `PeerCompositeEnv::peer_operations` (ADR-068).
### HubError
```rust
#[derive(Debug, thiserror::Error)]
#[non_exhaustive]
pub enum HubError {
#[error("worker connection has no resolved peer identity")]
NoPeerIdentity,
#[error("from_call discovery failed: {0}")]
Discovery(#[from] AdapterError),
#[error("call client error: {0}")]
Client(#[from] ClientError),
}
```
## What the hub does NOT do
- **Worker authentication policy.** The hub resolves the worker's identity
via `IdentityProvider` (the existing mechanism). Whether a given identity is
*allowed* to connect as a worker is an `AccessControl` decision on the
assembly layer's curated ops — the hub does not add a separate worker-auth
layer.
- **Worker-specific routing policy.** `PeerRef::Any` uses insertion-order
first-match (ADR-029 §2). A richer `RoutingPolicy` (round-robin,
least-loaded) is a future extension behind the same `PeerRef` enum.
- **Multi-hop federation.** The hub is one-hop: workers connect to the hub,
the hub composes their ops. Worker A does not transitively see worker B's
ops through the hub unless the hub explicitly re-exports them (ADR-029
Assumption 5).
- **Ownership reaping.** Stale ownership entries from autonomously-dead
containers are tolerated (ADR-050 §4b). The hub does not add a reaper.
- **Worker provisioning.** The hub does not spawn workers, configure them, or
manage their lifecycle beyond connection supervision. Worker provisioning
is an assembly-layer concern.
## Crate dependencies
```
alknet-hub
├── alknet-call (CallClient, CallAdapter, Dispatcher, PeerCompositeEnv,
│ from_call, FromCallConfig, AdapterError, ClientError)
├── alknet-core (IdentityProvider, Connection, OperationRegistry)
├── tokio (spawn, time::sleep)
└── tracing (logging)
```
`alknet-hub` depends on `alknet-call`, which depends on `alknet-core`. No new
dependency edges. The crate is a consumer of the call protocol types, not a
new protocol handler.
## Assembly layer integration
A downstream hub (alkapi) uses `alknet-hub` like this:
```rust
// 1. Build the curated registry (Layer 0)
let registry = OperationRegistryBuilder::new()
.with_local(agent_chat_spec(), agent_chat_handler, ...)
.with_local(services_list_spec(), services_list_handler, ...)
// ... other curated ops
.build();
let registry = Arc::new(registry);
// 2. Create the hub
let hub = Arc::new(Hub::new(
Arc::clone(&registry),
Arc::clone(&identity_provider),
).with_ownership_provider(ownership_provider));
// 3. Start the inbound call-protocol listener (workers dial the hub)
let adapter = CallAdapter::new(Arc::clone(&registry), Arc::clone(&identity_provider))
.with_aggregated_env(hub.aggregated_env().clone());
// ... register adapter on the QUIC endpoint
// 4. Dial outbound workers (hub dials workers)
hub.supervise_worker(
dev1_addr,
dev1_credentials,
FromCallConfig::new(),
BackoffConfig::default(),
);
// 5. Start the HTTP listener (clients dial the hub)
// ... HttpAdapter with the same registry and identity_provider
```
The hub's `aggregated_env()` accessor returns the shared `Arc<RwLock<PeerCompositeEnv>>`
so the assembly layer can wire it into `CallAdapter::with_aggregated_env`.
## Design Decisions
| Decision | ADR | Summary |
|----------|-----|---------|
| Aggregated peer-env wiring | [ADR-067](../../decisions/067-aggregated-peer-env-wiring.md) | `Dispatcher::with_aggregated_env` hook; `compose_root_env` reads shared env |
| PeerCompositeEnv::peer_operations | [ADR-068](../../decisions/068-peer-composite-env-peer-operations.md) | `list_operation_names` trait method; `PeerCompositeEnv` override |
| from_call is manual | [ADR-069](../../decisions/069-from-call-manual-free-function.md) | `from_call` is a free function; the hub calls it after connect |
| Peer-graph routing model | [ADR-029](../../decisions/029-peer-graph-routing-model.md) | Peer-keyed overlays, `PeerRef` routing, `AccessControl`-based peer auth |
| PeerEntry and Identity.id | [ADR-030](../../decisions/030-peerentry-and-identity-id-decoupling.md) | `PeerId` = `Identity.id` = `PeerEntry.peer_id` (stable) |
## Open Questions
See [open-questions.md](../../open-questions.md) for full details.
- **OQ-52** (open): `CallConnection::wait_for_close()` — the supervision loop
needs a way to await connection close. Today `CallConnection` exposes
`connection()` (the underlying `Connection`) but not a "wait for run_loop
exit" future. The committed interim is polling `connection().accept_bi()`
in a loop until it returns `ConnectionClosed`. A `closed()` method on
`CallConnection` (signaled by the dispatcher on `run_loop` exit) is the
target resolution. See OQ-52.
- **OQ-53** (open): `BackoffConfig` defaults — the committed policy is 1s
initial, 60s max, 2x multiplier. OQ-53 tracks whether operational
experience from the alkapi deployment warrants a change before the first
release. The `BackoffConfig` struct shape is committed; the defaults are
the starting point.
- **OQ-54** (open): Inbound worker `on_worker_connected` hook placement —
the committed design is the explicit approach: the assembly layer calls
`hub.on_worker_connected()` after `CallAdapter::handle` accepts the
connection. A `HubCallAdapter` wrapper that calls the hook automatically
is additive if needed. See OQ-54.
## References
- [client-and-adapters.md](../call/client-and-adapters.md) — `CallClient`,
`from_call`, `OperationAdapter`
- [call-protocol.md](../call/call-protocol.md) — `CallAdapter`, `Dispatcher`,
`CallConnection`
- [operation-registry.md](../call/operation-registry.md) — `OperationRegistry`,
`OperationRegistryBuilder`, `OperationEnv`
- ADR-029: Peer-Graph Routing Model
- ADR-067: Aggregated Peer-Environment Wiring
- ADR-068: PeerCompositeEnv::peer_operations Override
- ADR-069: from_call Is a Manual Free Function
- alkapi [hub.md](/workspace/@alkdev/alkapi/docs/architecture/hub.md) — the
first hub consumer, the concrete use case that informed this crate
- alkapi [ADR-011](/workspace/@alkdev/alkapi/docs/architecture/decisions/011-aggregated-peer-env.md) —
the downstream aggregation decision
@@ -390,15 +390,19 @@ where `AdapterError` is a crate-level enum. The *presence* of the error type
is recorded in [client-and-adapters.md](../crates/call/client-and-adapters.md);
the exact variants are the two-way-door remainder, tracked as OQ-26.
### DC-2 — from_call re-import on reconnection: default set
### DC-2 — from_call re-import on reconnection: manual free function
Assumption 4 noted re-import "happens on reconnection or is triggered
explicitly." The v1 default is **auto-re-import on connection establishment**.
The overlay is per-connection (Layer 2, ADR-024), so re-import is naturally
scoped; a stale overlay dies with the connection. Explicit re-import via a
future `CallConnection::refresh()` is additive. Two-way door; recorded in
explicitly." The decision is **manual**: `from_call` is a free function; the
assembly layer calls it after `connect()`. The overlay is per-connection
(Layer 2, ADR-024), so re-import on reconnect is naturally scoped; a stale
overlay dies with the connection. A `CallConnection::refresh()` method for
mid-connection re-discovery is a genuine feature addition — non-breaking,
additive — if a deployment needs manual re-discovery without
drop-and-reconnect. Two-way door; recorded in
[client-and-adapters.md](../crates/call/client-and-adapters.md); tracked as
OQ-27.
OQ-27. See [ADR-069](069-from-call-manual-free-function.md) for the full
rationale.
### DC-3 — from_call namespace collision: default set
@@ -0,0 +1,208 @@
# ADR-067: Aggregated Peer-Environment Wiring for Hub Deployments
## Status
Proposed
## Context
`Dispatcher::compose_root_env` (`protocol/dispatch.rs:134-154`) constructs a
**fresh** `PeerCompositeEnv` per call and attaches **only the current call's
own connection** as a peer overlay. It does not aggregate the hub's other live
worker connections into that call's environment.
Consequence: on a hub with N connected workers, a handler composing
`env.invoke_peer(&PeerRef::Specific("dev1"), "docker", "container/exec",
input, &ctx, policy)` from a call that arrived on an HTTP connection (or on a
*different* worker's connection) will **not** find dev1's overlay — the routing
falls through to the curated base and returns `NOT_FOUND`.
The `PeerCompositeEnv` *type* and the `invoke_peer` routing logic are built for
multi-peer aggregation (`attach_peer`/`detach_peer` with insertion-order
preservation, `PeerRef::Specific`/`Any` routing — `registry/env.rs:155-301`).
The per-call `compose_root_env` does not use that capability. The ADR-029
*model* is committed; the implementation is incomplete for the head→N-workers
case.
This is the single highest-impact gap a first hub consumer (alkapi) surfaces.
A hub is *defined* by composing ops across its connected workers. Without an
aggregated env shared across all calls, the hub pattern does not work: a call
arriving on one transport cannot reach a worker connected on another.
The alkapi project identified this as OQ-08 and committed to the aggregation
decision in their ADR-011. The decision to aggregate is made; the question is
where the wiring lives — alknet (reusable by any hub) or a hub-side wrapper
(alkapi-only). This ADR resolves that question: the wiring lives in alknet.
## Decision
### 1. `Dispatcher` gains a `with_aggregated_env` builder method
A new optional field on `Dispatcher` holds a shared aggregated
`PeerCompositeEnv`:
```rust
pub struct Dispatcher {
pub registry: Arc<OperationRegistry>,
pub identity_provider: Arc<dyn IdentityProvider>,
pub session_source: Option<Arc<dyn SessionOverlaySource + Send + Sync>>,
pub ownership_provider: Option<Arc<dyn OwnershipProvider>>,
pub aggregated_env: Option<Arc<std::sync::RwLock<PeerCompositeEnv>>>,
pub default_timeout: Duration,
}
impl Dispatcher {
pub fn with_aggregated_env(
mut self,
env: Arc<std::sync::RwLock<PeerCompositeEnv>>,
) -> Self {
self.aggregated_env = Some(env);
self
}
}
```
The builder method mirrors `with_session_source` and `with_ownership_provider`
— an optional hook the assembly layer wires at construction time. A deployment
that does not set an aggregated env gets today's `compose_root_env` behavior
unchanged.
### 2. `CallAdapter` gains a matching `with_aggregated_env` builder method
`CallAdapter` delegates to `Dispatcher`:
```rust
impl CallAdapter {
pub fn with_aggregated_env(
mut self,
env: Arc<std::sync::RwLock<PeerCompositeEnv>>,
) -> Self {
self.dispatcher = self.dispatcher.with_aggregated_env(env);
self
}
}
```
### 3. `compose_root_env` reads the aggregated env when set
When `aggregated_env` is `Some`, `compose_root_env` reads the shared env,
attaches the current connection's overlay as an override for the current call
only, and returns the result. When `None`, the existing per-call behavior is
preserved:
```rust
pub fn compose_root_env(
&self,
connection: &CallConnection,
context: &OperationContext,
) -> Arc<dyn OperationEnv + Send + Sync> {
let base: Arc<dyn OperationEnv + Send + Sync> =
Arc::new(LocalOperationEnv::new(Arc::clone(&self.registry)));
let session = self
.session_source
.as_ref()
.and_then(|s| s.overlay_for(context));
if let Some(aggregated) = &self.aggregated_env {
// Clone the shared aggregated env (cheap — all fields are Arc).
let mut env = aggregated.read().unwrap().clone();
// Attach the current connection's overlay as an override for this
// call only. The current connection's overlay is the authoritative
// view of *that* peer; the aggregated env is the authoritative view
// of *all other* peers. This avoids a race where the aggregated env
// has not yet picked up a new op the current peer just registered.
if let Some(peer_id) = connection.identity().map(|identity| identity.id.clone()) {
env.attach_peer(peer_id, connection.overlay_env());
}
Arc::new(env)
} else {
let mut env = PeerCompositeEnv::new(base);
if let Some(session) = session {
env = env.with_session(session);
}
if let Some(peer_id) = connection.identity().map(|identity| identity.id.clone()) {
env.attach_peer(peer_id, connection.overlay_env());
}
Arc::new(env)
}
}
```
The clone of the aggregated env is cheap: `PeerCompositeEnv`'s fields are all
`Arc` (the `HashMap` values are `Arc<dyn OperationEnv>`, the `Vec` is
`Vec<PeerId>` which is a `String` clone). The `RwLock::read()` is held only
for the clone, not for the duration of the call.
### 4. The hub owns the aggregated env lifecycle
The hub (assembly layer) constructs the aggregated env once at startup:
```rust
let base: Arc<dyn OperationEnv + Send + Sync> =
Arc::new(LocalOperationEnv::new(Arc::clone(&registry)));
let aggregated = Arc::new(RwLock::new(PeerCompositeEnv::new(base)));
let adapter = CallAdapter::new(registry, identity_provider)
.with_aggregated_env(Arc::clone(&aggregated));
```
The hub calls `aggregated.write().unwrap().attach_peer(peer_id, overlay)` on
every worker connection-establish (after `from_call` populates the overlay)
and `detach_peer(&peer_id)` on every disconnect. The `RwLock` write is held
only for the `HashMap` insert/remove — connection-rate, not call-rate.
### 5. `PeerCompositeEnv` gains `Clone`
`PeerCompositeEnv` is made `Clone` (all fields are `Arc` or `Clone` already).
This is a one-line derive addition.
## Consequences
**Positive:**
- The hub pattern works. A call arriving on any transport can reach any
connected worker's ops via `PeerRef::Specific` or `PeerRef::Any`.
- The existing single-connection behavior is preserved. A deployment that does
not set an aggregated env gets today's `compose_root_env` unchanged.
- The hook is additive — a new optional field, a new builder method, a branch
in `compose_root_env`. No existing code path changes.
- The capability is reusable by any future hub, not just alkapi. The alkapi
project's ADR-011 fallback (hub-side wrapper) is no longer needed.
**Negative:**
- A `RwLock<PeerCompositeEnv>` on the read hot path of every dispatch. The
lock is held only for a clone (all `Arc` fields — cheap). An `ArcSwap`
copy-on-write variant could avoid the lock on reads entirely, at the cost
of a clone on `attach_peer`/`detach_peer` (infrequent). The `RwLock` is the
simpler starting point; `ArcSwap` is an additive optimization.
- `PeerCompositeEnv` gains `Clone`. The derive is mechanical; all fields are
already `Clone`.
- The hub must manage the aggregated env lifecycle (`attach_peer`/`detach_peer`
on connection events). This is assembly-layer code, not alknet-call code.
The hooks exist; the hub wires them.
## Assumptions
1. **`PeerCompositeEnv` clone is cheap.** All fields are `Arc` or `Clone` of
small types (`String`, `Vec<String>`). The clone does not copy the
operation registries or the connection overlays — it copies `Arc` pointers.
2. **The current connection's overlay is authoritative for that peer.** A call
arriving on worker-a uses worker-a's live overlay as the view of worker-a
(not the aggregated env's possibly-stale snapshot), and the aggregated env
for all other peers. This avoids a race where the aggregated env has not
yet picked up a new op worker-a just registered.
3. **The `RwLock` is not a contention point.** Reads (clones) are call-rate
but the lock is held only for the clone duration (microseconds). Writes
(`attach_peer`/`detach_peer`) are connection-rate (seconds to minutes). If
profiling shows contention, `ArcSwap` is the additive optimization.
## References
- ADR-029: Peer-Graph Routing Model (the model this wiring completes)
- ADR-030: PeerEntry and Identity.id Decoupling (the `PeerId` source)
- ADR-068: PeerCompositeEnv::peer_operations Override (sibling hub-wiring decision)
- ADR-069: from_call Is a Manual Free Function (sibling hub-wiring decision)
- alkapi ADR-011: Aggregated Peer Environment (the downstream commitment)
- alkapi OQ-08: alknet aggregated peer-env wiring (the blocking question)
- `crates/alknet-call/src/protocol/dispatch.rs:134-154` — current
`compose_root_env`
- `crates/alknet-call/src/registry/env.rs:155-301` — `PeerCompositeEnv`
@@ -0,0 +1,154 @@
# ADR-068: PeerCompositeEnv::peer_operations Override
## Status
Proposed
## Context
`OperationEnv::peer_operations` (defined in `registry/env.rs:63-65`) has a
default implementation returning `Vec::new()`. `PeerCompositeEnv` overrides
`invoke_with_policy`, `contains`, `invoke_peer`, `peer_contains`, and
`peer_ids` — but does **not** override `peer_operations`. This means
`peer_operations` on a `PeerCompositeEnv` always returns an empty `Vec`.
The `services/list-peers` handler (`registry/discovery.rs:245-296`) calls
`ctx.env.peer_operations(&peer_id)` to discover what operations each peer
serves. Since `PeerCompositeEnv` does not override this, non-local peers
always show empty operation lists in the `list-peers` response. The
`peer_ids()` method correctly returns the peer IDs, but the operations for
each peer are always empty.
This is a pure gap — the `services/list-peers` handler is specced to enumerate
each peer's operations (ADR-029 §6), and the `PeerCompositeEnv` type has all
the data needed to implement it (each peer's `OverlayOperationEnv` holds a
`HashMap<String, HandlerRegistration>`). The override is one method collecting
each peer overlay's registered op names.
The alkapi project identified this as gap G.6: a hub consumer calling
`services/list-peers` gets `peers: [{peer_id: "dev1", operations: []}]` until
this is fixed.
## Decision
`PeerCompositeEnv` overrides `peer_operations` to collect the operation names
from each peer's connection overlay:
```rust
fn peer_operations(&self, peer: &PeerId) -> Vec<String> {
match self.connections.get(peer) {
Some(overlay) => {
// The overlay is an OverlayOperationEnv wrapping a
// HashMap<String, HandlerRegistration>. We need the op names.
// Rather than adding a method to OperationEnv (which would
// require every impl to add it), we use the existing `contains`
// method — but that requires knowing the name to check.
//
// The correct approach: iterate the overlay's known names.
// OverlayOperationEnv already has the data (the HashMap keys).
// We add a `list_operation_names(&self) -> Vec<String>` method
// to OperationEnv with a default returning Vec::new(), and
// OverlayOperationEnv overrides it to return the keys.
overlay.list_operation_names()
}
None => Vec::new(),
}
}
```
### 1. `OperationEnv` gains `list_operation_names` with a default impl
```rust
fn list_operation_names(&self) -> Vec<String> {
Vec::new()
}
```
The default returns empty — existing impls (`LocalOperationEnv`, test-only
envs) don't need to change. Only `OverlayOperationEnv` overrides it.
### 2. `OverlayOperationEnv` overrides `list_operation_names`
```rust
impl OperationEnv for OverlayOperationEnv {
fn list_operation_names(&self) -> Vec<String> {
self.overlay.read().keys().cloned().collect()
}
// ... existing impl unchanged
}
```
### 3. `PeerCompositeEnv::peer_operations` uses `list_operation_names`
The override delegates to each peer's overlay:
```rust
fn peer_operations(&self, peer: &PeerId) -> Vec<String> {
self.connections
.get(peer)
.map(|overlay| overlay.list_operation_names())
.unwrap_or_default()
}
```
### Why a new trait method instead of a different approach
Alternatives considered:
- **Add `fn operations(&self) -> Vec<String>` to `OperationEnv`**: Same
concept, different name. `list_operation_names` is chosen to match the
existing `list_operations` naming on `OperationRegistry`.
- **Make `peer_operations` on `PeerCompositeEnv` reach into
`OverlayOperationEnv`'s internals**: Requires `OverlayOperationEnv` to
expose its `HashMap` or a method. The trait method is cleaner — it keeps
the abstraction boundary intact.
- **Have `services/list-peers` iterate `ctx.env.peer_ids()` and call
`contains` for every known op name**: Requires knowing all possible op
names (from the registry), which is a cross-layer coupling. The trait
method keeps the data where it lives.
The trait method is the smallest surface change: one new method with a
default impl, one override on `OverlayOperationEnv`, one override on
`PeerCompositeEnv`. No existing code changes.
## Consequences
**Positive:**
- `services/list-peers` returns actual operation lists for each peer. A hub
consumer calling `services/list-peers` gets `peers: [{peer_id: "dev1",
operations: [{name: "docker/container/exec", ...}, ...]}]` — the specced
behavior.
- The fix is small: one trait method, two overrides. No existing code paths
change.
- The `list_operation_names` method is generally useful — any future code
that needs to enumerate an env's operations can use it.
**Negative:**
- `OperationEnv` gains a method. The default impl preserves back-compat for
all existing implementors. Only `OverlayOperationEnv` and
`PeerCompositeEnv` override it.
- The `OverlayOperationEnv` override holds the `RwLock` read for the
duration of the `keys().cloned().collect()`. This is a `Vec<String>`
allocation — cheap for typical peer operation counts (tens, not thousands).
## Assumptions
1. **`OverlayOperationEnv`'s `RwLock<HashMap<String, HandlerRegistration>>`
read is cheap.** The lock is held only for the `keys()` iteration and
`collect()`. Typical peer operation counts are small (tens of ops).
2. **`list_operation_names` is the right name.** It matches the existing
`list_operations` naming on `OperationRegistry` and avoids confusion with
`peer_operations` (which takes a `PeerId` parameter).
## References
- ADR-029 §6: `services/list-peers` opt-in peer-attributed re-export listing
- ADR-067: Aggregated Peer-Environment Wiring (sibling hub-wiring decision)
- ADR-069: from_call Is a Manual Free Function (sibling hub-wiring decision)
- `crates/alknet-call/src/registry/env.rs:63-65` — default `peer_operations`
- `crates/alknet-call/src/registry/env.rs:155-301` — `PeerCompositeEnv`
- `crates/alknet-call/src/protocol/connection.rs:305-397` —
`OverlayOperationEnv`
- `crates/alknet-call/src/registry/discovery.rs:245-296` —
`services_list_peers_handler`
- alkapi gap G.6: `PeerCompositeEnv::peer_operations` unimplemented
@@ -0,0 +1,122 @@
# ADR-069: from_call Is a Manual Free Function, Not Auto-Wired
## Status
Proposed
## Context
OQ-27 resolved (2026-06-27): "The decision is **auto-re-import on connection
establishment**. The overlay is per-connection (Layer 2, ADR-024), so a stale
overlay dies with the connection; re-import on reconnect is naturally scoped to
the new connection."
The spec in `client-and-adapters.md` §"from_call" (line 358) states: "This is
the v1 default; explicit re-import via a future `CallConnection::refresh()` is
additive."
The implementation does not match. `from_call` is a standalone free function
(`client/from_call.rs:80`). `CallClient::connect()` does not call it. The
assembly layer must call `from_call()` + `register_imported_all()` explicitly
after every `connect()`. There is no `CallConnection::refresh()`.
The "v1 default" language is hedging — it makes a committed-but-not-implemented
feature sound like a deliberate phase. The spec says "auto-re-import on
connection establishment" but the code says "the assembly layer calls
`from_call` immediately after `connect()`" (the doc comment on `from_call`,
line 76). These are different things: auto-wiring means `connect()` calls
`from_call()` internally; manual means the caller does it.
The alkapi project identified this as gap G.4: the hedging language in the
spec, and the question of whether `from_call` should be auto-wired into
`connect()`.
## Decision
**`from_call` is a manual free function. The assembly layer calls it after
`connect()`. It is not auto-wired into `CallClient::connect()`.**
### Why manual is correct
1. **The hub controls discovery timing.** A hub may want to verify the
connection, resolve the peer's identity, check authorization, and *then*
discover operations. Auto-wiring `from_call` into `connect()` would run
discovery before the assembly layer has a chance to inspect the connection.
2. **Discovery is not always wanted.** A pure-client connection to a public
X.509 endpoint (ADR-034) has no `PeerEntry` and no `PeerId` — the remote
is not in the peer graph. Auto-discovering ops on such a connection would
register them in a connection overlay that has no peer key, making them
unreachable via `PeerRef`. The assembly layer decides whether to run
`from_call` based on whether the remote is a known peer.
3. **The `from_call` function is already the right API.** It takes a
`&CallConnection` and a `FromCallConfig`, returns
`Result<Vec<HandlerRegistration>, AdapterError>`, and the caller registers
the bundles. This is a clean separation: connect, discover, register. Each
step is independently testable and independently controllable.
4. **Auto-wiring would require `from_call` to know about the registry.**
`CallClient` holds an `Arc<OperationRegistry>`, but `from_call` produces
`HandlerRegistration` bundles that the caller registers — the caller
decides *where* to register them (the connection's overlay, a session
overlay, or not at all). Auto-wiring would hardcode the registration
target.
### What changes in the spec
The "v1 default" language in `client-and-adapters.md` and ADR-017 is replaced
with an honest statement: `from_call` is a free function; the assembly layer
calls it after `connect()`; there is no `CallConnection::refresh()` for
mid-connection re-discovery. A `CallConnection::refresh()` method is a
genuine feature addition — non-breaking, additive — if a deployment needs
manual re-discovery without drop-and-reconnect.
OQ-27 is updated: the resolution changes from "auto-re-import on connection
establishment" to "manual — the assembly layer calls `from_call` after
`connect()`." The door type remains two-way (auto-wiring is additive).
### What does NOT change
- The `from_call` function signature, behavior, and tests are unchanged.
- `CallClient::connect()` is unchanged.
- The re-import-on-reconnect pattern is unchanged: the assembly layer's
supervision loop calls `from_call` after each `connect()`. The overlay is
per-connection, so a stale overlay dies with the connection; re-import on
reconnect is naturally scoped. This is the correct behavior — it just
isn't automatic.
## Consequences
**Positive:**
- The spec matches the implementation. No hedging language.
- The assembly layer has full control over discovery timing and registration
target.
- The separation of concerns (connect / discover / register) is clean and
testable.
- No code changes needed — this is a spec correction, not an implementation
change.
**Negative:**
- The assembly layer must remember to call `from_call` after `connect()`.
This is a documentation concern, not a correctness concern — forgetting to
call `from_call` means the peer's ops are not imported, which is immediately
visible (calls to those ops return `NOT_FOUND`).
- The "auto-re-import on connection establishment" resolution of OQ-27 was
aspirational and is now corrected. The resolution was written before the
implementation existed; the implementation made the right call (manual),
and the spec is catching up.
## References
- ADR-017 §3: `from_call` adapter specification
- ADR-017 Amendments (DC-2): the amended `from_call` re-import resolution
(manual free function)
- ADR-067: Aggregated Peer-Environment Wiring (sibling hub-wiring decision)
- ADR-068: PeerCompositeEnv::peer_operations Override (sibling hub-wiring decision)
- OQ-27: from_call re-import trigger (amended 2026-07-09)
- `client-and-adapters.md` §"from_call" (updated)
- `crates/alknet-call/src/client/from_call.rs:80` — `from_call` free function
- `crates/alknet-call/src/client/call_client.rs:142-168` — `connect()` does
not call `from_call`
- alkapi gap G.4: `from_call` wiring + "v1 default" hedging cleanup
+8
View File
@@ -158,6 +158,14 @@ Door type is separate from whether a decision is made. A two-way door is a decis
| [OQ-50](questions/050-docker-system-events-subscription.md) | Docker System Events Subscription | deferred(scope) | two | low |
| [OQ-51](questions/051-container-create-options-surface.md) | Container Create Options Surface | deferred(scope) | two | med |
### alknet-hub
| OQ | Title | Status | Door | Pri |
|----|-------|--------|------|-----|
| [OQ-52](questions/052-callconnection-wait-for-close.md) | CallConnection::wait_for_close() for supervision loop | open | two | med |
| [OQ-53](questions/053-backoff-config-defaults.md) | BackoffConfig default policy | open | two | low |
| [OQ-54](questions/054-inbound-worker-hook-placement.md) | Inbound worker on_worker_connected hook placement | open | two | low |
## Deferred / Blocked
The safe-exit visibility surface. These questions are parked because the
@@ -1,15 +1,20 @@
# OQ-27: from_call Re-Import Trigger
- **Origin**: [client-and-adapters.md](crates/call/client-and-adapters.md), ADR-017 Assumption 4
- **Status**: **resolved** (2026-06-27)
- **Status**: **resolved** (2026-07-09 — amended from 2026-06-27 resolution)
- **Door type**: Two-way
- **Priority**: low
- **Resolution**: The decision is **auto-re-import on connection
establishment**. The overlay is per-connection (Layer 2, ADR-024), so a
stale overlay dies with the connection; re-import on reconnect is
naturally scoped to the new connection. This is the right default for the
runner pattern (a worker reconnects → the hub re-discovers the worker's
ops automatically). An explicit `CallConnection::refresh()` method is a
genuine feature addition — non-breaking, additive — if a deployment
needs manual control.
- **Resolution**: `from_call` is a **manual free function**; the assembly layer
calls it after `connect()`. The overlay is per-connection (Layer 2, ADR-024),
so a stale overlay dies with the connection; re-import on reconnect is
naturally scoped to the new connection. A `CallConnection::refresh()` method
for mid-connection re-discovery is a genuine feature addition —
non-breaking, additive — if a deployment needs manual re-discovery without
drop-and-reconnect. See [ADR-069](../decisions/069-from-call-manual-free-function.md)
for the full rationale.
The original 2026-06-27 resolution ("auto-re-import on connection
establishment") was aspirational — written before the implementation
existed. The implementation made the right call (manual free function);
this amendment aligns the spec with the implementation.
- **Cross-references**: ADR-017, ADR-024, [client-and-adapters.md](crates/call/client-and-adapters.md)
@@ -0,0 +1,26 @@
# OQ-52: CallConnection::wait_for_close() for Supervision Loop
- **Origin**: [crates/hub/README.md](crates/hub/README.md)
- **Status**: open
- **Door type**: Two-way
- **Priority**: medium
- **Resolution**: Not yet decided.
The hub's worker supervision loop needs a way to await connection close so it
can call `detach_peer` and retry. Today `CallConnection` exposes `connection()`
(the underlying `Connection`) but not a "wait for run_loop exit" future.
Options:
- **(a)** Add a `closed()` method to `CallConnection` that returns a
`Future<Output = ()>` resolving when `run_loop` exits. The dispatcher
signals a `tokio::sync::Notify` on exit; `closed()` awaits it.
- **(b)** Use a `tokio::sync::oneshot` channel created by the caller and
passed to the dispatcher, signaled on `run_loop` exit.
- **(c)** Poll `connection().accept_bi()` in a loop until it returns
`ConnectionClosed` — works but is polling, not event-driven.
Option (a) is the cleanest: a method on `CallConnection` that any caller
(not just the hub) can use to await connection close. It is a small additive
change to `alknet-call`.
- **Cross-references**: ADR-067, [crates/hub/README.md](crates/hub/README.md)
@@ -0,0 +1,19 @@
# OQ-53: BackoffConfig Default Policy
- **Origin**: [crates/hub/README.md](crates/hub/README.md)
- **Status**: open
- **Door type**: Two-way
- **Priority**: low
- **Resolution**: Not yet decided.
The `BackoffConfig` struct provides configurable backoff for worker
reconnection. The default policy (1s initial, 60s max, 2x multiplier) is a
starting point. Whether this is the right policy for production deployments
is an operational question, not an architectural one.
The `BackoffConfig` struct shape is committed; the defaults are a starting
point that can be changed without breaking the API. The question is whether
to change the defaults before the first release, based on operational
experience from the alkapi deployment.
- **Cross-references**: [crates/hub/README.md](crates/hub/README.md)
@@ -0,0 +1,26 @@
# OQ-54: Inbound Worker on_worker_connected Hook Placement
- **Origin**: [crates/hub/README.md](crates/hub/README.md)
- **Status**: open
- **Door type**: Two-way
- **Priority**: low
- **Resolution**: Not yet decided.
When a worker connects inbound (via `CallAdapter::handle`), the hub needs to
call `from_call`, register the bundles, and `attach_peer`. Two options:
- **(a) Explicit**: The assembly layer calls `hub.on_worker_connected()` after
`CallAdapter::handle` accepts the connection. The hub provides the method;
the assembly layer wires it.
- **(b) Wrapper**: The hub provides a `HubCallAdapter` wrapper that calls the
hook automatically inside `handle()`. The assembly layer uses the wrapper
instead of the raw `CallAdapter`.
Option (a) is simpler and more flexible — the assembly layer controls when
the hook fires and can add its own logic (authorization checks, logging)
before discovery. Option (b) is more convenient but hides the hook from the
assembly layer.
The explicit approach is the committed interim. A wrapper is additive.
- **Cross-references**: ADR-067, ADR-068, [crates/hub/README.md](crates/hub/README.md)