🔌 gRPC 交付

Decoded Shred Stream 可以通过 gRPC 消费,以替代 UDP:一条出站、面向连接的流,在同样亚毫秒级延迟的交易之上,额外提供服务器端账户过滤有序的 HTTP/2 交付——无需可公网访问的 UDP 端点。你的仪表盘会显示每个流的 gRPC 端点访问令牌

  • 端点 —— <host>:50051(所有流共用一个端口)。
  • 传输 —— gRPC over HTTP/2,明文 h2c(无 TLS)。这是有意为之的延迟取舍:没有 TLS 握手开销需要付出。请在网络层隔离该链路——私有网络、VPN 或 peering——绝不要将其暴露到公网。
  • 路由 —— 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,无需重连,流也不会出现空档。

边界情形:

  • 空 map{})⇒ 不投递任何内容(你订阅的是过滤器,而不是 firehose)。
  • 空的命名过滤器{ "all": {} })⇒ 所有交易通过,标注为 "all"
  • 过滤器发送之后,你无需再发送任何内容:流会一直按这些过滤器投递,直到你发送新的过滤器或关闭连接。

🎯 账户过滤器

每个命名过滤器由三个 base58 公钥列表组成,以逻辑**与(AND)**组合:

字段含义
account_include交易必须触及这些账户中的至少一个(空 = 无约束)
account_exclude交易必须不触及这些账户中的任何一个
account_required交易必须触及这些账户中的全部

对于交易,"触及"针对交易的静态账户键(static account keys)判定(含签名者)。通过 Address Lookup Tables (ALT) 解析出的地址存在于交易信封中,因此无法用于过滤——请只对静态键过滤。


📦 载荷

每条响应以原始字节形式携带交易,并附带其 slot:

proto
message SubscribeUpdateBinaryTransaction {
BinaryTransaction transaction = 1;
uint64 slot = 2; // Solana slot
}
message BinaryTransaction {
repeated bytes signatures = 1; // 签名(每个 64 字节)
bytes binary_transaction = 3; // VersionedTransaction,bincode,VERBATIM(原样)
}

binary_transactionSolana 标准线格式——一个 bincode 序列化的 VersionedTransaction,未经改动。用任何 Solana SDK 反序列化,再喂给你现有的解析器,与 UDP 模式完全一致。投票交易默认排除。


💻 官方客户端

官方 decoded-shredstream 客户端会替你对接这条路由——连接、令牌元数据、过滤器 map、重连并重新发送当前过滤器,以及附带 slot 和签名的交易。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 客户端(生成的 stub)

对于偏好使用自有 gRPC 技术栈的团队:契约是 decoded.proto(我们的路由,shredstream.com.DecodedShredStreamService),它导入 shreder_binary.proto 以获取消息定义——这两个文件随每个官方客户端一同分发。仅从 shreder_binary.proto 生成 stub,即可得到下文所用的、可直接替换的兼容路由。

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_ARGUMENTbase58 无效、pubkey 过多、命名过滤器过多修正出错的过滤器
DATA_LOSS你的消费者太慢——服务器发送队列溢出,流已被关闭重连并重新发送你的过滤器 map
PERMISSION_DENIED被运营方断开 / 吊销不要循环重试;请联系运营方

该服务仅提供实时数据——没有回放或补拉(backfill)。遇到 DATA_LOSS 或传输中断时,请以指数退避(并加入少量抖动 jitter)重连,然后重新发送你的过滤器 map;重连期间出现短暂的数据空档是正常的。由于每个令牌只允许一条连接,请确保每个令牌只运行一个重连循环。

要从源头上避免 DATA_LOSS,请不阻塞地消费流:把繁重处理(包括 bincode 反序列化)交给队列或工作线程池,让你的接收循环永不停顿。gRPC over HTTP/2 保证所有已发出内容的顺序与完整性——唯一可能的丢失就是"客户端太慢"这一情形,且它总是通过 DATA_LOSS 显式发出信号。


⚖️ gRPC 还是 UDP?

gRPCUDP
延迟低——h2c,无 TLS 握手最低——二进制数据报,无连接
交付有序、可靠;"客户端太慢"通过 DATA_LOSS 发出信号尽力而为的推送;可能丢包,通过 seq 缺口检测
过滤服务器端账户过滤(静态键)无——接收后在客户端过滤
可达性出站连接——在 NAT/防火墙之后也能工作需要可公网访问的 UDP 端点
载荷VersionedTransaction(bincode)VersionedTransaction(bincode),位于 16 字节信封中

当你需要服务器端过滤、有序交付以及 NAT 友好的连接方式时,选择 gRPC;当你拥有可达端点、并自行监控丢包,且追求绝对最低延迟时,选择 UDP。两种方式的载荷完全相同。


➡️ 后续步骤

gRPC 交付 — Docs | ShredStream.com