diff --git a/docs/architecture/README.md b/docs/architecture/README.md index 1198768..30b1513 100644 --- a/docs/architecture/README.md +++ b/docs/architecture/README.md @@ -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 diff --git a/docs/architecture/crates/call/README.md b/docs/architecture/crates/call/README.md index 58c2722..40b641e 100644 --- a/docs/architecture/crates/call/README.md +++ b/docs/architecture/crates/call/README.md @@ -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 | diff --git a/docs/architecture/crates/call/call-protocol.md b/docs/architecture/crates/call/call-protocol.md index 3ce5397..3ca0537 100644 --- a/docs/architecture/crates/call/call-protocol.md +++ b/docs/architecture/crates/call/call-protocol.md @@ -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). diff --git a/docs/architecture/crates/call/client-and-adapters.md b/docs/architecture/crates/call/client-and-adapters.md index 09fb495..34c4f91 100644 --- a/docs/architecture/crates/call/client-and-adapters.md +++ b/docs/architecture/crates/call/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>`, 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. diff --git a/docs/architecture/crates/call/operation-registry.md b/docs/architecture/crates/call/operation-registry.md index 3d1ccad..c07756a 100644 --- a/docs/architecture/crates/call/operation-registry.md +++ b/docs/architecture/crates/call/operation-registry.md @@ -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` (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`). 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 diff --git a/docs/architecture/crates/hub/README.md b/docs/architecture/crates/hub/README.md new file mode 100644 index 0000000..e86f699 --- /dev/null +++ b/docs/architecture/crates/hub/README.md @@ -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, + aggregated_env: Arc>, + dispatcher: Dispatcher, + identity_provider: Arc, +} +``` + +Construction: + +```rust +impl Hub { + pub fn new( + registry: Arc, + identity_provider: Arc, + ) -> Self { + let base: Arc = + Arc::new(LocalOperationEnv::new(Arc::clone(®istry))); + let aggregated_env = Arc::new(RwLock::new(PeerCompositeEnv::new(base))); + let dispatcher = Dispatcher::new(Arc::clone(®istry), 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> { + &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 { + 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, + 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(®istry), + 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(®istry), 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>` +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 diff --git a/docs/architecture/decisions/017-call-protocol-client-and-adapter-contract.md b/docs/architecture/decisions/017-call-protocol-client-and-adapter-contract.md index 4e860e1..44c8850 100644 --- a/docs/architecture/decisions/017-call-protocol-client-and-adapter-contract.md +++ b/docs/architecture/decisions/017-call-protocol-client-and-adapter-contract.md @@ -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 diff --git a/docs/architecture/decisions/067-aggregated-peer-env-wiring.md b/docs/architecture/decisions/067-aggregated-peer-env-wiring.md new file mode 100644 index 0000000..cb2886c --- /dev/null +++ b/docs/architecture/decisions/067-aggregated-peer-env-wiring.md @@ -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, + pub identity_provider: Arc, + pub session_source: Option>, + pub ownership_provider: Option>, + pub aggregated_env: Option>>, + pub default_timeout: Duration, +} + +impl Dispatcher { + pub fn with_aggregated_env( + mut self, + env: Arc>, + ) -> 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>, + ) -> 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 { + let base: Arc = + 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`, the `Vec` is +`Vec` 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 = + Arc::new(LocalOperationEnv::new(Arc::clone(®istry))); +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` 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`). 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` diff --git a/docs/architecture/decisions/068-peer-composite-env-peer-operations.md b/docs/architecture/decisions/068-peer-composite-env-peer-operations.md new file mode 100644 index 0000000..6a844f6 --- /dev/null +++ b/docs/architecture/decisions/068-peer-composite-env-peer-operations.md @@ -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`). 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 { + match self.connections.get(peer) { + Some(overlay) => { + // The overlay is an OverlayOperationEnv wrapping a + // HashMap. 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` 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 { + 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 { + 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 { + 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` 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` + allocation — cheap for typical peer operation counts (tens, not thousands). + +## Assumptions + +1. **`OverlayOperationEnv`'s `RwLock>` + 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 diff --git a/docs/architecture/decisions/069-from-call-manual-free-function.md b/docs/architecture/decisions/069-from-call-manual-free-function.md new file mode 100644 index 0000000..0bfdd8a --- /dev/null +++ b/docs/architecture/decisions/069-from-call-manual-free-function.md @@ -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, 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`, 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 diff --git a/docs/architecture/open-questions.md b/docs/architecture/open-questions.md index 78051a4..d183180 100644 --- a/docs/architecture/open-questions.md +++ b/docs/architecture/open-questions.md @@ -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 diff --git a/docs/architecture/questions/027-from-call-re-import-trigger.md b/docs/architecture/questions/027-from-call-re-import-trigger.md index 69e1e47..6593d02 100644 --- a/docs/architecture/questions/027-from-call-re-import-trigger.md +++ b/docs/architecture/questions/027-from-call-re-import-trigger.md @@ -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) diff --git a/docs/architecture/questions/052-callconnection-wait-for-close.md b/docs/architecture/questions/052-callconnection-wait-for-close.md new file mode 100644 index 0000000..55a2404 --- /dev/null +++ b/docs/architecture/questions/052-callconnection-wait-for-close.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` 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) diff --git a/docs/architecture/questions/053-backoff-config-defaults.md b/docs/architecture/questions/053-backoff-config-defaults.md new file mode 100644 index 0000000..83e1a19 --- /dev/null +++ b/docs/architecture/questions/053-backoff-config-defaults.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) diff --git a/docs/architecture/questions/054-inbound-worker-hook-placement.md b/docs/architecture/questions/054-inbound-worker-hook-placement.md new file mode 100644 index 0000000..33646c7 --- /dev/null +++ b/docs/architecture/questions/054-inbound-worker-hook-placement.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)