fix: Unit 1 — make publish() work end-to-end (P-01, P-02, P-05, P-08, P-09, P-12 #1)

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<Box<dyn Stream<Item = Value> + Send>> per ADR-046 §8, not Vec<Value>.
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<Instant> 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)
This commit is contained in:
glm-5.2 committed 2026-08-12 13:29:58 +00:00
1 parent f1d6796436
commit 502488cc72
9 files changed
+400 -188

No files matched your search

+16 -6
View File
@@ -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<Value> = publish_stream
.filter_map(|item| async move { item.ok() })
.collect()
.await;
connection.publish_with_payload(payload, chunks).await
let value_stream: Pin<Box<dyn Stream<Item = Value> + 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
}
})
}
+1 -1
View File
@@ -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()),
);
}
+3 -3
View File
@@ -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,
);
}
+227 -64
View File
@@ -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<Value>,
stream: Pin<Box<dyn Stream<Item = Value> + 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<Value>,
stream: Pin<Box<dyn Stream<Item = Value> + 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<Value>,
) -> 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<W>(
&self,
send: W,
@@ -398,6 +368,36 @@ impl CallConnection {
}
}
async fn pump_publish_to_wire<W>(
send: W,
request_id: &str,
payload: Value,
mut stream: Pin<Box<dyn Stream<Item = Value> + 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<R>(recv: R, pending: &Arc<Mutex<PendingRequestMap>>)
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<Identity> {
None
}
fn resolve_from_token(&self, _: &crate::core::auth::AuthToken) -> Option<Identity> {
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(&registry), 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<Box<dyn Stream<Item = serde_json::Value> + 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<Identity> {
None
}
fn resolve_from_token(&self, _: &crate::core::auth::AuthToken) -> Option<Identity> {
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(&registry), 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<Box<dyn Stream<Item = serde_json::Value> + 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;
}
}
+12 -64
View File
@@ -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 = &registration.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 &registration.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::<Result<Value, CallError>>(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()),
);
}
+1 -1
View File
@@ -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};
+21 -13
View File
@@ -15,7 +15,7 @@ pub struct PendingRequestMap {
pub(crate) enum PendingEntry {
Call {
tx: oneshot::Sender<Result<Value, CallError>>,
timeout: Instant,
timeout: Option<Instant>,
parent_request_id: Option<String>,
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<Instant>,
parent_request_id: Option<String>,
) -> oneshot::Receiver<Result<Value, CallError>> {
let (tx, rx) = oneshot::channel();
@@ -175,7 +180,10 @@ impl PendingRequestMap {
let mut to_remove: Vec<String> = 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,
);
+66 -1
View File
@@ -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<Option<BiStream>>,
addr: Option<SocketAddr>,
}
impl SingleStreamSource {
pub(crate) fn new(stream: BiStream, addr: Option<SocketAddr>) -> Self {
Self {
stream: Mutex::new(Some(stream)),
addr,
}
}
}
#[async_trait::async_trait]
impl BidiStreamSource for SingleStreamSource {
async fn accept_bi(&self) -> Result<BiStream, StreamError> {
match self.stream.lock().expect("source mutex poisoned").take() {
Some(stream) => Ok(stream),
None => Err(StreamError::ConnectionClosed),
}
}
async fn open_bi(&self) -> Result<BiStream, StreamError> {
match self.stream.lock().expect("source mutex poisoned").take() {
Some(stream) => Ok(stream),
None => Err(StreamError::StreamClosed),
}
}
fn remote_addr(&self) -> Option<SocketAddr> {
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)
}
+53 -35
View File
@@ -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<SinkHandler, ResponseEnvelope> {
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 = &registration.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 &registration.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 &registration.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,
}
}
}