-
Notifications
You must be signed in to change notification settings - Fork 0
Transports
In Vyra, a Transport is a dumb byte courier.
It is responsible only for delivering raw byte[] arrays to a destination. The transport knows nothing about Java types, routing semantics, or request/response correlation; it only needs the MessageKind (REQUEST, BROADCAST_REQUEST, RESPONSE, or EVENT) to select the best delivery primitive.
Every transport implements the simple Transport interface:
public interface Transport extends AutoCloseable {
// Hand raw envelope bytes to the network
CompletionStage<Void> send(MessageKind kind, String destination, byte[] envelope);
// Register interest in messages of a specific kind and destination
void subscribe(MessageKind kind, String destination);
// Provide the callback that receives incoming raw envelopes
void setMessageHandler(Function<byte[], CompletionStage<Void>> handler);
// Close transport connections
void close();
}The Redis transport uses the asynchronous, thread-safe Lettuce driver to deliver messages across distributed nodes.
Add vyra-redis to your build.gradle.kts:
implementation("com.marioded.vyra:vyra-redis:0.1.0")Initialize your Vyra instance with your existing Lettuce RedisClient:
import io.lettuce.core.RedisClient;
import com.marioded.vyra.redis.RedisVyra;
import com.marioded.vyra.jackson.JacksonSerializer;
RedisClient redisClient = RedisClient.create("redis://localhost:6379");
Vyra vyra = RedisVyra.redis(redisClient)
.serializer(JacksonSerializer.jackson())
.nodeId("order-service-1")
.build();Vyra uses the right Redis tool for each messaging pattern:
| Message kind | Mechanism | Behavior |
|---|---|---|---|
| REQUEST | push + blocking dequeue | Worker queue: distributed load balancing. One worker picks up each request. |
| BROADCAST_REQUEST | PUBLISH + SUBSCRIBE | 1-to-N broadcast: delivered to all workers on topic (bcast:<target>). Non-owners return null and stay silent; first non-null response replies to res:<nodeId>. |
| RESPONSE | PUBLISH + SUBSCRIBE | Point-to-point: routed directly to the requesting node's unique topic (res:<nodeId>). |
| EVENT | PUBLISH + SUBSCRIBE | Broadcast: delivered to all active subscribers on the topic (evt:<channel>). |
-
Dedicated blocking connections: each subscribed
REQUESTworker queue allocates its own dedicated connection for blocking dequeue calls, ensuring that blocking operations never stall broker threads. - Auto-reconnection: the blocking worker loop automatically reconnects if network hiccups or Redis failovers occur.
-
Client ownership (supplied = yours): when
vyra.close()is called, Vyra closes only the connections it opened. It never shuts down yourRedisClient. You can safely share a singleRedisClientacross multiple Vyra instances or other DAOs.
The in-memory transport requires zero dependencies and routes messages entirely within the JVM.
Add vyra-inmemory to your dependencies:
testImplementation("com.marioded.vyra:vyra-inmemory:0.1.0")To connect multiple instances in a test or local application, share a single InMemoryBus:
import com.marioded.vyra.inmemory.InMemoryVyra;
import com.marioded.vyra.inmemory.InMemoryBus;
import com.marioded.vyra.jackson.JacksonSerializer;
// The shared bus acts as the local network
InMemoryBus bus = new InMemoryBus();
Vyra authService = InMemoryVyra.inMemory(bus)
.serializer(JacksonSerializer.jackson())
.nodeId("auth-service")
.build();
Vyra clientService = InMemoryVyra.inMemory(bus)
.serializer(JacksonSerializer.jackson())
.nodeId("client-service")
.build();Tip
InMemoryVyra is ideal for unit tests and local development. It provides the exact same asynchronous behavior, serialization, and error handling as Redis without needing Docker or Redis servers running.
Need to run Vyra over NATS, RabbitMQ, or Kafka? You can implement the Transport interface yourself:
-
Opaque bytes: treat the
byte[] envelopeas completely opaque. Never try to parse or modify it. -
Namespace destinations: use
kind.prefix()to prevent collision between requests, responses, and events:-
MessageKind.REQUEST->req:<destination> -
MessageKind.BROADCAST_REQUEST->bcast:<destination> -
MessageKind.RESPONSE->res:<destination> -
MessageKind.EVENT->evt:<destination>
-
-
Dispatch safely: when bytes arrive from your network driver, pass them directly to the function supplied via
setMessageHandler(handler). Core will automatically offload deserialization and user handling away from your network thread. -
Acknowledgement: the message handler returns a
CompletionStage<Void>. If your broker supports manual ACKs (like RabbitMQ), ACK when the stage succeeds and NACK when it fails.
Then pass your transport to the builder:
Vyra vyra = Vyra.builder()
.transport(new MyCustomTransport())
.serializer(JacksonSerializer.jackson())
.build();- Explore wire formats in Serialization.
- Learn about thread pool sizing in Concurrency.
- Review the complete Configuration options.
Vyra documentation