🔌 Entrega gRPC
Un Decoded Shred Stream puede consumirse por gRPC en lugar de UDP: un stream saliente, orientado a conexión, que añade filtrado por cuentas del lado del servidor y una entrega ordenada sobre HTTP/2 por encima de las mismas transacciones con latencia submilisegundo — sin necesidad de un endpoint UDP accesible públicamente. Tu dashboard muestra el endpoint gRPC y el token de acceso de cada stream.
- Endpoint —
<host>:50051(un solo puerto para todos los flujos). - Transporte — gRPC sobre HTTP/2, en claro h2c (sin TLS). Es una decisión de latencia deliberada: no hay handshake TLS que pagar. Aísla el enlace a nivel de red — red privada, VPN o peering — nunca lo expongas a la internet abierta.
- Rutas —
shredstream.com.DecodedShredStreamService/SubscribeDecodedTransactions, o la ruta compatible drop-inshreder_binary.ShrederBinaryService/SubscribeBinaryTransactions. Ambas son estrictamente equivalentes (mismos mensajes protobuf, contenido idéntico byte a byte) — un cliente que venga de un endpoint Shreder / Raiden Pulse solo tiene que cambiar la dirección.
🔑 Autenticación
Cada llamada debe llevar tu token en los metadatos gRPC, en cualquiera de estas dos formas (ambas se aceptan):
authorization: Bearer <TOKEN> x-token: <TOKEN>
Un token ausente o inválido se rechaza con un estado uniforme UNAUTHENTICATED y el mensaje authentication refused. Comprueba tu token.
Recuerda: una conexión por token. Abrir un segundo stream con el mismo token expulsa al más antiguo.
🧭 Modelo de suscripción
La ruta es streaming bidireccional (stream request → stream response):
- Abre el stream y envía al menos una request que lleve una map de filtros nombrados:
{ "<name>": <filter>, … }. - El servidor entrega solo las transacciones que coinciden con al menos un filtro. Cada respuesta va etiquetada con el/los nombre(s) del/los filtro(s) que satisfizo (el campo
filters) — así puedes enrutar varias estrategias sobre un único stream. - Puedes enviar una nueva map en cualquier momento: reemplaza a la anterior en caliente, sin reconexión y sin hueco en el flujo.
Casos límite:
- Map vacía (
{}) ⇒ no se entrega nada (te suscribes a filtros, no al firehose). - Filtro nombrado vacío (
{ "all": {} }) ⇒ toda transacción pasa, etiquetada como"all". - Una vez enviados tus filtros, no tienes nada más que enviar: el stream sigue entregando con esos filtros hasta que envíes otros nuevos o cierres la conexión.
🎯 Filtros por cuentas
Cada filtro nombrado son tres listas de claves públicas base58, combinadas con Y lógico (AND):
| Campo | Significado |
|---|---|
account_include | la transacción debe tocar al menos una de estas cuentas (vacío = sin restricción) |
account_exclude | la transacción no debe tocar ninguna de estas cuentas |
account_required | la transacción debe tocar todas estas cuentas |
Para las transacciones, «tocar» se evalúa sobre las claves de cuenta estáticas de la transacción (firmantes incluidos). Las direcciones resueltas a través de Address Lookup Tables (ALTs) no están presentes en el sobre de la transacción y, por tanto, no son filtrables — filtra solo sobre las claves estáticas.
📦 El payload
Cada respuesta lleva la transacción como bytes en bruto más su slot:
message SubscribeUpdateBinaryTransaction {BinaryTransaction transaction = 1;uint64 slot = 2; // slot de Solana}message BinaryTransaction {repeated bytes signatures = 1; // firmas (64 bytes cada una)bytes binary_transaction = 3; // VersionedTransaction, bincode, VERBATIM}
binary_transaction es el formato wire estándar de Solana — una VersionedTransaction serializada con bincode, sin tocar. Deserialízala con cualquier SDK de Solana y pásala a tu parser existente, exactamente igual que en el modo UDP. Las transacciones de voto se excluyen por defecto.
💻 Cliente oficial
Los clientes oficiales decoded-shredstream se encargan de esta ruta por ti — conexión, token en los metadatos, map de filtros, reconexión con los filtros actuales enviados de nuevo, y la transacción con su slot y sus firmas. El mismo paquete para UDP y gRPC.
// npm install decoded-shredstreamimport { DecodedShredStream, FilterAll } from "decoded-shredstream";const client = await DecodedShredStream.grpc({endpoint: "<host>:50051",token: process.env.DECODED_SHREDSTREAM_TOKEN!,filters: { all: FilterAll }, // or: { "watched-wallet": { include: [wallet] } }});for await (const tx of client.transactions()) {console.log(`slot=${tx.slot} sig=${tx.signature.toBase58()} matched=${JSON.stringify(tx.filters)}`);}
# pip install decoded-shredstreamfrom decoded_shredstream import Client, Filter, GrpcConfigwith Client.grpc(GrpcConfig(endpoint="<host>:50051", token="<TOKEN>",filters={"all": Filter()})) as client:for update in client:print(update.slot, len(update.data), list(update.filters))
// cargo add decoded-shredstream tokio --features tokio/macros,tokio/rt-multi-threaduse decoded_shredstream::{Filter, GrpcClient, GrpcConfig};let mut client = GrpcClient::connect(GrpcConfig::new("<host>:50051", "<TOKEN>").filter("all", Filter::all()),).await?;while let Some(update) = client.next_update().await {let update = update?;println!("slot={} sig={} matched={:?}", update.slot(), update.signature(), update.filters());}
// go get github.com/shredstream/decoded-shredstream-goclient, err := decodedshredstream.NewGRPC(decodedshredstream.GRPCConfig{Endpoint: "<host>:50051",Token: "<TOKEN>",Filters: decodedshredstream.Filters{"all": decodedshredstream.FilterAll()},})if err != nil { log.Fatal(err) }defer client.Close()err = client.Run(ctx, func(u *decodedshredstream.TransactionUpdate) {sig, _ := u.Signature()fmt.Println(u.Slot, sig, u.Filters)})
Los filtros pueden reemplazarse en un stream en vivo (updateFilters / update_filters) sin reconexión. Las interrupciones recuperables nunca llegan hasta ti — el cliente se reconecta y envía de nuevo la map de filtros actual; solo un token rechazado, una sesión cerrada por el servidor o una map de filtros rechazada terminan el stream, mediante un error que gestionas una sola vez.
💻 Clientes gRPC estándar (stubs generados)
Para los equipos que prefieren su propia pila gRPC: el contrato es decoded.proto (nuestra ruta, shredstream.com.DecodedShredStreamService), que importa shreder_binary.proto para sus mensajes — ambos ficheros vienen incluidos en cada cliente oficial. Generar los stubs solo a partir de shreder_binary.proto te da la ruta compatible drop-in que se usa a continuación.
Rust (tonic)
use std::collections::HashMap;use shreder_binary::shreder_binary_service_client::ShrederBinaryServiceClient;use shreder_binary::{SubscribeBinaryTransactionsRequest, SubscribeRequestFilterBinaryTransactions};use solana_transaction::versioned::VersionedTransaction;use tonic::Request;let mut client = ShrederBinaryServiceClient::connect("http://<host>:50051").await?;let filters = HashMap::from([("pumpfun".to_string(),SubscribeRequestFilterBinaryTransactions {account_include: vec!["6EF8rrecthR5Dkzon8Nwu78hRvfCKubJ14M5uBEwF6P".into()],account_exclude: vec![],account_required: vec![],},)]);// Filtros estáticos: se emite una request y luego se hace half-close (el iterador termina).let outbound = tokio_stream::iter(vec![SubscribeBinaryTransactionsRequest { transactions: filters }]);let mut req = Request::new(outbound);req.metadata_mut().insert("authorization", "Bearer <TOKEN>".parse()?);let mut stream = client.subscribe_binary_transactions(req).await?.into_inner();while let Some(resp) = stream.message().await? {if let Some(tx) = resp.transaction.and_then(|u| u.transaction) {let vtx: VersionedTransaction = bincode::deserialize(&tx.binary_transaction)?;// → enruta según resp.filters, procesa vtx …}}
Nuestra ruta equivalente es
DecodedShredStreamServiceClient::subscribe_decoded_transactions, con exactamente los mismos tipos de request/response. Para actualizar los filtros en mitad del stream, reemplazatokio_stream::iterpor un canal que mantengas abierto y sobre el que hagassendde una nueva request — sin necesidad de reconexión.
TypeScript (@grpc/grpc-js)
import { Metadata } from "@grpc/grpc-js";// … cliente generado desde shreder_binary.proto …const meta = new Metadata();meta.set("x-token", "<TOKEN>");const call = client.subscribeBinaryTransactions(meta);call.write({transactions: { pumpfun: { accountInclude: ["6EF8rrec…"], accountExclude: [], accountRequired: [] } },});// call.end(); // half-close si los filtros no cambiancall.on("data", (resp) => {const raw = resp.transaction?.transaction?.binaryTransaction; // Buffer// deserializa una VersionedTransaction del lado del cliente …});call.on("error", (e) => { /* DATA_LOSS / UNAUTHENTICATED → reconexión / log */ });
Python (grpcio)
import grpcimport shreder_binary_pb2 as pb, shreder_binary_pb2_grpc as rpcchan = grpc.insecure_channel("<host>:50051")stub = rpc.ShrederBinaryServiceStub(chan)def requests():yield pb.SubscribeBinaryTransactionsRequest(transactions={"pumpfun": pb.SubscribeRequestFilterBinaryTransactions(account_include=["6EF8rrec…"])})md = (("x-token", "<TOKEN>"),)for resp in stub.SubscribeBinaryTransactions(requests(), metadata=md):raw = resp.transaction.transaction.binary_transaction # bytes bincode# deserializa una VersionedTransaction …
Prueba rápida con grpcurl
grpcurl -plaintext -proto shreder_binary.proto \-H 'x-token: <TOKEN>' \-d '{"transactions":{"all":{}}}' \<host>:50051 shreder_binary.ShrederBinaryService/SubscribeBinaryTransactions
🧯 Gestión de errores y reconexión
| Estado gRPC | Significado | Qué hacer |
|---|---|---|
UNAUTHENTICATED | token ausente / inválido / flujo equivocado | corrige el token o el flujo — no reintentes a ciegas |
INVALID_ARGUMENT | base58 inválido, demasiadas pubkeys, demasiados filtros nombrados | corrige el filtro problemático |
DATA_LOSS | tu consumidor es demasiado lento — la cola de envío del servidor desbordó y el stream se cerró | reconéctate y reenvía tu map de filtros |
PERMISSION_DENIED | desconectado/revocado por el operador | no hagas bucle; contacta con el operador |
El servicio es solo en tiempo real — no hay replay ni backfill. Ante un DATA_LOSS o una caída del transporte, reconéctate con backoff exponencial (y algo de jitter) y luego reenvía tu map de filtros; un breve hueco de datos durante la reconexión es normal. Por la regla de una-conexión-por-token, asegúrate de que solo corre un bucle de reconexión por token.
Para evitar DATA_LOSS de entrada, vacía el stream sin bloquear: delega el procesamiento pesado (incluida la deserialización bincode) a una cola o un pool de workers para que tu bucle de recepción nunca se atasque. gRPC sobre HTTP/2 garantiza el orden y la integridad de todo lo emitido — la única pérdida posible es el drop «cliente demasiado lento», y siempre se señaliza explícitamente con DATA_LOSS.
⚖️ ¿gRPC o UDP?
| gRPC | UDP | |
|---|---|---|
| Latencia | Baja — h2c, sin handshake TLS | La más baja — datagramas binarios, sin conexión |
| Entrega | Ordenada, fiable; «cliente demasiado lento» señalizado con DATA_LOSS | Push best-effort; pérdida posible, detectada por saltos de seq |
| Filtrado | Filtros por cuentas del lado del servidor (claves estáticas) | Ninguno — filtra del lado del cliente tras recibir |
| Accesibilidad | Conexión saliente — funciona detrás de NAT/firewalls | Requiere un endpoint UDP accesible públicamente |
| Payload | VersionedTransaction (bincode) | VersionedTransaction (bincode), en un sobre de 16 bytes |
Elige gRPC cuando quieras filtrado del lado del servidor, entrega ordenada y conectividad compatible con NAT. Elige UDP para la latencia absolutamente más baja cuando tengas un endpoint accesible y monitorees tú mismo las pérdidas. El payload es idéntico en ambos casos.
➡️ Próximos pasos
- Recibir transacciones decodificadas — el sobre UDP, los offsets, la detección de pérdidas y un decodificador mínimo.
- Decoded Shred Stream — posicionamiento, latencia y modos de entrega.