@pamoja/ladder
v0.1.18
Published
Cheapest reachable link first, buffering to a store when every link is down.
Readme
@pamoja/ladder
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#.
Install
npm install @pamoja/ladderThis pulls in @pamoja/native, the compiled engine. npm install pamoja is the whole framework in one package.
Example
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, '21.5')
const arrived = (await gateway.recv())!
console.log(`first reading: ${first}, gateway got ${arrived.text!}`)
// Now nothing will take a send, so the next reading is buffered rather than lost.
const second = await ladder.send(TOPIC, '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.text!}`)
// The ladder is a link both ways. A subscription placed on it goes onto every rung that
// listens, and a receive takes whichever rung delivers, so a command reaches the node
// over whatever link is up. This one comes back over the backhaul.
await ladder.subscribe('actuators/1/valve')
await gateway.send('actuators/1/valve', 'open')
const command = (await ladder.recv())!
console.log(`command back over the ladder: ${command.text!}`)
const left = await ladder.buffered()
return { first, second, waiting, whileDown, whenUp, left, late, command }
}
main()The same capability in every language
| 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 |
Documentation
@pamoja/ladderreference, every class, function, and type this package exports.- The Transport ladder guide, with the same example in Rust, Python, and C#.
- Every capability, and the install page.
License
MIT
