Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions md/SUMMARY.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@

# Reference

- [Migrating Connection Drivers](./migration-connection-drivers.md)
- [Migrating the rmcp Integration to v4](./migration-rmcp-v4.md)
- [Migrating to v2.0](./migration_v2.0.md)
- [Migrating to v0.11](./migration_v0.11.x.md)
Expand Down
221 changes: 221 additions & 0 deletions md/migration-connection-drivers.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,221 @@
# Migrating Connection Drivers

`ConnectTo::into_channel_and_future` now returns
`(Channel, Option<ConnectionDriver>)` instead of
`(Channel, BoxFuture<'static, Result<()>>)`. The same change applies when
accessing a component through `DynConnectTo`.

This is a source-breaking transport-adapter change. It does not change ACP wire
messages, the raw `Channel` sender/receiver types, or `unbounded_send`, and it
introduces no new frame-size, queue, or task limits.

## Components using the default conversion

If your component implements only `connect_to`, no change is needed. The
default conversion still creates a channel pair and drives your component, now
returning `Some(ConnectionDriver)`. That default wraps opaque work; it cannot
infer a physical finish hook. A buffered transport that needs a finite
foreground to await physical flush should override normalization with
`with_finish`, as described below.

Low-level callers must handle the optional work explicitly. The optional value
is not a future: awaiting it directly no longer compiles. For a component that
is known to own work, extract its driver before polling it:

```rust,ignore
let (channel, driver) = component.into_channel_and_future();
let driver = driver.expect("this component owns connection work");
// Use channel while continuing to poll the driver.
driver.await?;
```

For a generic component, handle both cases: poll `Some(driver)` alongside
traffic and drain accepted output on completion; for `None`, retain the
channel's independent halves until they close. Absence of work is not EOF.
Do not replace `None` with a ready-success future in a shutdown race.

## Custom conversion overrides

Import `ConnectionDriver` from `agent_client_protocol` and change the return
type. Wrap a future that owns the connection work with `ConnectionDriver::new`:

```rust,ignore
fn into_channel_and_future(self) -> (Channel, Option<ConnectionDriver>) {
let (channel, future) = self.into_channel_transport();
(channel, Some(ConnectionDriver::new(future)))
}
```

For an endpoint whose work is driven elsewhere, return `None` instead of
wrapping a ready no-op future:

```rust,ignore
fn into_channel_and_future(self) -> (Channel, Option<ConnectionDriver>) {
(self.channel, None)
}
```

An existing `Channel` has no driver. There is no awaitable passive sentinel,
and no finish hook belongs to the `None` case. `ConnectionDriver` always holds
real owned work; cooperative drivers additionally support a finish hook.
Passive bridges retain each read/write half until its
own closure; an input half-close can still be followed by a final response.

If a wrapper simply exposes another component's endpoint, return its original
`(channel, optional_driver)` pair. Re-boxing an owned driver and wrapping it
with `new` would hide its finish capability. Inventing a ready
driver for `None` would also turn absence into a false completion signal.

For tracing, error annotation, or completion cleanup, decorate the future with
`map_future`. This preserves both the finish capability and any already-issued
request; opaque work stays opaque:

```rust,ignore
use futures::FutureExt;

let (channel, driver) = component.into_channel_and_future();
let driver = driver.map(|driver| {
driver.map_future(|work| {
work.inspect(|result| eprintln!("transport completed: {result:?}"))
})
});
(channel, driver)
```

The transformed future must still drive the original work and must not report
success before its accepted output has drained.

## Completion and drain responsibilities

Poll owned work and outbound forwarding concurrently. A driver may need its
outbound request to be delivered before it can receive a response and finish.

An adapter must not report success before flushing output it already accepted.
Use `ConnectionDriver::with_finish(future, finish)` for a custom normalized
transport that needs to flush during finite foreground shutdown. The
nonblocking `FnOnce()` hook requests graceful completion; the future proves
completion and reports any I/O error.

```rust,ignore
let (finish_tx, finish_rx) = futures::channel::oneshot::channel();
let future = async move {
// Keep processing input and output while waiting for the finish request.
// A dropped sender is not a finish request; it may simply mean that
// finish control was abandoned while normal half-closes remain in use.
//
// After a successful signal, seal the outgoing queue, drain every accepted
// frame, and flush/close the physical write half. Do not wait for remote
// read EOF; continue observing genuine I/O errors during the drain.
run_custom_adapter(outgoing_rx, physical_io, finish_rx).await
};
let driver = ConnectionDriver::with_finish(future, move || {
let _ = finish_tx.send(());
});
(channel, Some(driver))
```

SDK shutdown coordination invokes this hook only after protocol output has
been handed off to the normalized transport, then awaits the driver. Low-level
callers can use `driver.request_finish()` themselves. A `true` return means
cooperative finish is supported and has been requested, including a request
already issued. Requests are idempotent, but the hook runs only once. A `false`
return means opaque work, not "already requested."

Requesting finish does not prove output has finished flushing; continue polling
or await the driver. Capability remains intact if that driver is handed to
another owner while flushing. Quiesce and hand off output before requesting
finish; idempotence does not permit new output after sealing. Dropping the
driver drops its owned future without a graceful request. Dropping only the
hook does not invoke it or necessarily stop that future.

There is no implicit timeout. A cooperative adapter that cannot flush keeps
the connection pending, so applications that need a deadline must impose one
and accept that cancelling it can truncate output. `with_finish` declares the
adapter's contract; it cannot make an arbitrary future or external buffer
flush automatically.

Built-in `Lines` and `ByteStreams` preserve normal half-close behavior. When
their owner explicitly finishes, they drain accepted output while continuing
to poll incoming I/O for errors, rather than waiting for unrelated remote
input to reach EOF. Errors may terminate the connection without graceful drain.

## Direct adapter entry point

Returning a cooperative driver from `into_channel_and_future` lets normalized
SDK consumers coordinate finish. A custom transport's direct `connect_to`
implementation must coordinate it too: `try_join!(bridge, driver)` alone can
wait forever after a finite peer has returned.

This scaffold follows the built-in `Lines` policy, using only public APIs:

```rust,ignore
use agent_client_protocol::{Channel, ConnectTo, ConnectionDriver, Result, UntypedRole};
use futures::{future::{select, Either}, FutureExt};

struct BufferedAdapter {
channel: Channel,
driver: ConnectionDriver,
}

impl ConnectTo<UntypedRole> for BufferedAdapter {
async fn connect_to(self, peer: impl ConnectTo<UntypedRole>) -> Result<()> {
let bridge = Box::pin(self.channel.connect_to(peer));
match select(bridge, self.driver).await {
Either::Left((result, mut driver)) => {
result?; // The peer's accepted output has been handed off.
if driver.request_finish() {
driver.await // Prove physical drain; propagate its errors.
} else {
// Preserve a ready error, then cancel opaque work.
driver.now_or_never().unwrap_or(Ok(()))
}
}
Either::Right((result, _bridge)) => result,
}
}

fn into_channel_and_future(self) -> (Channel, Option<ConnectionDriver>) {
(self.channel, Some(self.driver))
}
}
```

The adapter's own future must own/close its physical producers before reporting
completion. Its finish implementation must stop forwarding successful input
to a completed peer while still observing genuine read errors during drain.
Otherwise, late input can fail against the dropped receiver and cancel final
output. Both direct and normalized entry points should be tested with output
backpressure and independently open input.

## Finite foreground shutdown

On successful `Builder::connect_with` foreground completion, routable queued
output is drained. Requests still blocked on unresolved readiness are failed
and removed rather than published after shutdown. Physical transport and
protocol progress continue during this drain; queued application tasks are
not started merely to finish the sink. The inherited cleanup coordinator can
still poll application tasks while protecting a close callback already
underway; a blocked close callback can delay completion.

Foreground success stops beginning new application delivery or close
callbacks. Physical reads remain driven without delivering their input to the
completed foreground. A close callback already underway finishes before the
outgoing drain boundary seals, and its errors retain precedence.

Protocol connectors and routers use the same ownership-aware rule. An owned
foreground's completion requests cooperative drain instead of waiting for
unrelated remote input. Initialization rejection also hands off its reply
before requesting finish. Passive half-closes alone do not request finish;
they preserve the other direction for a final response.

Cooperative drivers, both built-in and custom, are awaited through physical
write shutdown.
An opaque driver constructed with `ConnectionDriver::new(future)` has no
externally requestable finish control: finite foreground shutdown transfers
protocol output into its normalized channel, then cancels that work without
guaranteeing custom physical flush. Reactive `connect_to` still joins owned
work after input EOF. Choose `new` for opaque cancellable work and `with_finish`
when the adapter can honor an explicit graceful-finish request.

See [Transport Architecture](./transport-architecture.md#component-boundary)
for the active/passive boundary and forwarding rules.
66 changes: 58 additions & 8 deletions md/transport-architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -45,7 +45,7 @@ by the JSON-RPC envelope types from `agent-client-protocol-schema`:
enum RawJsonRpcMessage {
Request(Request<RawJsonRpcParams>),
Notification(Notification<RawJsonRpcParams>),
Response(Response<serde_json::Value>),
Response(RawJsonRpcResponse),
}
```

Expand Down Expand Up @@ -262,16 +262,64 @@ Ordering](./conductor.md#routing-and-ordering).
is the common component and transport abstraction. `connect_to` joins a
component to its counterpart and drives the connection until completion.
`into_channel_and_future` exposes the canonical low-level boundary as a
`Channel` plus the future that drives the component:
`Channel` plus an explicit connection driver:

```rust,ignore
fn into_channel_and_future(self) -> (Channel, BoxFuture<'static, Result<()>>);
fn into_channel_and_future(self) -> (Channel, Option<ConnectionDriver>);
```

The returned future owns transport failures and lifecycle completion. The
channel carries only `TransportFrame` wire events. Most components implement
only `connect_to`; direct transports override `into_channel_and_future` to avoid
an intermediate copy.
The channel carries only `TransportFrame` wire events. The optional driver
distinguishes owned work from a passive endpoint:

| Returned work | Lifetime rule |
| --- | --- |
| `Some(ConnectionDriver::new(future))` | Poll the owned work alongside traffic; successful completion ends it after accepted output is drained. |
| `None` | No work is owned here; each channel half determines its own lifetime. |

A `ConnectionDriver` always contains a real future; the optional return value
cannot itself be awaited. There is no ready-successful passive driver and no
finish hook in the `None` case. This makes the ownership decision explicit
rather than requiring callers to recognize a special future.

A raw `Channel` returns `None`. Its bridge preserves both directions
independently:
one sender closing must not prevent a final response in the reverse direction.
Owned completion lets a bridge stop accepting new output, drain frames already
accepted, and finish without waiting for unrelated remote input to close.
Outbound forwarding must remain polled while owned work is running; otherwise
a component waiting for a response to its own request could deadlock.

Buffered adapters are responsible for flushing their accepted output before
reporting successful completion. The built-in line and byte-stream adapters
keep the read half moving during write drain and propagate incoming errors;
their explicit finish handling does not require remote read EOF.
`ConnectionDriver::with_finish(future, finish)` lets custom adapters declare
the same cooperative contract. Its one-shot, nonblocking hook requests finish
after protocol output handoff; the still-polled future performs the drain and
reports completion or I/O errors. `request_finish()` exposes this request to
low-level callers and is idempotent: a supported request remains supported
after it has been issued. The callback still runs at most once. Neither
requesting finish nor dropping the hook proves a successful flush, and no
implicit timeout is imposed.

`ConnectionDriver::new(future)` remains appropriate for opaque cancellable
work. A finite foreground does not wait indefinitely for such work after
handing off protocol output. Merely wrapping an arbitrary future cannot make
an opaque custom adapter drain safely.

Most components implement only `connect_to`; default normalization supplies
`Some(driver)` containing the owned work. Direct transports override
`into_channel_and_future` to avoid an intermediate copy. Wrappers that expose an
existing endpoint should forward
its driver unchanged so absence and cooperative completion handling are
not lost. Wrappers that decorate execution use `map_future` to transform the
owned future while preserving its finish capability and requested state.
Direct custom transport entry points must also request and await cooperative
finish after a finite peer completes; a bare join does not provide that step.

See [Migrating Connection Drivers](./migration-connection-drivers.md) for custom
override changes. This lifecycle distinction does not change the existing raw
channel types or introduce frame-size, queue, or task limits.

## Transport Implementations

Expand Down Expand Up @@ -346,7 +394,9 @@ Split the socket and pass compatible read/write halves to `ByteStreams::new`.
Embedders supply and drive their own runtime and host transport:

- Exchange `TransportFrame` values through an in-component `Channel`. A caller
using `ConnectTo::into_channel_and_future` must poll the returned future.
using `ConnectTo::into_channel_and_future` polls a present owned driver
alongside traffic. When no driver is returned, preserve the channel halves'
independent lifetimes; absence is not EOF.
- Exchange newline-delimited JSON through `Lines`, using a
`futures::Sink<String>` and
`futures::Stream<Item = std::io::Result<String>>`.
Expand Down
Loading
Loading