Skip to content
Docs

Reference shred receiver and deshredder

Provide one annotated control flow for a production deshredder built against a pinned Agave release.

Before you start

  • All receiver and shred pages
  • A pinned Agave checkout
  • A prefetched leader schedule
  • Optional slot-correct lookup-table state

Pin the protocol adapter first

Build the receiver against the Agave revision used by the network. The public solana-ledger surface parses, sanitizes, verifies, reads data, and deshreds. Current Agave keeps erasure-slice access, coding counts, local erasure index, and Merkle recovery behind crate-private methods. A correct external receiver needs a thin reviewed adapter in a pinned fork or must insert shreds through Blockstore.

A stable standalone recovery API and the active entry-envelope revision are not currently specified by the product contract. Do not recreate private offsets. Export this adapter from the pinned shred module:

rust

pub struct RecoveryMeta {
    pub position: usize,
    pub data_shards: Option<usize>,
    pub coding_shards: Option<usize>,
}

pub fn recovery_meta(shred: &Shred) -> Result<RecoveryMeta, Error> {
    let (data_shards, coding_shards) = if shred.is_code() {
        (Some(usize::from(shred.num_data_shreds()?)),
         Some(usize::from(shred.num_coding_shreds()?)))
    } else {
        (None, None)
    };
    Ok(RecoveryMeta {
        position: shred.erasure_shard_index()?,
        data_shards,
        coding_shards,
    })
}

pub fn recover_verified(
    shreds: Vec<Shred>, cache: &ReedSolomonCache,
) -> Result<Vec<Shred>, Error> {
    merkle::recover(shreds, cache)?.collect()
}

Place the wrapper inside ledger/src/shred.rs, where the private methods are visible. Keep the patch small, versioned, and covered by upstream recovery tests.

Use one explicit pipeline

The application below shows every state transition. It omits metrics plumbing and bounded-channel declarations to keep the protocol path readable. leaders, lookup_resolver, and sink are prefetched dependencies. No network request occurs in handle.

rust

use solana_entry::block_component::BlockComponent;
use solana_ledger::shred::{
    recover_verified, recovery_meta, ReedSolomonCache, Shred, Shredder,
};
use solana_pubkey::Pubkey;
use solana_transaction::versioned::VersionedTransaction;
use std::{collections::{BTreeMap, HashMap}, io, net::{Ipv4Addr, SocketAddr, UdpSocket}};

#[derive(Clone, Copy, Eq, Hash, PartialEq)]
struct Id { slot: u64, index: u32, code: bool }
#[derive(Clone, Copy, Eq, Hash, PartialEq)]
struct FecKey { slot: u64, start: u32 }

struct Set {
    members: BTreeMap<usize, Shred>,
    k: Option<usize>,
    m: Option<usize>,
}
#[derive(Default)]
struct SlotBuffer { next: u32, pending: BTreeMap<u32, Shred>, range: Vec<Shred> }

trait Leaders { fn leader(&self, slot: u64) -> Option<Pubkey>; }
trait Lookups {
    type Resolved;
    fn resolve(&self, slot: u64, tx: &VersionedTransaction) -> Result<Self::Resolved, String>;
}
trait Sink<R> { fn proposed(&mut self, slot: u64, tx: VersionedTransaction, keys: Result<R, String>); }

struct Decoder<L, A, O> {
    leaders: L,
    lookups: A,
    sink: O,
    seen: HashMap<Id, Shred>,
    sets: HashMap<FecKey, Set>,
    slots: HashMap<u64, SlotBuffer>,
    rs: ReedSolomonCache,
}

impl<L: Leaders, A: Lookups, O: Sink<A::Resolved>> Decoder<L, A, O> {
    fn handle(&mut self, bytes: &[u8]) -> Result<(), String> {
        // 1. Parse and sanitize exactly the received datagram.
        let shred = Shred::new_from_serialized_shred(bytes.to_vec())
            .map_err(|e| e.to_string())?;
        shred.sanitize().map_err(|e| e.to_string())?;

        // 2. Resolve the scheduled leader and verify proof plus signature.
        let leader = self.leaders.leader(shred.slot())
            .ok_or_else(|| "leader missing for slot".to_owned())?;
        if shred.verify(&leader) == false { return Err("leader signature invalid".into()); }

        // 3. Remove repeats and reject different content at one identity.
        let id = Id { slot: shred.slot(), index: shred.index(), code: shred.is_code() };
        if let Some(old) = self.seen.get(&id) {
            return if old.is_shred_duplicate(&shred) {
                Err("conflicting shred identity".into())
            } else { Ok(()) };
        }
        self.seen.insert(id, shred.clone());

        // 4. Group one authenticated member into its explicit FEC set.
        let key = FecKey { slot: shred.slot(), start: shred.fec_set_index() };
        let meta = recovery_meta(&shred).map_err(|e| e.to_string())?;
        let set = self.sets.entry(key).or_insert_with(|| Set {
            members: BTreeMap::new(), k: None, m: None,
        });
        if let (Some(k), Some(m)) = (meta.data_shards, meta.coding_shards) {
            if set.k.is_some_and(|old| (old == k) == false)
                || set.m.is_some_and(|old| (old == m) == false) {
                return Err("conflicting erasure configuration".into());
            }
            set.k = Some(k); set.m = Some(m);
        }
        if set.members.contains_key(&meta.position) {
            return Err("duplicate erasure position with a different identity".into());
        }
        set.members.insert(meta.position, shred);

        // 5. Wait for k unique members, recover missing data, then finish this set.
        let Some(k) = set.k else { return Ok(()) };
        if set.members.len() < k { return Ok(()); }
        let have_data = set.members.values().filter(|s| s.is_data()).count();
        let mut data: Vec<Shred> = if have_data == k {
            set.members.values().filter(|s| s.is_data()).cloned().collect()
        } else {
            let input = set.members.values().cloned().collect();
            let recovered = recover_verified(input, &self.rs)
                .map_err(|e| e.to_string())?;
            set.members.values().chain(recovered.iter())
                .filter(|s| s.is_data()).cloned().collect()
        };
        data.sort_unstable_by_key(Shred::index);
        data.dedup_by_key(|s| s.index());
        if (data.len() == k) == false { return Err("recovery did not produce every data shard".into()); }
        self.sets.remove(&key);

        // 6. Insert data into the per-slot consecutive ordering buffer.
        for shred in data { self.order(shred)?; }
        Ok(())
    }

    fn order(&mut self, shred: Shred) -> Result<(), String> {
        let slot = shred.slot();
        let ready = {
            let state = self.slots.entry(slot).or_default();
            if shred.index() < state.next { return Ok(()); }
            state.pending.insert(shred.index(), shred);
            let mut ready = Vec::new();
            while let Some(shred) = state.pending.remove(&state.next) {
                state.next = state.next.checked_add(1).ok_or("data index overflow")?;
                let complete = shred.data_complete();
                state.range.push(shred);
                if complete { ready.push(std::mem::take(&mut state.range)); }
            }
            ready
        };
        for range in ready {
            self.decode(slot, range)?;
        }
        Ok(())
    }

    fn decode(&mut self, slot: u64, range: Vec<Shred>) -> Result<(), String> {
        // 7. Reassemble ledger bytes and decode the release-matched component.
        let bytes = Shredder::deshred(range.iter().map(Shred::payload))
            .map_err(|e| e.to_string())?;
        let component: BlockComponent = wincode::deserialize(&bytes)
            .map_err(|e| e.to_string())?;
        let BlockComponent::EntryBatch(entries) = component else { return Ok(()) };

        // 8. Extract every versioned transaction. Preserve lookup references,
        // and attach slot-labelled resolution when the state provider can do it.
        for entry in entries {
            for tx in entry.transactions {
                let keys = self.lookups.resolve(slot, &tx);
                self.sink.proposed(slot, tx, keys);
            }
        }
        Ok(())
    }
}

fn receive(mut handle: impl FnMut(&[u8]) -> Result<(), String>, bind: &str) -> io::Result<()> {
    let socket = UdpSocket::bind(bind)?;
    let allowed = Ipv4Addr::new(64, 130, 40, 90);
    let mut buf = [0u8; 1228];
    loop {
        let (len, peer) = socket.recv_from(&mut buf)?;
        let allowed_source = match peer {
            SocketAddr::V4(v4) => *v4.ip() == allowed,
            SocketAddr::V6(_) => false,
        };
        if allowed_source {
            if let Err(error) = handle(&buf[..len]) { let _sampled_error = error; }
        }
    }
}

Harden the reference before production

Move receive, verification, FEC ownership, ordering, decoding, and output into bounded stages. The synchronous example makes control flow visible but can drop packets while recovery or output runs. Add SO_RCVBUF, SO_RXQ_OVFL, recvmmsg, metrics, state TTLs, maximum active slots and sets, worker supervision, and sampled logging from the receiver pages.

Validate every set member's local position and compatible root in the adapter. The reference relies on authenticated typed shreds and upstream recovery checks, but production state should reject conflicting configuration before it reaches recovery.

The lookup resolver must preserve v0 order: static keys, all loaded writable keys, then all loaded readonly keys. If slot-correct table state is unavailable, return a labelled unresolved result and still emit the original VersionedTransaction. Never block UDP receipt on RPC.

Treat every emitted transaction as proposed. Reconcile it with later confirmed state. Replay captured fixtures after every Agave, codec, compiler, or adapter change.

Parameters

NameTypeDefaultNotes
bindSocketAddrV4noneVerified destination address and UDP port.
sourceIpv4Addr64.130.40.90Application and firewall source allowlist.
agave_revisionstringnoneExact parser, recovery, deshred, entry, and transaction implementation.
leader_cacheLeadersnonePrefetched slot-to-leader keys for the configured cluster.
lookup_resolverLookupsnoneNonblocking slot-labelled ALT state, or an unresolved-result implementation.

When it goes wrong

leader missing: slot=N

Cause. The prefetched leader schedule does not cover the observed slot.

Fix. Refresh ahead of epoch boundaries and keep a bounded pending path outside the receive loop.

recovery: Invalid Merkle root

Cause. Members, reconstructed bytes, or adapter logic do not rebuild the authenticated set root.

Fix. Discard the set and compare member roots, configuration, and pinned Agave recovery code.

entry decode: codec error

Cause. The range boundary, bytes, or configured entry envelope does not match the network release.

Fix. Verify consecutive completion flags and roll back to the tested codec revision.

lookup resolution unavailable

Cause. Slot-correct address lookup table state is not present yet.

Fix. Emit the original versioned transaction with unresolved lookup references and reconcile later.

Questions

Why does the reference require a pinned Agave adapter?
Current Agave keeps variant-dependent erasure slices, coding metadata, and recovery internals crate-private. Copying guessed offsets would be unsafe. A small versioned wrapper inside the pinned ledger crate reuses the same reconstruction, sanitization, Merkle-tree, and proof logic as the validator implementation.
Is the synchronous reference loop production-ready?
No. It documents the complete protocol control flow in one place. Production code must isolate UDP receipt from cryptography, recovery, decoding, lookup resolution, and output with bounded queues. It also needs socket overflow reporting, batching, state eviction, metrics, and supervised workers.
What happens when address lookup tables cannot be resolved yet?
Emit the original VersionedTransaction immediately with its static keys and lookup references, plus a labelled unresolved result. Resolve later against an explicit state context. Do not block the receive path or claim that latest confirmed table state is exact for an unconfirmed proposed fork.