Cheapest reachable link first, buffering to a store when every link is down. One capability of pamoja, one memory-safe Rust core with bindings for TypeScript, Python, and C#.
npm install @pamoja/ladder
This pulls in @pamoja/native, the compiled engine. npm install pamoja is the whole framework in one package.
The test that runs in CI, spliced here as it ran.
From bindings/node/guides/ladder.ts:
import { Transport } from '@pamoja/core'
import { Delivery, Ladder } from '@pamoja/ladder'
import { LoopbackBroker } from '@pamoja/loopback'
import { Store } from '@pamoja/sync'
const TOPIC = 'sensors/1/temperature'
async function main() {
// Two links off the same node: a near mesh hop and a metered backhaul. Each is a
// separate broker, so which one carried a reading is visible from its subscriber.
const mesh = new LoopbackBroker()
const backhaul = new LoopbackBroker()
const gateway = backhaul.link()
await gateway.connect()
await gateway.subscribe(TOPIC)
// Rungs are tried in the order they are added, cheapest first. The mesh hop loses every
// packet here; the backhaul carries one send, then drops the next two.
const ladder = new Ladder(Store.memory())
await ladder.rung(Transport.degraded(mesh.rung(), { dropEvery: 1 }))
await ladder.rung(Transport.degraded(backhaul.rung(), { up: 1, down: 2 }))
await ladder.connect()
// The mesh hop refuses, so the reading goes out over the backhaul and arrives on the
// broker only that rung publishes to.
const first = await ladder.send(TOPIC, Buffer.from('21.5'))
const arrived = (await gateway.recv())!
console.log(`first reading: ${first}, gateway got ${arrived.payload.toString()}`)
// Now nothing will take a send, so the next reading is buffered rather than lost.
const second = await ladder.send(TOPIC, Buffer.from('21.6'))
const waiting = await ladder.buffered()
console.log(`second reading: ${second}, ${waiting} waiting in the queue`)
// A flush while the links are still down forwards nothing and leaves the backlog
// intact, because a record is removed only once a rung has accepted it.
const whileDown = await ladder.flush()
console.log(`flush while down forwarded ${whileDown}, queue still ${await ladder.buffered()}`)
// The backhaul is reachable again, so the buffered reading goes out exactly once.
const whenUp = await ladder.flush()
const late = (await gateway.recv())!
console.log(`flush when up forwarded ${whenUp}, gateway got ${late.payload.toString()}`)
return { first, second, waiting, whileDown, whenUp, left: await ladder.buffered(), late }
}
main()
| Language | Package | Reference |
|---|---|---|
| Rust | pamoja-ladder |
reference, docs.rs, install |
| TypeScript | @pamoja/ladder |
reference, install |
| Python | pamoja-ladder |
reference, install |
| C# | Pamoja.Ladder |
reference, install |
@pamoja/ladder reference, every class, function, and type this package exports.MIT
Ergonomic facade over the generated ladder binding.
A ladder is the answer to a node with more than one way to reach the network and no single one that always works: rungs are tried in the order they were added, cheapest first, and a message no rung accepts goes into a buffer rather than being lost.
The delivery outcome is re-exported as a runtime Delivery object, because the generated enum is types-only.