🔌 التسليم عبر gRPC
يمكن استهلاك Decoded Shred Stream عبر gRPC بدلاً من UDP: بث صادر قائم على الاتصال يضيف تصفية حسابات على جانب الخادم وتسليماً مرتّباً عبر HTTP/2 فوق المعاملات نفسها بزمن استجابة دون المللي ثانية — دون الحاجة إلى نقطة UDP متاحة للعموم. تعرض لوحة التحكم نقطة نهاية gRPC ورمز الوصول (token) لكل بث.
- نقطة النهاية —
<host>:50051(منفذ واحد لكل التدفقات). - النقل — gRPC عبر HTTP/2، بنص صريح h2c (بلا TLS). هذا خيار زمن استجابة مقصود: لا مصافحة TLS تدفع ثمنها. اعزل الوصلة على مستوى الشبكة — شبكة خاصة أو VPN أو peering — ولا تعرضها أبداً للإنترنت المفتوح.
- المسارات —
shredstream.com.DecodedShredStreamService/SubscribeDecodedTransactions، أو المسار المتوافق كبديل مباشرshreder_binary.ShrederBinaryService/SubscribeBinaryTransactions. كلاهما متكافئ تماماً (نفس رسائل protobuf، محتوى مطابق بايتاً ببايت) — العميل القادم من نقطة نهاية Shreder / Raiden Pulse لا يحتاج سوى تغيير العنوان.
🔑 المصادقة
يجب أن يحمل كل نداء رمزك في الـ metadata الخاص بـ gRPC، بـإحدى الصيغتين التاليتين (كلتاهما مقبولة):
authorization: Bearer <TOKEN> x-token: <TOKEN>
يُرفض الرمز المفقود أو غير الصالح بحالة موحّدة UNAUTHENTICATED مع الرسالة authentication refused. تحقّق من رمزك.
تذكّر: اتصال واحد لكل رمز. فتح بث ثانٍ بالرمز نفسه يطرد الأقدم.
🧭 نموذج الاشتراك
المسار بث ثنائي الاتجاه (stream request → stream response):
- افتح البث وأرسل طلباً واحداً على الأقل يحمل map من الفلاتر المسمّاة:
{ "<name>": <filter>, … }. - يسلّم الخادم فقط المعاملات التي تطابق فلتراً واحداً على الأقل. كل استجابة موسومة باسم (أو أسماء) الفلتر (الفلاتر) الذي حققته (الحقل
filters) — فتوجّه عدة استراتيجيات عبر بث واحد. - يمكنك إرسال map جديدة في أي وقت: تحلّ محلّ السابقة على الساخن، دون إعادة اتصال ودون فجوة في التدفق.
الحالات الحدّية:
- map فارغة (
{}) ⇒ لا يُسلَّم شيء (أنت تشترك في فلاتر، لا في firehose). - فلتر مسمّى فارغ (
{ "all": {} }) ⇒ تمرّ كل معاملة، موسومة بـ"all". - بعد إرسال فلاترك لا يبقى لديك ما ترسله: يواصل البث التسليم بهذه الفلاتر حتى ترسل فلاتر جديدة أو تغلق الاتصال.
🎯 فلاتر الحسابات
كل فلتر مسمّى هو ثلاث قوائم من المفاتيح العامة base58، مجموعة بـ**«و» المنطقية (AND)**:
| الحقل | المعنى |
|---|---|
account_include | يجب أن تلمس المعاملة حساباً واحداً على الأقل من هذه الحسابات (فارغة = بلا قيد) |
account_exclude | يجب ألا تلمس المعاملة أياً من هذه الحسابات |
account_required | يجب أن تلمس المعاملة كل هذه الحسابات |
بالنسبة إلى المعاملات، تُقيَّم «الملامسة» على مفاتيح الحسابات الثابتة (static account keys) للمعاملة (بما فيها الموقّعون). العناوين المُحلَّلة عبر جداول البحث عن العناوين (Address Lookup Tables — ALTs) ليست موجودة في مغلّف المعاملة ومن ثمّ لا يمكن التصفية عليها — صفِّ على المفاتيح الثابتة فقط.
📦 الحمولة
تحمل كل استجابة المعاملة بايتات خام مع الـ slot الخاص بها:
message SubscribeUpdateBinaryTransaction {BinaryTransaction transaction = 1;uint64 slot = 2; // slot في Solana}message BinaryTransaction {repeated bytes signatures = 1; // توقيعات (64 بايت لكل واحد)bytes binary_transaction = 3; // VersionedTransaction, bincode, VERBATIM}
binary_transaction هو تنسيق wire القياسي في Solana — VersionedTransaction مُسلسلة بـ bincode، دون أي تعديل. فكّ تسلسلها بأي SDK Solana ومرّرها إلى محلّلك الحالي، تماماً كما في وضع UDP. معاملات التصويت (vote) مُستبعدة افتراضياً.
💻 العميل الرسمي
يتولّى عملاء decoded-shredstream الرسميون هذا المسار نيابةً عنك — الاتصال، والرمز في الـ metadata، وmap الفلاتر، وإعادة الاتصال مع إعادة إرسال الفلاتر الحالية، وتسليم المعاملة مع الـ slot والتوقيعات الخاصة بها. الحزمة نفسها لـ UDP و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)})
يمكن استبدال الفلاتر على بث حيّ (updateFilters / update_filters) دون إعادة اتصال. الانقطاعات القابلة للاستعادة لا تصلك أبداً — يعيد العميل الاتصال ويعيد إرسال map الفلاتر الحالية؛ لا يُنهي البث سوى رمز مرفوض، أو جلسة أغلقها الخادم، أو map فلاتر مرفوضة، وذلك عبر خطأ تعالجه مرة واحدة.
💻 عملاء gRPC القياسيون (stubs مولّدة)
للفرق التي تفضّل مكدّس gRPC الخاص بها: العقد هو decoded.proto (مسارنا، shredstream.com.DecodedShredStreamService)، والذي يستورد shreder_binary.proto من أجل رسائله — ويُرفق الملفان مع كل عميل رسمي. توليد stubs من shreder_binary.proto وحده يمنحك المسار المتوافق كبديل مباشر المستخدم في الأمثلة أدناه.
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![],},)]);// فلاتر ثابتة: نرسل طلباً واحداً ثم نغلق النصف المُرسِل 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)
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)
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 # بايتات bincode# فكّ تسلسل VersionedTransaction …
اختبار سريع بـ grpcurl
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 | فصله/إلغاؤه من قبل المشغّل | لا تُكرّر في حلقة؛ تواصل مع المشغّل |
الخدمة في الوقت الفعلي فقط — لا إعادة تشغيل (replay) ولا backfill. عند DATA_LOSS أو انقطاع النقل، أعد الاتصال بتراجع أُسّي (exponential backoff) مع قليل من العشوائية (jitter)، ثم أعد إرسال map الفلاتر؛ فجوة بيانات قصيرة أثناء إعادة الاتصال أمر متوقع. وبسبب قاعدة اتصال واحد لكل رمز، تأكد من أن حلقة إعادة اتصال واحدة فقط تعمل لكل رمز.
لتجنّب DATA_LOSS من الأساس، استهلك البث دون حجب: أحِل المعالجة الثقيلة (بما فيها فكّ تسلسل bincode) إلى طابور أو مجمّع عمّال (worker pool) كي لا تتوقف حلقة الاستقبال أبداً. يضمن gRPC عبر HTTP/2 ترتيب وسلامة كل ما يُبَثّ — الفقد الوحيد الممكن هو إسقاط «العميل بطيء جداً»، وهو دائماً مُشار إليه صراحةً بـ DATA_LOSS.
⚖️ gRPC أم UDP؟
| gRPC | UDP | |
|---|---|---|
| زمن الاستجابة | منخفض — h2c، بلا مصافحة TLS | الأدنى — datagrams ثنائية، بلا اتصال |
| التوصيل | مرتّب وموثوق؛ «العميل بطيء جداً» يُشار إليه بـ DATA_LOSS | دفع بأفضل جهد؛ الفقد ممكن، يُكتشف عبر فجوات seq |
| التصفية | فلاتر حسابات على جانب الخادم (مفاتيح ثابتة) | لا شيء — صفِّ على جانب العميل بعد الاستقبال |
| القابلية للوصول | اتصال صادر — يعمل خلف NAT/الجدران النارية | يتطلب نقطة UDP متاحة للعموم |
| الحمولة | VersionedTransaction (bincode) | VersionedTransaction (bincode)، داخل مغلّف من 16 بايت |
اختر gRPC عندما تريد التصفية على جانب الخادم والتوصيل المرتّب والاتصال المتوافق مع NAT. واختر UDP لأدنى زمن استجابة على الإطلاق حين تملك نقطة قابلة للوصول وتراقب الفقد بنفسك. الحمولة متطابقة في الحالتين.
➡️ الخطوات التالية
- استقبال المعاملات المفكوكة — مغلّف UDP، والإزاحات، واكتشاف الفجوات، ومفكك ترميز بسيط.
- Decoded Shred Stream — الموقع، وزمن الاستجابة، وأوضاع التوصيل.