Skip to content

Request Response

Mario edited this page Sep 23, 2026 · 1 revision

Request / response messaging

Vyra provides type-safe, asynchronous request/response messaging across distributed services without heavyweight RPC frameworks, code generation, or fragile string mappings.


The request/response lifecycle

sequenceDiagram
    autonumber
    actor Caller as Caller thread
    participant RM as RequestManager
    participant TransA as Client transport
    participant Queue as Worker queue (req:target)
    participant TransB as Backend transport
    participant Exec as Handler executor
    participant Handler as Service handler

    Caller->>RM: request(target, req, Response.class, timeout)
    Note over RM: 1. Generate messageId<br/>2. Store pending future<br/>3. Schedule timeout
    RM->>TransA: send(REQUEST, target, envelope)
    TransA->>Queue: push to target queue
    Queue->>TransB: dequeue request
    TransB->>Exec: submit deserialization & handler
    Exec->>Handler: invoke handler(request)
    Handler-->>Exec: CompletableFuture<Response>
    Exec->>TransB: send(RESPONSE, callerNodeId, envelope)
    TransB->>TransA: deliver response
    TransA->>RM: match correlationId & complete future
    RM-->>Caller: complete CompletionStage<Response>
Loading

Step-by-step lifecycle

  1. Invocation: you call vyra.request("billing", request, InvoiceResponse.class, Duration.ofSeconds(3)).
  2. Correlation ID generation: the RequestManager generates a unique messageId (UUID) and creates a pending completion future.
  3. Safe timeout scheduling: the pending future is registered in a thread-safe map before scheduling the timeout, preventing race conditions where a fast response or timeout could fire into an empty map.
  4. Envelope serialization: the Protocol component creates a Message (kind = REQUEST, destination = "billing"), serializes the envelope using your configured serializer, and hands the raw bytes to the transport.
  5. Worker queue load balancing: requests are pushed to the target's worker queue (req:billing). If multiple nodes handle the "billing" target, the transport delivers the request to one available worker.
  6. Off-thread execution: the receiving node deserializes the envelope on the handler executor (never blocking transport network threads) and executes your handler function.
  7. Response delivery: when your handler's CompletionStage completes, the backend sends a Message (kind = RESPONSE, destination = <originNodeId>, correlationId = <originalMessageId>).
  8. Completion: the originating client matches the incoming response by its correlationId, converts the payload to InvoiceResponse.class, and completes the caller's CompletionStage.

Client usage: sending requests

The request(...) method returns a standard Java CompletionStage. You can compose it asynchronously or block on it when necessary.

Asynchronous pipeline (recommended)

client.request("billing", new GenerateInvoiceRequest("cust-42", 150.00), InvoiceResponse.class, Duration.ofSeconds(3))
        .thenApply(InvoiceResponse::invoiceId)
        .thenAccept(id -> System.out.println("Generated invoice: " + id))
        .exceptionally(ex -> {
            System.err.println("Could not generate invoice: " + ex.getMessage());
            return null;
        });

Synchronous / blocking

If you are running in a traditional blocking controller or thread:

try {
    InvoiceResponse response = client.request("billing", request, InvoiceResponse.class, Duration.ofSeconds(3))
            .toCompletableFuture()
            .join();
    System.out.println("Invoice ID: " + response.invoiceId());
} catch (CompletionException ex) {
    Throwable cause = ex.getCause();
    System.err.println("Request failed: " + cause.getMessage());
}

Server usage: handling requests

Handlers are registered using handle(...). Registering a handler automatically subscribes to the target's worker queue:

backend.handle("billing", GenerateInvoiceRequest.class, req -> {
    System.out.println("Processing invoice for: " + req.customerId());
    
    return database.saveInvoiceAsync(req.customerId(), req.amount())
            .thenApply(saved -> new InvoiceResponse(saved.id(), saved.status()));
});

Note

Handlers never run on transport event loops/threads. They are automatically dispatched to the handler executor so slow database calls or HTTP queries won't starve network threads.


Timeouts and late responses

Every request requires an explicit timeout Duration:

client.request("orders", req, OrderResponse.class, Duration.ofSeconds(2));
  • When a timeout triggers: the pending request future is removed from the active map and completes exceptionally with a VyraTimeoutException.
  • Late responses: if a response arrives from the server after the timeout has fired, the correlation lookup fails safely and the late message is silently discarded. No memory leaks occur.

Error propagation with RemoteError

If your handler throws an exception or returns a failed stage, Vyra catches it and sends back a structured RemoteError envelope instead of a raw stack trace:

public record RemoteError(Code code, String message) {
    public enum Code {
        HANDLER_ERROR,          // Handler threw an uncaught exception
        UNKNOWN_MESSAGE_TYPE,   // No handler was registered for this message type
        SERVER_BUSY,            // Handler executor queue is full
        SERIALIZATION_ERROR     // Payload could not be converted or deserialized
    }
}

On the client side, this surfaces cleanly as a VyraRemoteException:

client.request("orders", req, OrderResponse.class, Duration.ofSeconds(3))
    .exceptionally(ex -> {
        if (ex.getCause() instanceof VyraRemoteException remote) {
            System.err.println("Server returned error code: " + remote.getCode());
            System.err.println("Error message: " + remote.getMessage());
        }
        return null;
    });

See Timeouts & errors for full error recovery and retry strategies.


Broadcast request / response (vyra.broadcastRequest)

In partitioned architectures, state often lives in memory on one of multiple worker nodes (e.g., player sessions in a game cluster, distributed room state, cache partitions). If the client does not yet know which node holds the entity, you can use broadcast requests:

CompletionStage<PlayerResponse> stage = client.broadcastRequest(
    "game-cluster",
    new GetPlayerRequest("player-100"),
    PlayerResponse.class,
    Duration.ofSeconds(2)
);

How it works

sequenceDiagram
    autonumber
    participant Client as Client Node
    participant Bcast as Broadcast (bcast:game-cluster)
    participant Node1 as Worker Node 1
    participant Node2 as Worker Node 2 (Owner)
    participant Node3 as Worker Node 3

    Client->>Bcast: broadcastRequest("game-cluster", req, ...)
    Bcast->>Node1: deliver at t=0
    Bcast->>Node2: deliver at t=0
    Bcast->>Node3: deliver at t=0
    Note over Node1: entity not found -> returns null (silent)
    Note over Node3: entity not found -> returns null (silent)
    Note over Node2: entity found, return response
    Node2->>Client: direct response to res:clientNodeId (1 RTT)
    Client-->>Client: complete CompletionStage<PlayerResponse>
Loading

Every worker gets the request; only the owner answers. Non-owners stay silent, so the first non-null response wins.

  1. Silent non-owners: if a node does not hold the entity, its handler returns CompletableFuture.completedFuture(null). In broadcast mode, returning null or throwing an exception remains completely silent without sending back RemoteError, leaving the channel open for the true owner.
  2. Clean timeout: if no node holds the entity, the client's request times out naturally with VyraTimeoutException.
// On worker nodes:
node.handle("game-cluster", GetPlayerRequest.class, req -> {
    Player player = localMemory.get(req.playerId());
    if (player != null) {
        return CompletableFuture.completedFuture(new PlayerResponse(player));
    }
    // Returning null keeps this node completely silent:
    return CompletableFuture.completedFuture(null);
});

Important rules to remember

  1. Concrete response types: the responseType parameter must be a concrete Class<T>. You cannot pass List<Player>.class. Instead, wrap collections in a record:
    // Won't work: generics erased at runtime
    // client.request("players", req, List.class, timeout);
    
    // Recommended: wrap in a concrete record
    public record PlayerListResponse(List<Player> players) {}
  2. Registration order: always register your message types with register(...) before calling request(...), broadcastRequest(...), or handle(...).
  3. Silent non-owners: in broadcast requests, return CompletableFuture.completedFuture(null) if this instance does not hold the entity. Do not throw exceptions for "not found" in broadcast handlers.

Next steps

  • Explore one-way fire-and-forget messaging in Events.
  • Understand error propagation and retry strategies in Timeouts & errors.
  • Learn how to deploy worker queues in Transports.
  • Review the complete Best practices for production deployments.

Vyra documentation

Getting started and messaging patterns

Transports and formats

Inside the framework

Reference

Clone this wiki locally