Skip to main content

ConnectTo

Trait ConnectTo 

Source
pub trait ConnectTo<R: Role>: Send + 'static {
    // Required method
    fn connect_to(
        self,
        client: impl ConnectTo<R::Counterpart>,
    ) -> impl Future<Output = Result<()>> + Send;

    // Provided method
    fn into_channel_and_future(self) -> (Channel, Option<ConnectionDriver>)
       where Self: Sized { ... }
}
Expand description

A component that can exchange JSON-RPC messages to an endpoint playing the role R (e.g., an ACP Agent or an MCP Server).

This trait represents anything that can communicate via JSON-RPC messages over channels - agents, proxies, in-process connections, or any ACP-speaking component.

The type parameter R is the role that this component connects to (its counterpart). For example:

  • An agent implements ConnectTo<Client> to connect to clients
  • A proxy implements ConnectTo<Conductor> to connect to conductors
  • Transports like Channel implement ConnectTo<R> for every R because they are role-agnostic

§Component Types

The trait is implemented by several built-in types representing different communication patterns:

  • Lines: A component communicating over asynchronous line streams
  • ByteStreams: A component communicating over byte streams (stdin/stdout, sockets, etc.)
  • Channel: A component communicating via in-process message channels (for testing or direct connections)
  • Custom components: Proxies, transformers, or any ACP-aware service
  • AcpAgent: An external agent running in a separate process with stdio communication

§Two Ways to Connect

Components can be used in two ways:

  1. connect_to(client) - Connect directly to another component (most components implement this)
  2. into_channel_and_future() - Obtain a channel endpoint and optional owned driver

Most components only need to implement connect_to(client). The into_channel_and_future() method has a default implementation that creates an intermediate channel and calls connect_to.

§Implementation Example

use agent_client_protocol::{Agent, Client, ConnectTo, Result};

struct MyAgent;

impl ConnectTo<Client> for MyAgent {
    async fn connect_to(self, client: impl ConnectTo<Agent>) -> Result<()> {
        Agent.builder()
            .name("my-agent")
            .connect_to(client)
            .await
    }
}

§Heterogeneous Collections

For storing different component types in the same collection, use DynConnectTo:

use agent_client_protocol::{Channel, Client, DynConnectTo};

let (first, _first_peer) = Channel::duplex();
let (second, _second_peer) = Channel::duplex();
let components: Vec<DynConnectTo<Client>> = vec![
    DynConnectTo::new(first),
    DynConnectTo::new(second),
];
assert_eq!(components.len(), 2);

Required Methods§

Source

fn connect_to( self, client: impl ConnectTo<R::Counterpart>, ) -> impl Future<Output = Result<()>> + Send

Connect this component to another component.

Most components implement this method to set up their connection and exchange messages with the provided component.

§Arguments
  • client - The component to connect to (implements ConnectTo<R::Counterpart>)
§Returns

A future that resolves when the connection ends, either successfully or with an error. The future must be Send.

A component that buffers outbound messages should not return Ok(()) merely because its client completed: it should first finish messages the client already transferred to it. This lets wrappers preserve graceful drain guarantees through to the physical transport sink. Errors may still terminate the connection immediately.

Provided Methods§

Source

fn into_channel_and_future(self) -> (Channel, Option<ConnectionDriver>)
where Self: Sized,

Convert this component into a channel endpoint and optional owned driver.

The returned Channel is the canonical frame-aware boundary. It carries complete TransportFrame values so default adapters preserve batch grouping.

This method returns:

  • A Channel that can be used to communicate with this component
  • Some(ConnectionDriver) when the component owns work to drive
  • None when the channel halves alone own the endpoint’s lifetime

The default implementation creates an intermediate channel pair and calls connect_to on one endpoint while returning the other endpoint for the caller to use.

Base cases like Channel and ByteStreams override this to avoid unnecessary copying.

§Returns

A tuple of (Channel, Option<ConnectionDriver>). Owned drivers must be polled concurrently with channel traffic. Successful owned completion ends the endpoint after draining accepted output. None is not EOF: preserve both independent channel half-closes.

Absence must be handled explicitly; the optional driver is not awaitable:

ⓘ
use agent_client_protocol::{Channel, ConnectTo, UntypedRole};

let (channel, _peer) = Channel::duplex();
let (_channel, driver) = ConnectTo::<UntypedRole>::into_channel_and_future(channel);
driver.await?;

Once present, the owned driver itself is awaitable:

use agent_client_protocol::{Channel, ConnectionDriver, Result};

async fn drive_owned_work((_channel, driver): (Channel, Option<ConnectionDriver>)) -> Result<()> {
    if let Some(driver) = driver {
        // In a real adapter, also poll the channel traffic concurrently.
        driver.await?;
    }
    Ok(())
}

Dyn Compatibility§

This trait is not dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§

Source§

impl ConnectTo<Client> for AgentProtocolRouter

Available on crate feature unstable_protocol_v2 only.
Source§

impl ConnectTo<Conductor> for ProxyProtocolRouter

Available on crate feature unstable_protocol_v2 only.
Source§

impl<Counterpart: AcpAgentCounterpartRole> ConnectTo<Counterpart> for AcpAgent

Available on crate feature process and non-target_family=wasm only.
Source§

impl<Counterpart: Role> ConnectTo<Counterpart> for Stdio

Available on crate feature stdio and non-target_family=wasm only.
Source§

impl<OB, IB, R: Role> ConnectTo<R> for ByteStreams<OB, IB>
where OB: AsyncWrite + Send + 'static, IB: AsyncRead + Send + 'static,

Source§

impl<OutgoingSink, IncomingStream, R: Role> ConnectTo<R> for Lines<OutgoingSink, IncomingStream>
where OutgoingSink: Sink<String, Error = Error> + Send + 'static, IncomingStream: Stream<Item = Result<String>> + Send + 'static,

Source§

impl<R, H, Run, Close, Context> ConnectTo<<R as Role>::Counterpart> for Builder<R, H, Run, Close, Context>
where R: Role, H: HandleDispatchFrom<R::Counterpart> + 'static, Run: RunWithConnectionTo<R::Counterpart> + 'static, Close: HandleConnectionClose<R::Counterpart> + 'static, Context: ConnectionContext,

Source§

impl<R: Role> ConnectTo<R> for Channel

Source§

impl<R: Role> ConnectTo<R> for DynConnectTo<R>

Source§

impl<Run> ConnectTo<Client> for McpServer<Client, Run>
where Run: RunWithConnectionTo<Client> + 'static,