feat(cassettes): consume operator operations instead of counts
The machine now owns its bay counts outright. The operator publishes what it did — a refill in notes added, an empty, a recount, a denomination change — and this process applies it to the total it already holds. Both sides used to write the same value over a transport that never tells a writer it lost. Addressable events order by created_at at second granularity with ties broken on event id, and a relay returns OK for an event it then discards, so a dashboard form loaded before a dispense silently discarded that dispense and neither side could detect it. A value with one writer cannot be clobbered. Schema v13 adds cassette_ops, the dedup ledger. A delta applied twice is wrong and addressable events are re-delivered on every reconnect, so the operator mints an id per operation and this table records the ones applied. That also retires the created_at watermark on this path: it was the only replay defence under absolute counts, but it drops an out-of-order event whole, operations included, where per-op ids let the unseen ones through and no-op the rest. A window is applied oldest-first by `at`, ties broken by id, in one transaction with the count mutation. A recount then a refill is not the same as the reverse, and a crash mid-apply must roll back to a coherent count rather than a partial one. A malformed op or one naming a bay this machine does not have is neither applied nor recorded, so it stays pending on the operator's dashboard. That is the honest outcome. Recording it as applied would stop the noise by telling the operator their refill landed. The state document gains applied_ops, seq and schema_version. applied_ops is the acknowledgement leg — echoing the ids back is the only way the operator can tell an operation that landed from one merely sent. seq is bumped on every local count change from any cause, so a reader can reject a regression without trusting either clock.
This commit is contained in:
parent
9077f9c299
commit
c889f7f0df
6 changed files with 537 additions and 165 deletions
|
|
@ -15,7 +15,7 @@ import fs from 'node:fs'
|
|||
|
||||
let db: Database.Database | null = null
|
||||
|
||||
const SCHEMA_VERSION = '12'
|
||||
const SCHEMA_VERSION = '13'
|
||||
|
||||
function getDbPath(): string {
|
||||
const prodDir = '/var/lib/bitspire'
|
||||
|
|
@ -57,6 +57,17 @@ export function initDatabase(dbPath?: string): void {
|
|||
count INTEGER NOT NULL DEFAULT 0
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS cassette_ops (
|
||||
id TEXT PRIMARY KEY,
|
||||
position INTEGER NOT NULL,
|
||||
op_type TEXT NOT NULL,
|
||||
bills INTEGER,
|
||||
count INTEGER,
|
||||
denomination INTEGER,
|
||||
op_at INTEGER NOT NULL,
|
||||
applied_at INTEGER NOT NULL
|
||||
);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS cashbox (
|
||||
id INTEGER PRIMARY KEY CHECK (id = 1),
|
||||
total_bills INTEGER NOT NULL DEFAULT 0,
|
||||
|
|
@ -373,12 +384,47 @@ export function initDatabase(dbPath?: string): void {
|
|||
console.log('[StateStore] Migrated schema v11 → v12 (bunker_binding transport config)')
|
||||
}
|
||||
|
||||
if (existing && existing.value === '12') {
|
||||
// Migration v12 → v13: operator OPERATIONS replace operator counts
|
||||
// (aiolabs/bitspire ADR-004).
|
||||
//
|
||||
// The operator used to publish absolute counts and this machine applied
|
||||
// them outright. Both sides wrote the same value over a transport that
|
||||
// never tells a writer it lost, so a dashboard form loaded before a
|
||||
// dispense silently discarded that dispense — and nothing on either side
|
||||
// could detect it afterwards. The operator now publishes what it DID and
|
||||
// this machine, which holds the notes, owns the running total.
|
||||
//
|
||||
// `cassette_ops` is the dedup ledger. A delta applied twice is wrong, and
|
||||
// addressable events are re-delivered on every reconnect, so the operator
|
||||
// mints an id per operation and we record the ones we have applied. The
|
||||
// operator's window is a slice of recent operations rather than just the
|
||||
// newest, so one we missed arrives with the next publish; dedup is what
|
||||
// makes re-delivery free instead of dangerous.
|
||||
db.exec(`
|
||||
CREATE TABLE IF NOT EXISTS cassette_ops (
|
||||
id TEXT PRIMARY KEY,
|
||||
position INTEGER NOT NULL,
|
||||
op_type TEXT NOT NULL,
|
||||
bills INTEGER,
|
||||
count INTEGER,
|
||||
denomination INTEGER,
|
||||
op_at INTEGER NOT NULL,
|
||||
applied_at INTEGER NOT NULL
|
||||
);
|
||||
`)
|
||||
db.prepare('UPDATE meta SET value = ? WHERE key = ?').run('13', 'schema_version')
|
||||
console.log('[StateStore] Migrated schema v12 → v13 (added cassette_ops)')
|
||||
existing.value = '13'
|
||||
}
|
||||
|
||||
// Defensive: a fresh install at SCHEMA_VERSION skips all migrations.
|
||||
// Seed the operator-config meta rows if they're missing (idempotent).
|
||||
const seedMeta = db.prepare('INSERT OR IGNORE INTO meta (key, value) VALUES (?, ?)')
|
||||
seedMeta.run('lastKnownConfigCreatedAt', '0')
|
||||
seedMeta.run('bootstrapPublishedAt', '')
|
||||
seedMeta.run('lastKnownFeeConfigCreatedAt', '0')
|
||||
seedMeta.run('cassetteStateSeq', '0')
|
||||
|
||||
const cashboxRow = db.prepare('SELECT id FROM cashbox WHERE id = 1').get()
|
||||
if (!cashboxRow) {
|
||||
|
|
@ -471,6 +517,36 @@ export function clearCountsUncertain(): void {
|
|||
).run('countsUncertainSince', '')
|
||||
}
|
||||
|
||||
/**
|
||||
* A counter bumped on every local change to a bay count, from any cause.
|
||||
*
|
||||
* It rides along in the state document so a reader can reject a regression
|
||||
* without trusting a clock. `created_at` cannot carry that: it has
|
||||
* second granularity, so two publishes in the same second are ordered by
|
||||
* whichever event id hashes lower — and a machine whose clock stepped
|
||||
* backwards would otherwise have every later report look older than the one
|
||||
* already on the relay.
|
||||
*/
|
||||
export function getCassetteStateSeq(): number {
|
||||
if (!db) throw new Error('Database not initialized')
|
||||
const row = db.prepare('SELECT value FROM meta WHERE key = ?').get('cassetteStateSeq') as
|
||||
| { value: string }
|
||||
| undefined
|
||||
return row ? Number(row.value) || 0 : 0
|
||||
}
|
||||
|
||||
/**
|
||||
* Bump the counter. Safe to call inside an open transaction — every caller
|
||||
* that mutates a count does, so the bump commits or rolls back with it.
|
||||
*/
|
||||
export function bumpCassetteStateSeq(): void {
|
||||
if (!db) throw new Error('Database not initialized')
|
||||
db.prepare(
|
||||
'INSERT INTO meta (key, value) VALUES (?, ?) ' +
|
||||
'ON CONFLICT(key) DO UPDATE SET value = CAST(CAST(meta.value AS INTEGER) + 1 AS TEXT)'
|
||||
).run('cassetteStateSeq', '1')
|
||||
}
|
||||
|
||||
/** Record the `created_at` just published, as the next publish's floor. */
|
||||
export function markStatePublished(unixTimestamp: number): void {
|
||||
if (!db) throw new Error('Database not initialized')
|
||||
|
|
@ -617,111 +693,190 @@ export function resetForRepair(): void {
|
|||
})()
|
||||
}
|
||||
|
||||
export type OperatorCassettesPayload = {
|
||||
positions: Record<string, { denomination: number; count: number }>
|
||||
}
|
||||
|
||||
/**
|
||||
* Outcome of applying an operator-authored absolute config. Still the right
|
||||
* shape for fee config, where the operator is the only writer of the value
|
||||
* and a later event simply supersedes an earlier one. Cassette counts left
|
||||
* this model in ADR-004 precisely because they had two writers.
|
||||
*/
|
||||
export type ApplyResult = { applied: true } | { applied: false; reason: string }
|
||||
|
||||
/** One operator-authored operation, as it arrives on the wire. */
|
||||
export type CassetteOp = {
|
||||
id: string
|
||||
at: number
|
||||
type: 'refill' | 'empty' | 'recount' | 'set_denomination'
|
||||
position: number
|
||||
bills?: number
|
||||
count?: number
|
||||
denomination?: number
|
||||
}
|
||||
|
||||
export type ApplyOpsResult = {
|
||||
/** Ids applied by this call. Empty when every op was already on file. */
|
||||
applied: string[]
|
||||
/** Ids rejected, with why. These stay unapplied and unrecorded. */
|
||||
rejected: { id: string; reason: string }[]
|
||||
}
|
||||
|
||||
const CASSETTE_OP_TYPES = new Set(['refill', 'empty', 'recount', 'set_denomination'])
|
||||
|
||||
/**
|
||||
* Atomic apply of an operator-published cassette config (aiolabs/lamassu-next#56).
|
||||
* Validate one operation in isolation. Returns null when it is well-formed.
|
||||
*
|
||||
* Caller has already verified the event signature and decrypted the
|
||||
* content. This function:
|
||||
*
|
||||
* 1. Rechecks replay-protection against `meta.lastKnownConfigCreatedAt`
|
||||
* (defense-in-depth — caller should have done this too).
|
||||
* 2. Validates the payload's `positions` key set is *exactly* the set of
|
||||
* positions currently in the `cassettes` table. The bay count is
|
||||
* hardware-determined and can't be added to or removed from via this
|
||||
* path; only the per-bay denomination and count are operator-mutable.
|
||||
* 3. Validates per-entry `denomination` is a positive int, `count` is a
|
||||
* non-negative int. **Duplicate denominations across positions are
|
||||
* intentionally permitted** — real machines load multiple cassettes
|
||||
* with the same denomination for cash-out throughput.
|
||||
* 4. In a single SQLite transaction: updates `cassettes` rows by position
|
||||
* (denomination + count both mutable per row) AND advances
|
||||
* `meta.lastKnownConfigCreatedAt` to `eventCreatedAt`.
|
||||
*
|
||||
* Mid-write crashes roll back cleanly; on restart the same event is
|
||||
* re-delivered by the relay and the watermark check drops it as already
|
||||
* consumed (or the watermark is pre-event because the tx rolled back,
|
||||
* and the apply runs again from scratch).
|
||||
* Shape errors and unknown positions are treated the same way by the caller:
|
||||
* the op is neither applied nor recorded, so it stays pending on the
|
||||
* operator's dashboard. That is the honest outcome — it did not happen — and
|
||||
* it beats recording it as applied to stop the noise, which would tell the
|
||||
* operator their refill landed when the notes are unaccounted for.
|
||||
*/
|
||||
export function applyOperatorCassettesConfig(
|
||||
payload: OperatorCassettesPayload,
|
||||
eventCreatedAt: number
|
||||
): ApplyResult {
|
||||
function validateCassetteOp(op: CassetteOp, knownPositions: Set<number>): string | null {
|
||||
if (typeof op.id !== 'string' || op.id.length === 0) return 'missing id'
|
||||
if (!CASSETTE_OP_TYPES.has(op.type)) return `unknown type ${String(op.type)}`
|
||||
if (!Number.isInteger(op.position)) return `position must be an integer (got ${op.position})`
|
||||
if (!knownPositions.has(op.position)) return `unknown position ${op.position}`
|
||||
if (!Number.isFinite(op.at)) return 'missing at'
|
||||
|
||||
if (op.type === 'refill') {
|
||||
if (!Number.isInteger(op.bills) || (op.bills as number) <= 0) {
|
||||
return `refill needs a positive integer bills (got ${op.bills})`
|
||||
}
|
||||
}
|
||||
if (op.type === 'recount') {
|
||||
if (!Number.isInteger(op.count) || (op.count as number) < 0) {
|
||||
return `recount needs a non-negative integer count (got ${op.count})`
|
||||
}
|
||||
}
|
||||
if (op.type === 'set_denomination') {
|
||||
if (!Number.isInteger(op.denomination) || (op.denomination as number) <= 0) {
|
||||
return `set_denomination needs a positive integer denomination (got ${op.denomination})`
|
||||
}
|
||||
}
|
||||
return null
|
||||
}
|
||||
|
||||
/**
|
||||
* Apply an operator's cassette operations, skipping any already on file.
|
||||
*
|
||||
* This replaces applying absolute counts. The operator authors what it DID —
|
||||
* a refill in notes added, an empty, a recount, a denomination change — and
|
||||
* this machine, which holds the physical notes, keeps the running total.
|
||||
* Nobody but this process writes a count any more, so there is no second
|
||||
* writer to lose a race to.
|
||||
*
|
||||
* Deltas are not idempotent and addressable events ARE re-delivered on every
|
||||
* relay reconnect, so idempotency is carried explicitly: the operator mints an
|
||||
* id per operation, `cassette_ops` records the ones applied, and a repeat is a
|
||||
* no-op. That is also why there is no `created_at` watermark here any more.
|
||||
* Under absolute counts the watermark was the only replay defence; with
|
||||
* per-op ids it is strictly weaker than the dedup and would do active harm,
|
||||
* because an event that arrives out of order may still carry an operation this
|
||||
* machine has never seen.
|
||||
*
|
||||
* Applied oldest-first by `at`, ties broken by id so two operations stamped in
|
||||
* the same second still order the same way on every machine. Ordering matters
|
||||
* because a recount followed by a refill is not the same as the reverse.
|
||||
*
|
||||
* The whole batch runs in one SQLite transaction with the sequence bump, so a
|
||||
* crash mid-apply rolls back to a coherent count and the next publish re-offers
|
||||
* every op in the window.
|
||||
*/
|
||||
export function applyOperatorCassetteOps(ops: CassetteOp[]): ApplyOpsResult {
|
||||
if (!db) throw new Error('Database not initialized')
|
||||
const database = db
|
||||
const result: ApplyOpsResult = { applied: [], rejected: [] }
|
||||
if (ops.length === 0) return result
|
||||
|
||||
const watermark = getLastKnownConfigCreatedAt()
|
||||
if (eventCreatedAt <= watermark) {
|
||||
return {
|
||||
applied: false,
|
||||
reason: `event.created_at (${eventCreatedAt}) <= lastKnownConfigCreatedAt (${watermark})`,
|
||||
}
|
||||
}
|
||||
|
||||
const currentRows = db.prepare('SELECT position FROM cassettes').all() as { position: number }[]
|
||||
const currentPositions = new Set(currentRows.map((r) => r.position))
|
||||
const payloadPositions = new Set(Object.keys(payload.positions).map((k) => Number(k)))
|
||||
|
||||
if (currentPositions.size !== payloadPositions.size) {
|
||||
return {
|
||||
applied: false,
|
||||
reason: `position count mismatch: state.db has ${currentPositions.size}, payload has ${payloadPositions.size}`,
|
||||
}
|
||||
}
|
||||
for (const p of currentPositions) {
|
||||
if (!payloadPositions.has(p)) {
|
||||
return { applied: false, reason: `payload missing position ${p}` }
|
||||
}
|
||||
}
|
||||
for (const p of payloadPositions) {
|
||||
if (!currentPositions.has(p)) {
|
||||
return { applied: false, reason: `payload includes unknown position ${p}` }
|
||||
}
|
||||
}
|
||||
|
||||
for (const [posKey, entry] of Object.entries(payload.positions)) {
|
||||
if (!Number.isInteger(entry.denomination) || entry.denomination <= 0) {
|
||||
return {
|
||||
applied: false,
|
||||
reason: `denomination must be positive int (position ${posKey}, got ${entry.denomination})`,
|
||||
}
|
||||
}
|
||||
if (!Number.isInteger(entry.count) || entry.count < 0) {
|
||||
return {
|
||||
applied: false,
|
||||
reason: `count must be non-negative int (position ${posKey}, got ${entry.count})`,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
const updateCassette = db.prepare(
|
||||
'UPDATE cassettes SET denomination = ?, count = ? WHERE position = ?'
|
||||
const knownPositions = new Set(
|
||||
(database.prepare('SELECT position FROM cassettes').all() as { position: number }[]).map(
|
||||
(r) => r.position
|
||||
)
|
||||
)
|
||||
const setWatermark = db.prepare('UPDATE meta SET value = ? WHERE key = ?')
|
||||
const seen = database.prepare('SELECT 1 FROM cassette_ops WHERE id = ?')
|
||||
|
||||
const clearUncertain = db.prepare(
|
||||
const pending: CassetteOp[] = []
|
||||
for (const op of ops) {
|
||||
if (op && typeof op.id === 'string' && seen.get(op.id)) continue
|
||||
const reason = validateCassetteOp(op, knownPositions)
|
||||
if (reason) {
|
||||
result.rejected.push({ id: op?.id ?? '<no id>', reason })
|
||||
continue
|
||||
}
|
||||
pending.push(op)
|
||||
}
|
||||
if (pending.length === 0) return result
|
||||
|
||||
pending.sort((a, b) => a.at - b.at || (a.id < b.id ? -1 : a.id > b.id ? 1 : 0))
|
||||
|
||||
const addBills = database.prepare(
|
||||
'UPDATE cassettes SET count = MAX(0, count + ?) WHERE position = ?'
|
||||
)
|
||||
const setCount = database.prepare('UPDATE cassettes SET count = ? WHERE position = ?')
|
||||
const setDenomination = database.prepare(
|
||||
'UPDATE cassettes SET denomination = ? WHERE position = ?'
|
||||
)
|
||||
const recordOp = database.prepare(
|
||||
'INSERT INTO cassette_ops (id, position, op_type, bills, count, denomination, op_at, applied_at) ' +
|
||||
'VALUES (?, ?, ?, ?, ?, ?, ?, ?)'
|
||||
)
|
||||
const upsertMeta = database.prepare(
|
||||
'INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value'
|
||||
)
|
||||
|
||||
const run = db.transaction(() => {
|
||||
for (const [posKey, entry] of Object.entries(payload.positions)) {
|
||||
updateCassette.run(entry.denomination, entry.count, Number(posKey))
|
||||
const appliedAt = Math.floor(Date.now() / 1000)
|
||||
let sawRecount = false
|
||||
|
||||
database.transaction(() => {
|
||||
for (const op of pending) {
|
||||
if (op.type === 'refill') addBills.run(op.bills, op.position)
|
||||
else if (op.type === 'empty') setCount.run(0, op.position)
|
||||
else if (op.type === 'recount') {
|
||||
setCount.run(op.count, op.position)
|
||||
sawRecount = true
|
||||
} else setDenomination.run(op.denomination, op.position)
|
||||
|
||||
recordOp.run(
|
||||
op.id,
|
||||
op.position,
|
||||
op.type,
|
||||
op.bills ?? null,
|
||||
op.count ?? null,
|
||||
op.denomination ?? null,
|
||||
Math.floor(op.at),
|
||||
appliedAt
|
||||
)
|
||||
result.applied.push(op.id)
|
||||
}
|
||||
setWatermark.run(String(eventCreatedAt), 'lastKnownConfigCreatedAt')
|
||||
// The operator just asserted real counts, which is what a recount is.
|
||||
// Whatever made the old numbers untrustworthy no longer applies.
|
||||
clearUncertain.run('countsUncertainSince', '')
|
||||
})
|
||||
bumpCassetteStateSeq()
|
||||
// A recount is an operator opening the bay and counting it, which is
|
||||
// exactly what resolves an unverified count. Nothing else does: a refill
|
||||
// adds to a number still known to be wrong.
|
||||
if (sawRecount) upsertMeta.run('countsUncertainSince', '')
|
||||
})()
|
||||
|
||||
run()
|
||||
console.log(
|
||||
`[StateStore] Applied operator cassettes config @ created_at=${eventCreatedAt} (${Object.keys(payload.positions).length} positions)`
|
||||
`[StateStore] Applied ${result.applied.length} cassette op(s)` +
|
||||
(result.rejected.length ? `, rejected ${result.rejected.length}` : '')
|
||||
)
|
||||
return { applied: true }
|
||||
return result
|
||||
}
|
||||
|
||||
/**
|
||||
* The ids most recently applied, newest first — the acknowledgement leg of
|
||||
* the protocol.
|
||||
*
|
||||
* An addressable event gives its publisher no failure signal at all: the relay
|
||||
* returns OK for an event it then discards, and a losing writer is never told.
|
||||
* Echoing the ids back in this machine's own state document is the only way
|
||||
* the operator can distinguish an operation that landed from one that was
|
||||
* merely sent.
|
||||
*/
|
||||
export function getAppliedOpIds(limit = 50): string[] {
|
||||
if (!db) throw new Error('Database not initialized')
|
||||
const rows = db
|
||||
.prepare('SELECT id FROM cassette_ops ORDER BY applied_at DESC, rowid DESC LIMIT ?')
|
||||
.all(limit) as { id: string }[]
|
||||
return rows.map((r) => r.id)
|
||||
}
|
||||
|
||||
// ---------------------------------------------------------------------------
|
||||
|
|
@ -899,6 +1054,7 @@ export function setCassettes(
|
|||
const row = rows[i]!
|
||||
upsert.run(row.position ?? i + 1, row.denomination, row.count)
|
||||
}
|
||||
bumpCassetteStateSeq()
|
||||
}
|
||||
)
|
||||
|
||||
|
|
@ -913,11 +1069,14 @@ export function setCassettes(
|
|||
*/
|
||||
export function updateCassetteCountByPosition(position: number, delta: number): void {
|
||||
if (!db) throw new Error('Database not initialized')
|
||||
const database = db
|
||||
|
||||
db.prepare('UPDATE cassettes SET count = MAX(0, count + ?) WHERE position = ?').run(
|
||||
delta,
|
||||
position
|
||||
)
|
||||
database.transaction(() => {
|
||||
database
|
||||
.prepare('UPDATE cassettes SET count = MAX(0, count + ?) WHERE position = ?')
|
||||
.run(delta, position)
|
||||
bumpCassetteStateSeq()
|
||||
})()
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -1120,6 +1279,10 @@ export function recordTransaction(tx: TransactionInput): void {
|
|||
}
|
||||
}
|
||||
}
|
||||
// The counts moved, so the sequence must move with them, inside this
|
||||
// same transaction. It rides in the state document as the operator's
|
||||
// way to reject a regression without trusting either clock.
|
||||
bumpCassetteStateSeq()
|
||||
}
|
||||
|
||||
if (t.type === 'cash_in') {
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue