From 502488cc72c0bab109a3d0f2cec50e23c59219ee Mon Sep 17 00:00:00 2001 From: "glm-5.2" Date: Wed, 12 Aug 2026 13:29:58 +0000 Subject: [PATCH] =?UTF-8?q?fix:=20Unit=201=20=E2=80=94=20make=20publish()?= =?UTF-8?q?=20work=20end-to-end=20(P-01,=20P-02,=20P-05,=20P-08,=20P-09,?= =?UTF-8?q?=20P-12=20#1)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fixes the Pub operation's client→responder path so a publish() actually works on any transport. Previously the wire was broken in three compounding ways and the unit tests couldn't see it. P-01 [critical] — publish chunks were written to a *fresh* open_bi stream (stream B) while the responder read chunks from the stream that carried call.requested (stream A). Stream B's frames were silently discarded by the dispatch loop's "ignoring non-requested event" branch. On single-stream transports (TCP+TLS, SSH) the second open_bi returns StreamClosed outright, so every publish failed. Fix: pump call.requested + call.published + call.completed on the *same* write half via the new pump_publish_to_wire free function; the read half is pumped concurrently for the single call.responded. P-02 [major] — publish() / publish_with_payload() now take Pin + Send>> per ADR-046 §8, not Vec. The Vec shape buffered the entire publish in memory and made the ADR's flagship use cases (telemetry ingest, live upload, drag-drop file streaming) impossible. The wire format is unchanged; this corrects API drift (no deployments exist). P-05 [major] — publish now registers with timeout: None (like subscribe), not DEFAULT_CALL_TIMEOUT (30s). A publish whose stream took >30s wall-clock got a spurious client-side TIMEOUT while the responder kept consuming. PendingEntry::Call.timeout is now Option so the sweeper never evicts unbounded calls; all register_call callers updated (Some(...) for call(), None for publish()). P-08 [major] — the dispatch Pub branch re-implemented invoke_sink's not-found / visibility / ACL / handler-kind checks inline; invoke_sink was only ever called from its own tests. Any future fix in invoke_sink (e.g. the missing publish_schema validation, P-03) wouldn't reach the wire path. Fix: extract OperationRegistry::resolve_sink_handler (pub(crate)) as the single source of truth for the sink dispatch checks; both invoke_sink and Dispatcher::dispatch (Pub branch) call it. The two paths can no longer diverge. P-09 [major, side-effect] — make_sink_forwarding_handler in from_call.rs was store-and-forward (collect into Vec) and silently discarded Err items (filter_map(|item| item.ok())), contradicting the SinkHandler contract that an Err terminates the stream. Now that publish_with_payload takes a Stream, the forwarding handler passes the stream through directly (streamed, not buffered) and uses take_while(Ok) + filter_map to terminate on Err. P-09 was listed in Unit 6 but depends on P-01/P-02, so it falls out naturally here. P-12 #1 [acceptance gate] — adds the end-to-end publish test that the review identifies as the gate for this unit: publish_end_to_end_delivers_chunks_and_returns_response wires CallConnection::publish ↔ Dispatcher::run_loop over a real tokio::io::duplex pair (new duplex_connection_pair test helper + SingleStreamSource BidiStreamSource impl), publishes 3 chunks, and asserts the responder saw all 3 and returned the right result. A second test covers the unknown-op → NOT_FOUND path. These tests would have caught P-01 immediately. Verification: - cargo test --lib → 434 passed, 0 failed (was 432; +2 e2e tests) - cargo clippy --all-targets -- -D warnings → clean - cargo fmt --check → clean - cargo doc --no-deps → 4 warnings (pre-existing C-09, unchanged) --- src/client/from_call.rs | 22 ++- src/protocol/abort.rs | 2 +- src/protocol/adapter.rs | 6 +- src/protocol/connection.rs | 291 +++++++++++++++++++++++++++-------- src/protocol/dispatch.rs | 76 ++------- src/protocol/mod.rs | 2 +- src/protocol/pending.rs | 34 ++-- src/protocol/test_support.rs | 67 +++++++- src/registry/registration.rs | 88 ++++++----- 9 files changed, 400 insertions(+), 188 deletions(-) diff --git a/src/client/from_call.rs b/src/client/from_call.rs index 676872f..7a3ca73 100644 --- a/src/client/from_call.rs +++ b/src/client/from_call.rs @@ -11,8 +11,10 @@ //! the spec and the v1 defaults (auto-on-reconnect, error-on-collision). use std::collections::HashSet; +use std::pin::Pin; use std::sync::Arc; +use futures::stream::Stream; use serde_json::{json, Value}; use crate::client::AdapterError; @@ -450,7 +452,14 @@ fn make_streaming_forwarding_handler( /// and forwards the local `PublishStream` to the remote as /// `call.published` events. The remote's single `call.responded` becomes /// the handler's `ResponseEnvelope`. No truncation — the full stream is -/// forwarded end-to-end. +/// forwarded end-to-end (streamed, not buffered: the local stream's +/// items are pumped to the remote as they arrive). +/// +/// An `Err` item in the local `PublishStream` (an initiator-side error) +/// terminates the forwarding — the stream is cut at the error, matching +/// the `SinkHandler` contract (`Err` terminates the stream). The +/// remote's `call.completed` is sent after the last `Ok` item; the +/// remote handler sees a truncated stream. /// /// `forwarded_for` is populated from `context.identity` (ADR-032 §3), /// exactly as the request/response and streaming forwarding handlers @@ -466,11 +475,12 @@ fn make_sink_forwarding_handler( let remote_name = remote_name.clone(); async move { let payload = build_forwarded_payload(&remote_name, input, &context); - let chunks: Vec = publish_stream - .filter_map(|item| async move { item.ok() }) - .collect() - .await; - connection.publish_with_payload(payload, chunks).await + let value_stream: Pin + Send>> = Box::pin( + publish_stream + .take_while(|item| futures::future::ready(item.is_ok())) + .filter_map(|item| futures::future::ready(item.ok())), + ); + connection.publish_with_payload(payload, value_stream).await } }) } diff --git a/src/protocol/abort.rs b/src/protocol/abort.rs index c446450..314fc2c 100644 --- a/src/protocol/abort.rs +++ b/src/protocol/abort.rs @@ -113,7 +113,7 @@ mod tests { fn register_call(map: &mut PendingRequestMap, id: &str, parent: Option<&str>) { map.register_call( id.to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), parent.map(|p| p.to_string()), ); } diff --git a/src/protocol/adapter.rs b/src/protocol/adapter.rs index cde7893..17b63e2 100644 --- a/src/protocol/adapter.rs +++ b/src/protocol/adapter.rs @@ -1170,12 +1170,12 @@ mod tests { let mut pending = conn.pending().lock(); pending.register_call( "parent-1".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); pending.register_call( "child-1".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), Some("parent-1".to_string()), ); } @@ -1214,7 +1214,7 @@ mod tests { let mut pending = conn.pending().lock(); pending.register_call( "unrelated-1".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); } diff --git a/src/protocol/connection.rs b/src/protocol/connection.rs index 9ba29bb..2c03243 100644 --- a/src/protocol/connection.rs +++ b/src/protocol/connection.rs @@ -12,7 +12,7 @@ use std::time::{Duration, Instant}; use crate::core::auth::Identity; use crate::core::types::Connection; -use futures::stream::Stream; +use futures::stream::{Stream, StreamExt}; use parking_lot::{Mutex, RwLock}; use serde_json::Value; use tokio::io::{AsyncRead, AsyncWrite}; @@ -142,7 +142,7 @@ impl CallConnection { let mut pending = self.pending.lock(); pending.register_call( request_id.clone(), - Instant::now() + DEFAULT_CALL_TIMEOUT, + Some(Instant::now() + DEFAULT_CALL_TIMEOUT), None, ) }; @@ -243,20 +243,24 @@ impl CallConnection { } /// Publish a stream to a remote `Pub` operation (ADR-046). Sends - /// `call.requested` (open), then one `call.published` per chunk in - /// `chunks`, then `call.completed` (stream end). Returns the + /// `call.requested` (open), then one `call.published` per item in + /// `stream`, then `call.completed` (stream end). Returns the /// responder's single `ResponseEnvelope`. + /// + /// The publish is long-lived and unbounded — no client-side deadline + /// (ADR-046 §7, "same as `Sub`"). The `PendingRequestMap` entry is + /// registered with `timeout: None`, so the sweeper never evicts it. pub async fn publish( &self, operation_id: &str, input: Value, - chunks: Vec, + stream: Pin + Send>>, ) -> ResponseEnvelope { let payload = serde_json::json!({ "operationId": operation_id, "input": input, }); - self.publish_with_payload(payload, chunks).await + self.publish_with_payload(payload, stream).await } /// Publish a stream to a remote `Pub` op with a caller-constructed @@ -265,10 +269,17 @@ impl CallConnection { /// `auth_token` for the hub forwarding path used by `from_call`'s /// sink forwarding handler. Mirrors /// [`subscribe_with_payload`](Self::subscribe_with_payload). + /// + /// All frames (`call.requested`, `call.published`, `call.completed`) + /// are written on the **same** bidirectional stream's write half — + /// the responder reads chunks from the same stream it received + /// `call.requested` on (ADR-046 §5 stream lifecycle). The read half + /// is pumped concurrently for the single `call.responded` / + /// `call.error` response. pub async fn publish_with_payload( &self, payload: Value, - chunks: Vec, + stream: Pin + Send>>, ) -> ResponseEnvelope { let request_id = generate_request_id(); @@ -282,7 +293,7 @@ impl CallConnection { } }; - let stream = match connection.open_bi().await { + let stream_a = match connection.open_bi().await { Ok(s) => s, Err(err) => { return ResponseEnvelope::error( @@ -291,39 +302,28 @@ impl CallConnection { ); } }; - let (recv, send) = tokio::io::split(stream); + let (recv, send) = tokio::io::split(stream_a); let receiver = { let mut pending = self.pending.lock(); - pending.register_call( - request_id.clone(), - Instant::now() + DEFAULT_CALL_TIMEOUT, - None, - ) + pending.register_call(request_id.clone(), None, None) }; - if let Err(err) = self.write_request(send, &request_id, payload).await { - let call_error = CallError::internal(err); - self.pending - .lock() - .handle_error(&request_id, call_error.clone()); - return ResponseEnvelope::error(request_id, call_error); - } + let pending_map = Arc::clone(&self.pending); + let read_handle = tokio::spawn(async move { + read_stream_until_closed(recv, &pending_map).await; + }); - let write_result = self.write_publish_chunks(&request_id, chunks).await; + let write_result = pump_publish_to_wire(send, &request_id, payload, stream).await; if let Err(err) = write_result { let call_error = CallError::internal(err); self.pending .lock() .handle_error(&request_id, call_error.clone()); + read_handle.abort(); return ResponseEnvelope::error(request_id, call_error); } - let pending = Arc::clone(&self.pending); - tokio::spawn(async move { - read_stream_until_closed(recv, &pending).await; - }); - match receiver.await { Ok(Ok(value)) => ResponseEnvelope::ok(request_id, value), Ok(Err(error)) => ResponseEnvelope::error(request_id, error), @@ -331,36 +331,6 @@ impl CallConnection { } } - async fn write_publish_chunks( - &self, - request_id: &str, - chunks: Vec, - ) -> Result<(), String> { - let connection = self - .connection - .as_ref() - .ok_or_else(|| "no underlying connection (overlay-only)".to_string())?; - let stream = connection - .open_bi() - .await - .map_err(|e| format!("failed to open stream: {e}"))?; - let (_recv, send) = tokio::io::split(stream); - let mut writer = FrameFramedWriter::new(send); - for chunk in chunks { - let envelope = EventEnvelope::published(request_id, chunk); - writer - .write_frame(&envelope) - .await - .map_err(|e| format!("failed to write published frame: {e}"))?; - } - let completed = EventEnvelope::completed(request_id); - writer - .write_frame(&completed) - .await - .map_err(|e| format!("failed to write completed frame: {e}"))?; - Ok(()) - } - async fn write_request( &self, send: W, @@ -398,6 +368,36 @@ impl CallConnection { } } +async fn pump_publish_to_wire( + send: W, + request_id: &str, + payload: Value, + mut stream: Pin + Send>>, +) -> Result<(), String> +where + W: AsyncWrite + Unpin, +{ + let mut writer = FrameFramedWriter::new(send); + let requested = EventEnvelope::requested(request_id, payload); + writer + .write_frame(&requested) + .await + .map_err(|e| format!("failed to write request frame: {e}"))?; + while let Some(chunk) = stream.next().await { + let envelope = EventEnvelope::published(request_id, chunk); + writer + .write_frame(&envelope) + .await + .map_err(|e| format!("failed to write published frame: {e}"))?; + } + let completed = EventEnvelope::completed(request_id); + writer + .write_frame(&completed) + .await + .map_err(|e| format!("failed to write completed frame: {e}"))?; + Ok(()) +} + async fn read_stream_until_closed(recv: R, pending: &Arc>) where R: AsyncRead + Unpin, @@ -593,7 +593,10 @@ mod tests { use super::*; use crate::core::types::Capabilities; use crate::registry::context::CompositionAuthority; - use crate::registry::registration::{make_handler, Handler, HandlerKind, OperationProvenance}; + use crate::registry::registration::{ + make_handler, Handler, HandlerKind, HandlerRegistration, OperationProvenance, + OperationRegistry, + }; use crate::registry::spec::{AccessControl, OperationSpec, OperationType, Visibility}; use std::collections::HashMap; use std::time::{Duration, Instant}; @@ -820,7 +823,7 @@ mod tests { let pending = empty_pending(); let rx = pending.lock().register_call( "req-1".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); let envelope = EventEnvelope::responded("req-1", serde_json::json!({"v": 42})); @@ -875,7 +878,7 @@ mod tests { let pending = empty_pending(); let _rx = pending.lock().register_call( "req-2".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); assert!(pending.lock().contains("req-2")); @@ -888,7 +891,7 @@ mod tests { let pending = empty_pending(); let rx = pending.lock().register_call( "req-3".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); let err = CallError::new("FILE_NOT_FOUND", "missing", false); @@ -928,7 +931,7 @@ mod tests { let pending = empty_pending(); let _rx = pending.lock().register_call( "req-4".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); let malformed = @@ -942,7 +945,7 @@ mod tests { let pending = empty_pending(); let _rx = pending.lock().register_call( "req-5".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); let unknown = EventEnvelope::new("call.mystery", "req-5", serde_json::json!({})); @@ -1053,7 +1056,7 @@ mod tests { assert!(pending.lock().is_empty()); let _rx = pending.lock().register_call( "req-overlay-1".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); assert!(pending.lock().contains("req-overlay-1")); @@ -1145,4 +1148,164 @@ mod tests { other => panic!("expected INVALID_OPERATION_TYPE, got {other:?}"), } } + + // --- end-to-end publish (P-12 #1, ADR-046 §8) ------------------------ + // + // Wires `CallConnection::publish` (client) ↔ `Dispatcher::run_loop` + // (responder) over a real `tokio::io::duplex` pair. This is the + // acceptance gate for Unit 1: it would have caught P-01 (chunks on + // the wrong stream) immediately — the responder would never see the + // chunks. + + fn pub_spec_e2e(name: &str) -> OperationSpec { + OperationSpec::new( + name, + OperationType::Pub, + Visibility::External, + serde_json::json!({}), + serde_json::json!({}), + vec![], + AccessControl::default(), + None, + ) + } + + #[tokio::test] + async fn publish_end_to_end_delivers_chunks_and_returns_response() { + use crate::core::auth::IdentityProvider; + use crate::protocol::dispatch::Dispatcher; + use crate::protocol::duplex_connection_pair; + use crate::registry::registration::{make_sink_handler, SinkHandler}; + use futures::stream::StreamExt; + + struct NoopIdProvider; + impl IdentityProvider for NoopIdProvider { + fn resolve_from_fingerprint(&self, _: &str) -> Option { + None + } + fn resolve_from_token(&self, _: &crate::core::auth::AuthToken) -> Option { + None + } + } + + let counting_sink: SinkHandler = make_sink_handler(|_input, ctx, mut stream| async move { + let mut count = 0u32; + let mut last = serde_json::Value::Null; + while let Some(item) = stream.next().await { + match item { + Ok(v) => { + count += 1; + last = v; + } + Err(_) => break, + } + } + ResponseEnvelope::ok( + ctx.request_id, + serde_json::json!({ "count": count, "last": last }), + ) + }); + + let mut registry = OperationRegistry::new(); + registry + .register(HandlerRegistration::new( + pub_spec_e2e("fs/upload"), + HandlerKind::Sink(counting_sink), + OperationProvenance::Local, + None, + None, + Capabilities::new(), + )) + .unwrap(); + let registry = Arc::new(registry); + + let (client_conn, server_conn) = duplex_connection_pair(64 * 1024); + let client = CallConnection::new(client_conn); + + let dp = Dispatcher::new(Arc::clone(®istry), Arc::new(NoopIdProvider)); + let server_call_conn = Arc::new(CallConnection::new(server_conn)); + let server_handle = tokio::spawn(async move { + dp.run_loop(server_call_conn).await; + }); + + let chunks = vec![ + serde_json::json!({"chunk": 1}), + serde_json::json!({"chunk": 2}), + serde_json::json!({"chunk": 3}), + ]; + let stream: Pin + Send>> = + Box::pin(futures::stream::iter(chunks.clone())); + let response = client + .publish("fs/upload", serde_json::json!({"path": "/x"}), stream) + .await; + + assert!( + response.result.is_ok(), + "publish should succeed, got {:?}", + response.result + ); + let out = response.result.unwrap(); + assert_eq!( + out["count"], + serde_json::json!(3), + "responder saw all 3 chunks" + ); + assert_eq!( + out["last"], + serde_json::json!({"chunk": 3}), + "last chunk matches" + ); + + drop(client); + let _ = server_handle.await; + } + + #[tokio::test] + async fn publish_end_to_end_unknown_op_returns_not_found() { + use crate::core::auth::IdentityProvider; + use crate::protocol::dispatch::Dispatcher; + use crate::protocol::duplex_connection_pair; + use crate::registry::registration::make_sink_handler; + + struct NoopIdProvider; + impl IdentityProvider for NoopIdProvider { + fn resolve_from_fingerprint(&self, _: &str) -> Option { + None + } + fn resolve_from_token(&self, _: &crate::core::auth::AuthToken) -> Option { + None + } + } + + let registry = Arc::new(OperationRegistry::new()); + + let (client_conn, server_conn) = duplex_connection_pair(64 * 1024); + let client = CallConnection::new(client_conn); + + let dp = Dispatcher::new(Arc::clone(®istry), Arc::new(NoopIdProvider)); + let server_call_conn = Arc::new(CallConnection::new(server_conn)); + let server_handle = tokio::spawn(async move { + dp.run_loop(server_call_conn).await; + }); + + let stream: Pin + Send>> = + Box::pin(futures::stream::iter(vec![ + serde_json::json!({"chunk": 1}), + serde_json::json!({"chunk": 2}), + ])); + let _unused_sink = make_sink_handler(|_, ctx, _| async move { + ResponseEnvelope::ok(ctx.request_id, serde_json::json!({})) + }); + let response = client + .publish("no/such/op", serde_json::json!({}), stream) + .await; + + match response.result { + Err(e) => assert_eq!(e.code, "NOT_FOUND"), + other => panic!("expected NOT_FOUND, got {other:?}"), + } + + drop(client); + let _ = server_handle.await; + } } diff --git a/src/protocol/dispatch.rs b/src/protocol/dispatch.rs index a9e99a3..59326d3 100644 --- a/src/protocol/dispatch.rs +++ b/src/protocol/dispatch.rs @@ -35,10 +35,8 @@ use super::wire::{ use crate::protocol::adapter::SessionOverlaySource; use crate::registry::context::{AbortPolicy, OperationContext, ScopedPeerEnv}; use crate::registry::env::{LocalOperationEnv, OperationEnv, PeerCompositeEnv}; -use crate::registry::registration::{ - extract_json_pointer, HandlerKind, OperationRegistry, PublishStream, ResponseStream, -}; -use crate::registry::spec::{AccessResult, OperationType, Visibility}; +use crate::registry::registration::{OperationRegistry, PublishStream, ResponseStream}; +use crate::registry::spec::OperationType; const DEFAULT_TIMEOUT: Duration = Duration::from_secs(30); const SWEEPER_INTERVAL: Duration = Duration::from_secs(10); @@ -308,64 +306,14 @@ impl Dispatcher { } OperationType::Pub => { context.deadline = None; - let registration = match self.registry.registration(&operation_name) { - Some(r) => r, - None => { - return DispatchResult::Once(ResponseEnvelope::not_found( - request_id.clone(), - &operation_name, - )); - } - }; - if registration.spec.visibility == Visibility::Internal && !context.internal { - return DispatchResult::Once(ResponseEnvelope::not_found( - request_id.clone(), - &operation_name, - )); - } - let acl = ®istration.spec.access_control; - let identity = if context.internal { - context - .handler_identity - .as_ref() - .and_then(|ca| ca.as_identity()) - } else { - context.identity.clone() - }; - let resource_id = registration - .spec - .resource_id_path - .as_ref() - .and_then(|path| extract_json_pointer(&input, path)); - if let AccessResult::Forbidden(message) = acl.check( - identity.as_ref(), - resource_id.as_deref(), - context.ownership.as_deref(), - ) { - return DispatchResult::Once(ResponseEnvelope::forbidden( - request_id.clone(), - message, - )); - } - let sink_handler = match ®istration.handler { - HandlerKind::Sink(h) => Arc::clone(h), - HandlerKind::Once(_) => { - return DispatchResult::Once(ResponseEnvelope::error( - request_id.clone(), - CallError::invalid_operation_type( - "invoke_sink() called on a Query/Mutation op; use invoke()", - ), - )); - } - HandlerKind::Stream(_) => { - return DispatchResult::Once(ResponseEnvelope::error( - request_id.clone(), - CallError::invalid_operation_type( - "invoke_sink() called on a Sub op; use invoke_streaming()", - ), - )); - } - }; + let sink_handler = + match self + .registry + .resolve_sink_handler(&operation_name, &input, &context) + { + Ok(h) => h, + Err(envelope) => return DispatchResult::Once(envelope), + }; let (chunk_tx, chunk_rx) = mpsc::channel::>(PUBLISH_CHANNEL_BUFFER); let publish_stream: PublishStream = Box::pin(chunk_rx); @@ -1111,12 +1059,12 @@ mod tests { let mut pending = conn.pending().lock(); pending.register_call( parent_id.clone(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); pending.register_call( child_id.clone(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), Some(parent_id.clone()), ); } diff --git a/src/protocol/mod.rs b/src/protocol/mod.rs index 0c3f3d9..cf1b8d4 100644 --- a/src/protocol/mod.rs +++ b/src/protocol/mod.rs @@ -14,4 +14,4 @@ pub mod wire; mod test_support; #[cfg(test)] -pub(crate) use test_support::sink_empty_connection; +pub(crate) use test_support::{duplex_connection_pair, sink_empty_connection}; diff --git a/src/protocol/pending.rs b/src/protocol/pending.rs index cb2e2f9..01904ea 100644 --- a/src/protocol/pending.rs +++ b/src/protocol/pending.rs @@ -15,7 +15,7 @@ pub struct PendingRequestMap { pub(crate) enum PendingEntry { Call { tx: oneshot::Sender>, - timeout: Instant, + timeout: Option, parent_request_id: Option, started: bool, }, @@ -54,10 +54,15 @@ impl PendingRequestMap { } } + /// Register a pending call (request/response or Pub). `timeout: None` + /// marks the entry as unbounded — the sweeper never evicts it. `Pub` + /// operations register with `timeout: None` (the stream may be + /// long-lived, same as `Sub` — ADR-046 §7); `call()` registers with + /// `Some(Instant::now() + DEFAULT_CALL_TIMEOUT)`. pub fn register_call( &mut self, request_id: String, - timeout: Instant, + timeout: Option, parent_request_id: Option, ) -> oneshot::Receiver> { let (tx, rx) = oneshot::channel(); @@ -175,7 +180,10 @@ impl PendingRequestMap { let mut to_remove: Vec = Vec::new(); for (id, entry) in self.pending.iter() { let expired = match entry { - PendingEntry::Call { timeout, .. } => *timeout <= now, + PendingEntry::Call { + timeout: Some(t), .. + } => *t <= now, + PendingEntry::Call { timeout: None, .. } => false, PendingEntry::Subscribe { timeout: Some(t), .. } => *t <= now, @@ -273,7 +281,7 @@ mod tests { let mut map = PendingRequestMap::new(); let rx = map.register_call( "req-1".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); @@ -333,7 +341,7 @@ mod tests { let mut map = PendingRequestMap::new(); let rx = map.register_call( "req-2".to_string(), - Instant::now() - Duration::from_millis(1), + Some(Instant::now() - Duration::from_millis(1)), None, ); @@ -388,7 +396,7 @@ mod tests { let mut map = PendingRequestMap::new(); let rx_call = map.register_call( "c-1".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); let mut rx_sub = map.register_subscribe( @@ -452,7 +460,7 @@ mod tests { let mut map = PendingRequestMap::new(); let rx = map.register_call( "req-3".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); @@ -471,7 +479,7 @@ mod tests { let mut map = PendingRequestMap::new(); let rx = map.register_call( "req-4".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); @@ -513,7 +521,7 @@ mod tests { let mut map = PendingRequestMap::new(); let rx = map.register_call( "req-stream-3".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); @@ -530,12 +538,12 @@ mod tests { let mut map = PendingRequestMap::new(); let _rx_old = map.register_call( "req-5".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); let rx_new = map.register_call( "req-5".to_string(), - Instant::now() + Duration::from_secs(30), + Some(Instant::now() + Duration::from_secs(30)), None, ); assert_eq!(map.len(), 1); @@ -553,12 +561,12 @@ mod tests { let mut map = PendingRequestMap::new(); let _rx_expired = map.register_call( "expired".to_string(), - Instant::now() - Duration::from_millis(1), + Some(Instant::now() - Duration::from_millis(1)), None, ); let _rx_alive = map.register_call( "alive".to_string(), - Instant::now() + Duration::from_secs(60), + Some(Instant::now() + Duration::from_secs(60)), None, ); diff --git a/src/protocol/test_support.rs b/src/protocol/test_support.rs index bf40262..bf8e931 100644 --- a/src/protocol/test_support.rs +++ b/src/protocol/test_support.rs @@ -6,9 +6,10 @@ use std::net::{IpAddr, Ipv4Addr, SocketAddr}; use std::pin::Pin; +use std::sync::Mutex; use std::task::{Context, Poll}; -use crate::core::types::Connection; +use crate::core::types::{BiStream, BidiStreamSource, Connection, StreamError}; use tokio::io::{AsyncRead, AsyncWrite, ReadBuf}; /// A test-only `AsyncRead + AsyncWrite` pair equivalent to @@ -64,3 +65,67 @@ pub(crate) fn sink_empty_connection() -> Connection { Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)), ) } + +/// A `BidiStreamSource` that yields one pre-built `BiStream` on +/// `accept_bi` (then `ConnectionClosed`) and one on `open_bi` (then +/// `StreamClosed`). Used by end-to-end tests that wire a +/// `CallConnection` (client) to a `Dispatcher` (responder) over a real +/// `tokio::io::duplex` pair: the client's `open_bi` yields one end of +/// the duplex, the server's `accept_bi` yields the other. +pub(crate) struct SingleStreamSource { + stream: Mutex>, + addr: Option, +} + +impl SingleStreamSource { + pub(crate) fn new(stream: BiStream, addr: Option) -> Self { + Self { + stream: Mutex::new(Some(stream)), + addr, + } + } +} + +#[async_trait::async_trait] +impl BidiStreamSource for SingleStreamSource { + async fn accept_bi(&self) -> Result { + match self.stream.lock().expect("source mutex poisoned").take() { + Some(stream) => Ok(stream), + None => Err(StreamError::ConnectionClosed), + } + } + + async fn open_bi(&self) -> Result { + match self.stream.lock().expect("source mutex poisoned").take() { + Some(stream) => Ok(stream), + None => Err(StreamError::StreamClosed), + } + } + + fn remote_addr(&self) -> Option { + self.addr + } + + fn close(&self, _code: u32, _reason: &str) { + let _ = self.stream.lock().expect("source mutex poisoned").take(); + } +} + +/// Construct a pair of `Connection`s wired by a `tokio::io::duplex`: +/// the client's `open_bi` yields one end, the server's `accept_bi` +/// yields the other. Returns `(client_connection, server_connection)`. +/// Used by end-to-end tests that need a real round-trip through the +/// call protocol's stream-per-request model (P-12 #1, C-25 #1). +pub(crate) fn duplex_connection_pair(buffer: usize) -> (Connection, Connection) { + let (client_end, server_end) = tokio::io::duplex(buffer); + let addr = Some(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 4321)); + let client = Connection::from_source( + SingleStreamSource::new(BiStream::from_bidi(client_end), addr), + b"alknet/call".to_vec(), + ); + let server = Connection::from_source( + SingleStreamSource::new(BiStream::from_bidi(server_end), addr), + b"alknet/call".to_vec(), + ); + (client, server) +} diff --git a/src/registry/registration.rs b/src/registry/registration.rs index 8c11d00..7f2937e 100644 --- a/src/registry/registration.rs +++ b/src/registry/registration.rs @@ -274,27 +274,32 @@ impl OperationRegistry { streaming_handler(input, context) } - /// Dispatch a `Pub` operation (ADR-046). The `publish_stream` is the - /// initiator's data stream — each `call.published` chunk's `input` as - /// one `Ok(Value)` item, or an initiator-side error as - /// `Err(CallError)`. Pre-handler errors (not-found, forbidden, - /// `INVALID_OPERATION_TYPE` for a non-Pub op) return a single error - /// `ResponseEnvelope`. - pub async fn invoke_sink( + /// Resolve the `SinkHandler` for a `Pub` operation, performing the + /// not-found / visibility / ACL / handler-kind checks that gate sink + /// dispatch (ADR-046 §6). Returns the resolved `SinkHandler` on + /// success, or a single error `ResponseEnvelope` on any pre-handler + /// failure. + /// + /// This is the single source of truth for the sink dispatch checks — + /// both `invoke_sink` (the registry-level dispatch) and the wire + /// dispatch path (`Dispatcher::dispatch` Pub branch) call it, so the + /// two paths cannot diverge (P-08: the `invoke_sink` path IS the + /// wire path). + #[allow(clippy::result_large_err)] + pub(crate) fn resolve_sink_handler( &self, name: &str, - input: Value, - publish_stream: PublishStream, - context: OperationContext, - ) -> ResponseEnvelope { + input: &Value, + context: &OperationContext, + ) -> Result { let request_id = context.request_id.clone(); let registration = match self.operations.get(name) { Some(r) => r, - None => return ResponseEnvelope::not_found(request_id, name), + None => return Err(ResponseEnvelope::not_found(request_id, name)), }; if registration.spec.visibility == Visibility::Internal && !context.internal { - return ResponseEnvelope::not_found(request_id, name); + return Err(ResponseEnvelope::not_found(request_id, name)); } let acl = ®istration.spec.access_control; @@ -311,37 +316,50 @@ impl OperationRegistry { .spec .resource_id_path .as_ref() - .and_then(|path| extract_json_pointer(&input, path)); + .and_then(|path| extract_json_pointer(input, path)); if let AccessResult::Forbidden(message) = acl.check( identity.as_ref(), resource_id.as_deref(), context.ownership.as_deref(), ) { - return ResponseEnvelope::forbidden(request_id, message); + return Err(ResponseEnvelope::forbidden(request_id, message)); } - let sink_handler = match ®istration.handler { - HandlerKind::Sink(h) => Arc::clone(h), - HandlerKind::Once(_) => { - return ResponseEnvelope::error( - request_id, - CallError::invalid_operation_type( - "invoke_sink() called on a Query/Mutation op; use invoke()", - ), - ); - } - HandlerKind::Stream(_) => { - return ResponseEnvelope::error( - request_id, - CallError::invalid_operation_type( - "invoke_sink() called on a Sub op; use invoke_streaming()", - ), - ); - } - }; + match ®istration.handler { + HandlerKind::Sink(h) => Ok(Arc::clone(h)), + HandlerKind::Once(_) => Err(ResponseEnvelope::error( + request_id, + CallError::invalid_operation_type( + "invoke_sink() called on a Query/Mutation op; use invoke()", + ), + )), + HandlerKind::Stream(_) => Err(ResponseEnvelope::error( + request_id, + CallError::invalid_operation_type( + "invoke_sink() called on a Sub op; use invoke_streaming()", + ), + )), + } + } - (sink_handler)(input, context, publish_stream).await + /// Dispatch a `Pub` operation (ADR-046). The `publish_stream` is the + /// initiator's data stream — each `call.published` chunk's `input` as + /// one `Ok(Value)` item, or an initiator-side error as + /// `Err(CallError)`. Pre-handler errors (not-found, forbidden, + /// `INVALID_OPERATION_TYPE` for a non-Pub op) return a single error + /// `ResponseEnvelope`. + pub async fn invoke_sink( + &self, + name: &str, + input: Value, + publish_stream: PublishStream, + context: OperationContext, + ) -> ResponseEnvelope { + match self.resolve_sink_handler(name, &input, &context) { + Err(envelope) => envelope, + Ok(sink_handler) => sink_handler(input, context, publish_stream).await, + } } }