Skip to content

proxystore.endpoint.dispatch

Dispatch requests to the storage of an endpoint or to its peers.

Requests from clients (see ClientHandler) and from peer endpoints (see PeerManager) are both handled by a Dispatcher. A request for this endpoint is performed on its storage. A request whose target is another endpoint is forwarded to that peer without the target so the peer performs the request itself. The response of the peer is returned to the client, except that the peer is named in an error message and the response to a PING is replaced with the latency of and path to the peer measured by this endpoint.

To add an operation, add an Op, add its case to Dispatcher._handle_local() in this module, and add a method to the EndpointClient. Requests for the operation are forwarded to peers without changes to this module.

Dispatcher

Dispatcher(
    endpoint_id: EndpointId,
    storage: Storage,
    peer_manager: PeerManager | None = None,
)

Dispatches requests to the storage of an endpoint or to its peers.

Example
dispatcher = Dispatcher(endpoint_id, MemoryStorage())
request = Message(Op.SET, Request(key='key').encode(), b'value')
response = await dispatcher.handle(request)
assert response.code == Status.OK

Parameters:

  • endpoint_id (EndpointId) –

    ID of the endpoint.

  • storage (Storage) –

    Storage of the endpoint.

  • peer_manager (PeerManager | None, default: None ) –

    Peer manager used to forward requests to peers or None if peering is disabled. The owner of the peer manager is responsible for starting it with handle_peer_request() as its handler and for closing it.

Source code in proxystore/endpoint/dispatch.py
def __init__(
    self,
    endpoint_id: EndpointId,
    storage: Storage,
    peer_manager: PeerManager | None = None,
) -> None:
    if peer_manager is not None and peer_manager.id != endpoint_id:
        raise ValueError(
            f'The ID of the peer manager ({peer_manager.id}) does not '
            f'match the ID of the endpoint ({endpoint_id}).',
        )
    self._id = endpoint_id
    self._storage = storage
    self._peer_manager = peer_manager

id property

ID of the endpoint.

storage property

storage: Storage

Storage of the endpoint.

peer_manager property

peer_manager: PeerManager | None

Peer manager used to forward requests to peers.

handle async

handle(request: Message) -> Message

Handle a request from a client.

A request whose target is another endpoint is forwarded to that peer.

Parameters:

  • request (Message) –

    Request message.

Returns:

Source code in proxystore/endpoint/dispatch.py
async def handle(self, request: Message) -> Message:
    """Handle a request from a client.

    A request whose target is another endpoint is forwarded to that
    peer.

    Args:
        request: Request message.

    Returns:
        The response message. Errors are returned as responses rather \
        than raised (see
        [`Message.from_error()`][proxystore.endpoint.protocol.Message.from_error]).
    """
    return await self._handle(request, source='client', forward=True)

handle_peer_request async

handle_peer_request(
    peer_id: EndpointId, request: Message
) -> Message

Handle a request from a peer endpoint.

Requests from peers are handled like requests from clients except they are never forwarded to another peer. This is the RequestHandler of the peer manager.

Parameters:

  • peer_id (EndpointId) –

    ID of the peer which sent the request.

  • request (Message) –

    Request message.

Source code in proxystore/endpoint/dispatch.py
async def handle_peer_request(
    self,
    peer_id: EndpointId,
    request: Message,
) -> Message:
    """Handle a request from a peer endpoint.

    Requests from peers are handled like requests from clients except
    they are never forwarded to another peer. This is the
    [`RequestHandler`][proxystore.endpoint.p2p.manager.RequestHandler]
    of the peer manager.

    Args:
        peer_id: ID of the peer which sent the request.
        request: Request message.
    """
    source = f'peer {self._peer_name(peer_id)}'
    return await self._handle(request, source=source, forward=False)