🔌 Доставка по gRPC

Decoded Shred Stream можно потреблять по gRPC вместо UDP: исходящий, ориентированный на соединение стрим, который добавляет фильтрацию по аккаунтам на стороне сервера и упорядоченную доставку по HTTP/2 поверх тех же транзакций с субмиллисекундной задержкой — без необходимости в публично доступной конечной точке UDP. В дашборде показаны gRPC-эндпоинт и токен доступа для каждого стрима.

  • Эндпоинт<host>:50051 (один порт для всех потоков).
  • Транспорт — gRPC поверх HTTP/2, в открытом виде h2c (без TLS). Это осознанный выбор ради задержки: не нужно платить за рукопожатие TLS. Изолируйте канал на сетевом уровне — приватная сеть, VPN или пиринг — и никогда не открывайте его в публичный интернет.
  • Маршрутыshredstream.com.DecodedShredStreamService/SubscribeDecodedTransactions или совместимый «из коробки» shreder_binary.ShrederBinaryService/SubscribeBinaryTransactions. Оба строго эквивалентны (те же protobuf-сообщения, побайтово идентичное содержимое) — клиенту, приходящему с эндпоинта Shreder / Raiden Pulse, достаточно сменить адрес.

🔑 Аутентификация

Каждый вызов должен передавать ваш токен в метаданных gRPC в одной из двух форм (принимаются обе):

authorization: Bearer <TOKEN>
x-token: <TOKEN>

Отсутствующий или недействительный токен отклоняется единым статусом UNAUTHENTICATED с сообщением authentication refused. Проверьте свой токен.

Помните: одно соединение на токен. Открытие второго стрима с тем же токеном вытесняет самый старый.


🧭 Модель подписки

Маршрут — это двунаправленный стриминг (stream request → stream response):

  1. Откройте стрим и отправьте хотя бы один запрос, содержащий map именованных фильтров: { "<name>": <filter>, … }.
  2. Сервер доставляет только те транзакции, что совпадают хотя бы с одним фильтром. Каждый ответ помечается именем(-ами) сработавших фильтров (поле filters) — так вы можете маршрутизировать несколько стратегий по одному стриму.
  3. Вы можете отправить новую map в любой момент: она заменяет предыдущую на лету, без переподключения и без разрыва в потоке.

Крайние случаи:

  • Пустая map ({}) ⇒ ничего не доставляется (вы подписываетесь на фильтры, а не на firehose).
  • Пустой именованный фильтр ({ "all": {} }) ⇒ проходит каждая транзакция, с меткой "all".
  • После отправки фильтров отправлять больше нечего: стрим продолжает доставку с этими фильтрами, пока вы не отправите новые или не закроете соединение.

🎯 Фильтры по аккаунтам

Каждый именованный фильтр — это три списка публичных ключей в base58, объединённых логическим И:

ПолеЗначение
account_includeтранзакция должна затрагивать хотя бы один из этих аккаунтов (пустой = без ограничений)
account_excludeтранзакция не должна затрагивать ни один из этих аккаунтов
account_requiredтранзакция должна затрагивать все эти аккаунты

Для транзакций «затрагивание» оценивается по статическим ключам аккаунтов транзакции (включая подписантов). Адреса, разрешаемые через Address Lookup Tables (ALT), отсутствуют в конверте транзакции и потому не поддаются фильтрации — фильтруйте только по статическим ключам.


📦 Полезная нагрузка

Каждый ответ несёт транзакцию в виде сырых байтов вместе с её слотом:

proto
message SubscribeUpdateBinaryTransaction {
BinaryTransaction transaction = 1;
uint64 slot = 2; // слот Solana
}
message BinaryTransaction {
repeated bytes signatures = 1; // подписи (по 64 байта каждая)
bytes binary_transaction = 3; // VersionedTransaction, bincode, VERBATIM
}

binary_transaction — это стандартный wire-формат Solana — сериализованная через bincode VersionedTransaction, без изменений. Десериализуйте её любым SDK Solana и передайте вашему существующему парсеру, ровно как в режиме UDP. Транзакции голосования по умолчанию исключены.


💻 Официальный клиент

Официальные клиенты decoded-shredstream берут этот маршрут на себя — соединение, токен в метаданных, map фильтров, переподключение с повторной отправкой текущих фильтров и транзакцию вместе с её слотом и подписями. Один и тот же пакет для UDP и gRPC.

ts
// npm install decoded-shredstream
import { 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)}`);
}
python
# pip install decoded-shredstream
from decoded_shredstream import Client, Filter, GrpcConfig
with 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))
rust
// cargo add decoded-shredstream tokio --features tokio/macros,tokio/rt-multi-thread
use 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());
}
Gogo
// go get github.com/shredstream/decoded-shredstream-go
client, 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)
})

Фильтры можно заменять на живом стриме (updateFilters / update_filters) без переподключения. Восстановимые сбои до вас не доходят — клиент переподключается и заново отправляет текущую map фильтров; стрим завершают только отклонённый токен, закрытая сервером сессия или отклонённая map фильтров — через ошибку, которую вы обрабатываете один раз.


💻 Стандартные gRPC-клиенты (сгенерированные стабы)

Для команд, предпочитающих собственный gRPC-стек: контракт — это decoded.proto (наш маршрут, shredstream.com.DecodedShredStreamService), который импортирует свои сообщения из shreder_binary.proto — оба файла поставляются с каждым официальным клиентом. Генерация стабов только из shreder_binary.proto даёт совместимый «из коробки» маршрут, используемый в примерах ниже.

Rust (tonic)

rust
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![],
},
)]);
// Статические фильтры: отправляем один запрос, затем half-close (итератор завершается).
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)?;
// → маршрутизируйте по resp.filters, обрабатывайте vtx …
}
}

Наш эквивалентный маршрут — DecodedShredStreamServiceClient::subscribe_decoded_transactions, с точно теми же типами запроса/ответа. Чтобы обновлять фильтры на лету, замените tokio_stream::iter на канал, который держите открытым, и send-ите в него новый запрос — переподключение не требуется.

TypeScript (@grpc/grpc-js)

ts
import { Metadata } from "@grpc/grpc-js";
// … клиент, сгенерированный из 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, если фильтры не меняются
call.on("data", (resp) => {
const raw = resp.transaction?.transaction?.binaryTransaction; // Buffer
// десериализуйте VersionedTransaction на стороне клиента …
});
call.on("error", (e) => { /* DATA_LOSS / UNAUTHENTICATED → переподключение / лог */ });

Python (grpcio)

python
import grpc
import shreder_binary_pb2 as pb, shreder_binary_pb2_grpc as rpc
chan = 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 # байты bincode
# десериализуйте VersionedTransaction …

Быстрый тест с grpcurl

bash
grpcurl -plaintext -proto shreder_binary.proto \
-H 'x-token: <TOKEN>' \
-d '{"transactions":{"all":{}}}' \
<host>:50051 shreder_binary.ShrederBinaryService/SubscribeBinaryTransactions

🧯 Обработка ошибок и переподключение

Статус gRPCЗначениеЧто делать
UNAUTHENTICATEDтокен отсутствует / недействителен / не тот потокисправьте токен или поток — не повторяйте вслепую
INVALID_ARGUMENTнекорректный base58, слишком много pubkeys, слишком много именованных фильтровисправьте проблемный фильтр
DATA_LOSSваш потребитель слишком медленный — очередь отправки сервера переполнилась и стрим закрытпереподключитесь и заново отправьте вашу map фильтров
PERMISSION_DENIEDотключение/отзыв операторомне зацикливайтесь; свяжитесь с оператором

Сервис работает только в реальном времени — без повтора и дозагрузки. При DATA_LOSS или обрыве транспорта переподключайтесь с экспоненциальным backoff (и небольшим джиттером), затем заново отправьте вашу map фильтров; короткий пропуск данных во время переподключения — это нормально. Из-за правила «одно соединение на токен» убедитесь, что на токен работает только один цикл переподключения.

Чтобы вообще не допускать DATA_LOSS, вычитывайте стрим не блокируясь: передавайте тяжёлую обработку (включая десериализацию bincode) в очередь или пул воркеров, чтобы ваш цикл приёма никогда не застревал. gRPC поверх HTTP/2 гарантирует порядок и целостность всего отправленного — единственная возможная потеря — это сброс «клиент слишком медленный», и он всегда явно сигнализируется через DATA_LOSS.


⚖️ gRPC или UDP?

gRPCUDP
ЗадержкаНизкая — h2c, без рукопожатия TLSМинимальная — бинарные датаграммы, без соединения
ДоставкаУпорядоченная, надёжная; «клиент слишком медленный» сигнализируется через DATA_LOSSPush по принципу best-effort; потери возможны, обнаруживаются по пропускам seq
ФильтрацияФильтры по аккаунтам на стороне сервера (статические ключи)Нет — фильтруйте на стороне клиента после получения
ДоступностьИсходящее соединение — работает за NAT/файрволамиТребует публично доступной конечной точки UDP
Полезная нагрузкаVersionedTransaction (bincode)VersionedTransaction (bincode), в конверте из 16 байт

Выбирайте gRPC, когда нужны серверная фильтрация, упорядоченная доставка и совместимость с NAT. Выбирайте UDP ради абсолютно минимальной задержки, если у вас есть доступная конечная точка и вы сами отслеживаете потери. Полезная нагрузка в обоих случаях идентична.


➡️ Дальнейшие шаги

Доставка по gRPC — Docs | ShredStream.com