Merge pull request 'Consume operator cassette operations instead of counts' (#106) from feat/cassette-ops-consumer into dev
Reviewed-on: #106
This commit is contained in:
commit
8d2509b092
7 changed files with 573 additions and 183 deletions
|
|
@ -12,8 +12,10 @@
|
||||||
|
|
||||||
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
|
import { afterEach, beforeEach, describe, expect, it } from 'vitest'
|
||||||
import {
|
import {
|
||||||
applyOperatorCassettesConfig,
|
applyOperatorCassetteOps,
|
||||||
closeDatabase,
|
closeDatabase,
|
||||||
|
getAppliedOpIds,
|
||||||
|
getCassetteStateSeq,
|
||||||
getCashbox,
|
getCashbox,
|
||||||
getCountsUncertainSince,
|
getCountsUncertainSince,
|
||||||
getInventory,
|
getInventory,
|
||||||
|
|
@ -337,31 +339,157 @@ describe('state-store: unverified counts after a silent dispense', () => {
|
||||||
expect(getCountsUncertainSince()).toBe(1000)
|
expect(getCountsUncertainSince()).toBe(1000)
|
||||||
})
|
})
|
||||||
|
|
||||||
it('clears when an operator asserts real counts', () => {
|
it('clears on a recount, because that is what a recount is', () => {
|
||||||
markCountsUncertain(1000)
|
markCountsUncertain(1000)
|
||||||
const applied = applyOperatorCassettesConfig(
|
const result = applyOperatorCassetteOps([
|
||||||
{
|
{ id: 'op-1', at: 1_700_000_000, type: 'recount', position: 1, count: 40 },
|
||||||
positions: {
|
])
|
||||||
'1': { denomination: 20, count: 40 },
|
expect(result.applied).toEqual(['op-1'])
|
||||||
'2': { denomination: 20, count: 40 },
|
|
||||||
'3': { denomination: 50, count: 25 },
|
|
||||||
},
|
|
||||||
},
|
|
||||||
1_700_000_000
|
|
||||||
)
|
|
||||||
expect(applied.applied).toBe(true)
|
|
||||||
// A recount is exactly the operator asserting authoritative counts.
|
|
||||||
expect(getCountsUncertainSince()).toBeNull()
|
expect(getCountsUncertainSince()).toBeNull()
|
||||||
})
|
})
|
||||||
|
|
||||||
it('leaves the flag alone when the operator config is rejected', () => {
|
it('does not clear on a refill', () => {
|
||||||
|
// A refill adds to a number still known to be wrong. Only someone
|
||||||
|
// opening the bay and counting it resolves that.
|
||||||
markCountsUncertain(1000)
|
markCountsUncertain(1000)
|
||||||
// Position key-set mismatch — the bay layout is hardware-determined.
|
const result = applyOperatorCassetteOps([
|
||||||
const applied = applyOperatorCassettesConfig(
|
{ id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 10 },
|
||||||
{ positions: { '1': { denomination: 20, count: 40 } } },
|
])
|
||||||
1_700_000_001
|
expect(result.applied).toEqual(['op-1'])
|
||||||
)
|
expect(getCountsUncertainSince()).toBe(1000)
|
||||||
expect(applied.applied).toBe(false)
|
})
|
||||||
|
|
||||||
|
it('leaves the flag alone when the op is rejected', () => {
|
||||||
|
markCountsUncertain(1000)
|
||||||
|
// Bay 9 does not exist — the layout is hardware-determined.
|
||||||
|
const result = applyOperatorCassetteOps([
|
||||||
|
{ id: 'op-1', at: 1_700_000_000, type: 'recount', position: 9, count: 40 },
|
||||||
|
])
|
||||||
|
expect(result.applied).toEqual([])
|
||||||
|
expect(result.rejected).toHaveLength(1)
|
||||||
expect(getCountsUncertainSince()).toBe(1000)
|
expect(getCountsUncertainSince()).toBe(1000)
|
||||||
})
|
})
|
||||||
})
|
})
|
||||||
|
|
||||||
|
describe('state-store: operator cassette operations (ADR-004)', () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
seedDuplicateDenomBays()
|
||||||
|
})
|
||||||
|
|
||||||
|
it('applies a refill as a delta, not a total', () => {
|
||||||
|
applyOperatorCassetteOps([
|
||||||
|
{ id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 30 },
|
||||||
|
])
|
||||||
|
expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(80)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('is a no-op on a re-delivered operation', () => {
|
||||||
|
// Addressable events are re-delivered on every relay reconnect and the
|
||||||
|
// operator republishes a WINDOW, so the same op arrives many times. A
|
||||||
|
// delta applied twice is simply wrong, which is why every op carries an
|
||||||
|
// id and this table records the ones already applied.
|
||||||
|
const op = {
|
||||||
|
id: 'op-1',
|
||||||
|
at: 1_700_000_000,
|
||||||
|
type: 'refill' as const,
|
||||||
|
position: 1,
|
||||||
|
bills: 30,
|
||||||
|
}
|
||||||
|
expect(applyOperatorCassetteOps([op]).applied).toEqual(['op-1'])
|
||||||
|
expect(applyOperatorCassetteOps([op]).applied).toEqual([])
|
||||||
|
expect(applyOperatorCassetteOps([op, op]).applied).toEqual([])
|
||||||
|
expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(80)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('applies only the unseen ops from a window that mixes both', () => {
|
||||||
|
applyOperatorCassetteOps([
|
||||||
|
{ id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 10 },
|
||||||
|
])
|
||||||
|
const result = applyOperatorCassetteOps([
|
||||||
|
{ id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 10 },
|
||||||
|
{ id: 'op-2', at: 1_700_000_001, type: 'refill', position: 1, bills: 5 },
|
||||||
|
])
|
||||||
|
expect(result.applied).toEqual(['op-2'])
|
||||||
|
expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(65)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('applies a window oldest-first regardless of arrival order', () => {
|
||||||
|
// A recount then a refill is not the same as the reverse, so ordering is
|
||||||
|
// load-bearing and cannot be left to however the array arrived.
|
||||||
|
applyOperatorCassetteOps([
|
||||||
|
{ id: 'op-b', at: 1_700_000_002, type: 'refill', position: 1, bills: 7 },
|
||||||
|
{ id: 'op-a', at: 1_700_000_001, type: 'recount', position: 1, count: 3 },
|
||||||
|
])
|
||||||
|
expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(10)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('empties a bay and sets a denomination', () => {
|
||||||
|
applyOperatorCassetteOps([
|
||||||
|
{ id: 'op-1', at: 1_700_000_000, type: 'empty', position: 2 },
|
||||||
|
{ id: 'op-2', at: 1_700_000_001, type: 'set_denomination', position: 2, denomination: 10 },
|
||||||
|
])
|
||||||
|
const bay = loadCassettes().find((c) => c.position === 2)!
|
||||||
|
expect(bay.count).toBe(0)
|
||||||
|
expect(bay.denomination).toBe(10)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('rejects a malformed op without applying or recording it', () => {
|
||||||
|
// Unrecorded on purpose: it stays pending on the operator's dashboard,
|
||||||
|
// which is the honest outcome. Recording it as applied would silence the
|
||||||
|
// noise by telling the operator their refill landed.
|
||||||
|
const result = applyOperatorCassetteOps([
|
||||||
|
{ id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: -5 },
|
||||||
|
])
|
||||||
|
expect(result.applied).toEqual([])
|
||||||
|
expect(result.rejected[0]!.id).toBe('op-1')
|
||||||
|
expect(loadCassettes().find((c) => c.position === 1)!.count).toBe(50)
|
||||||
|
expect(getAppliedOpIds()).not.toContain('op-1')
|
||||||
|
})
|
||||||
|
|
||||||
|
it('echoes applied ids back, newest first', () => {
|
||||||
|
applyOperatorCassetteOps([
|
||||||
|
{ id: 'op-1', at: 1_700_000_000, type: 'refill', position: 1, bills: 1 },
|
||||||
|
{ id: 'op-2', at: 1_700_000_001, type: 'refill', position: 1, bills: 1 },
|
||||||
|
])
|
||||||
|
expect(getAppliedOpIds()).toContain('op-1')
|
||||||
|
expect(getAppliedOpIds()).toContain('op-2')
|
||||||
|
})
|
||||||
|
|
||||||
|
it('advances the sequence on an applied op but not on a duplicate', () => {
|
||||||
|
const op = {
|
||||||
|
id: 'op-1',
|
||||||
|
at: 1_700_000_000,
|
||||||
|
type: 'refill' as const,
|
||||||
|
position: 1,
|
||||||
|
bills: 1,
|
||||||
|
}
|
||||||
|
const before = getCassetteStateSeq()
|
||||||
|
applyOperatorCassetteOps([op])
|
||||||
|
const after = getCassetteStateSeq()
|
||||||
|
expect(after).toBeGreaterThan(before)
|
||||||
|
applyOperatorCassetteOps([op])
|
||||||
|
expect(getCassetteStateSeq()).toBe(after)
|
||||||
|
})
|
||||||
|
|
||||||
|
it('advances the sequence on a dispense', () => {
|
||||||
|
const before = getCassetteStateSeq()
|
||||||
|
recordTransaction({
|
||||||
|
...TX_BASE,
|
||||||
|
txid: 'tx-seq',
|
||||||
|
type: 'cash_out',
|
||||||
|
status: 'complete',
|
||||||
|
bills: [{ denomination: 20, count: 1 }],
|
||||||
|
cassettes: [
|
||||||
|
{
|
||||||
|
name: 'cassette1',
|
||||||
|
position: 1,
|
||||||
|
denomination: 20,
|
||||||
|
provisioned: 1,
|
||||||
|
dispensed: 1,
|
||||||
|
rejected: 0,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
})
|
||||||
|
expect(getCassetteStateSeq()).toBeGreaterThan(before)
|
||||||
|
})
|
||||||
|
})
|
||||||
|
|
|
||||||
|
|
@ -30,14 +30,17 @@ import {
|
||||||
markStatePublished,
|
markStatePublished,
|
||||||
resetStatePublishWatermark,
|
resetStatePublishWatermark,
|
||||||
resetForRepair,
|
resetForRepair,
|
||||||
applyOperatorCassettesConfig,
|
applyOperatorCassetteOps,
|
||||||
|
getAppliedOpIds,
|
||||||
|
getCassetteStateSeq,
|
||||||
getFeeConfig,
|
getFeeConfig,
|
||||||
getLastKnownFeeConfigCreatedAt,
|
getLastKnownFeeConfigCreatedAt,
|
||||||
applyFeeConfig,
|
applyFeeConfig,
|
||||||
getBunkerBinding,
|
getBunkerBinding,
|
||||||
saveBunkerBinding,
|
saveBunkerBinding,
|
||||||
clearBunkerBinding,
|
clearBunkerBinding,
|
||||||
type OperatorCassettesPayload,
|
type CassetteOp,
|
||||||
|
type ApplyOpsResult,
|
||||||
type FeeConfigPayload,
|
type FeeConfigPayload,
|
||||||
type FeeConfigRow,
|
type FeeConfigRow,
|
||||||
type ApplyResult,
|
type ApplyResult,
|
||||||
|
|
@ -566,10 +569,13 @@ ipcMain.handle('state:mark-state-published', (_event, unixTimestamp: number): vo
|
||||||
markStatePublished(unixTimestamp)
|
markStatePublished(unixTimestamp)
|
||||||
})
|
})
|
||||||
ipcMain.handle(
|
ipcMain.handle(
|
||||||
'state:apply-operator-cassettes-config',
|
'state:apply-operator-cassette-ops',
|
||||||
(_event, payload: OperatorCassettesPayload, eventCreatedAt: number): ApplyResult =>
|
(_event, ops: CassetteOp[]): ApplyOpsResult => applyOperatorCassetteOps(ops)
|
||||||
applyOperatorCassettesConfig(payload, eventCreatedAt)
|
|
||||||
)
|
)
|
||||||
|
ipcMain.handle('state:get-applied-op-ids', (_event, limit?: number): string[] =>
|
||||||
|
getAppliedOpIds(limit)
|
||||||
|
)
|
||||||
|
ipcMain.handle('state:get-cassette-state-seq', (): number => getCassetteStateSeq())
|
||||||
|
|
||||||
// Operator-fees consumer (aiolabs/lamassu-next#57) — persisted singleton
|
// Operator-fees consumer (aiolabs/lamassu-next#57) — persisted singleton
|
||||||
// fee config + per-d-tag replay watermark + atomic apply for kind-30078
|
// fee config + per-d-tag replay watermark + atomic apply for kind-30078
|
||||||
|
|
|
||||||
|
|
@ -163,13 +163,23 @@ contextBridge.exposeInMainWorld('electronAPI', {
|
||||||
}): Promise<{ ok: boolean; bolt11?: string; reason?: string }> =>
|
}): Promise<{ ok: boolean; bolt11?: string; reason?: string }> =>
|
||||||
ipcRenderer.invoke('lnurl:pay-session', args),
|
ipcRenderer.invoke('lnurl:pay-session', args),
|
||||||
|
|
||||||
applyOperatorCassettesConfig: (
|
applyOperatorCassetteOps: (
|
||||||
payload: {
|
ops: {
|
||||||
positions: Record<string, { denomination: number; count: number }>
|
id: string
|
||||||
},
|
at: number
|
||||||
eventCreatedAt: number
|
type: 'refill' | 'empty' | 'recount' | 'set_denomination'
|
||||||
): Promise<{ applied: true } | { applied: false; reason: string }> =>
|
position: number
|
||||||
ipcRenderer.invoke('state:apply-operator-cassettes-config', payload, eventCreatedAt),
|
bills?: number
|
||||||
|
count?: number
|
||||||
|
denomination?: number
|
||||||
|
}[]
|
||||||
|
): Promise<{
|
||||||
|
applied: string[]
|
||||||
|
rejected: { id: string; reason: string }[]
|
||||||
|
}> => ipcRenderer.invoke('state:apply-operator-cassette-ops', ops),
|
||||||
|
getAppliedOpIds: (limit?: number): Promise<string[]> =>
|
||||||
|
ipcRenderer.invoke('state:get-applied-op-ids', limit),
|
||||||
|
getCassetteStateSeq: (): Promise<number> => ipcRenderer.invoke('state:get-cassette-state-seq'),
|
||||||
|
|
||||||
// Operator-fees consumer (aiolabs/lamassu-next#57)
|
// Operator-fees consumer (aiolabs/lamassu-next#57)
|
||||||
getFeeConfig: (): Promise<{
|
getFeeConfig: (): Promise<{
|
||||||
|
|
@ -314,10 +324,22 @@ declare global {
|
||||||
lnurlw: string
|
lnurlw: string
|
||||||
amountMsat: number
|
amountMsat: number
|
||||||
}) => Promise<{ ok: boolean; bolt11?: string; reason?: string }>
|
}) => Promise<{ ok: boolean; bolt11?: string; reason?: string }>
|
||||||
applyOperatorCassettesConfig: (
|
applyOperatorCassetteOps: (
|
||||||
payload: { positions: Record<string, { denomination: number; count: number }> },
|
ops: {
|
||||||
eventCreatedAt: number
|
id: string
|
||||||
) => Promise<{ applied: true } | { applied: false; reason: string }>
|
at: number
|
||||||
|
type: 'refill' | 'empty' | 'recount' | 'set_denomination'
|
||||||
|
position: number
|
||||||
|
bills?: number
|
||||||
|
count?: number
|
||||||
|
denomination?: number
|
||||||
|
}[]
|
||||||
|
) => Promise<{
|
||||||
|
applied: string[]
|
||||||
|
rejected: { id: string; reason: string }[]
|
||||||
|
}>
|
||||||
|
getAppliedOpIds: (limit?: number) => Promise<string[]>
|
||||||
|
getCassetteStateSeq: () => Promise<number>
|
||||||
getFeeConfig: () => Promise<{
|
getFeeConfig: () => Promise<{
|
||||||
cashInFeeFraction: number
|
cashInFeeFraction: number
|
||||||
cashOutFeeFraction: number
|
cashOutFeeFraction: number
|
||||||
|
|
|
||||||
|
|
@ -15,7 +15,7 @@ import fs from 'node:fs'
|
||||||
|
|
||||||
let db: Database.Database | null = null
|
let db: Database.Database | null = null
|
||||||
|
|
||||||
const SCHEMA_VERSION = '12'
|
const SCHEMA_VERSION = '13'
|
||||||
|
|
||||||
function getDbPath(): string {
|
function getDbPath(): string {
|
||||||
const prodDir = '/var/lib/bitspire'
|
const prodDir = '/var/lib/bitspire'
|
||||||
|
|
@ -57,6 +57,17 @@ export function initDatabase(dbPath?: string): void {
|
||||||
count INTEGER NOT NULL DEFAULT 0
|
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 (
|
CREATE TABLE IF NOT EXISTS cashbox (
|
||||||
id INTEGER PRIMARY KEY CHECK (id = 1),
|
id INTEGER PRIMARY KEY CHECK (id = 1),
|
||||||
total_bills INTEGER NOT NULL DEFAULT 0,
|
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)')
|
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.
|
// Defensive: a fresh install at SCHEMA_VERSION skips all migrations.
|
||||||
// Seed the operator-config meta rows if they're missing (idempotent).
|
// Seed the operator-config meta rows if they're missing (idempotent).
|
||||||
const seedMeta = db.prepare('INSERT OR IGNORE INTO meta (key, value) VALUES (?, ?)')
|
const seedMeta = db.prepare('INSERT OR IGNORE INTO meta (key, value) VALUES (?, ?)')
|
||||||
seedMeta.run('lastKnownConfigCreatedAt', '0')
|
seedMeta.run('lastKnownConfigCreatedAt', '0')
|
||||||
seedMeta.run('bootstrapPublishedAt', '')
|
seedMeta.run('bootstrapPublishedAt', '')
|
||||||
seedMeta.run('lastKnownFeeConfigCreatedAt', '0')
|
seedMeta.run('lastKnownFeeConfigCreatedAt', '0')
|
||||||
|
seedMeta.run('cassetteStateSeq', '0')
|
||||||
|
|
||||||
const cashboxRow = db.prepare('SELECT id FROM cashbox WHERE id = 1').get()
|
const cashboxRow = db.prepare('SELECT id FROM cashbox WHERE id = 1').get()
|
||||||
if (!cashboxRow) {
|
if (!cashboxRow) {
|
||||||
|
|
@ -471,6 +517,36 @@ export function clearCountsUncertain(): void {
|
||||||
).run('countsUncertainSince', '')
|
).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. */
|
/** Record the `created_at` just published, as the next publish's floor. */
|
||||||
export function markStatePublished(unixTimestamp: number): void {
|
export function markStatePublished(unixTimestamp: number): void {
|
||||||
if (!db) throw new Error('Database not initialized')
|
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 }
|
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
|
* Shape errors and unknown positions are treated the same way by the caller:
|
||||||
* content. This function:
|
* 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
|
||||||
* 1. Rechecks replay-protection against `meta.lastKnownConfigCreatedAt`
|
* it beats recording it as applied to stop the noise, which would tell the
|
||||||
* (defense-in-depth — caller should have done this too).
|
* operator their refill landed when the notes are unaccounted for.
|
||||||
* 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).
|
|
||||||
*/
|
*/
|
||||||
export function applyOperatorCassettesConfig(
|
function validateCassetteOp(op: CassetteOp, knownPositions: Set<number>): string | null {
|
||||||
payload: OperatorCassettesPayload,
|
if (typeof op.id !== 'string' || op.id.length === 0) return 'missing id'
|
||||||
eventCreatedAt: number
|
if (!CASSETTE_OP_TYPES.has(op.type)) return `unknown type ${String(op.type)}`
|
||||||
): ApplyResult {
|
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')
|
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()
|
const knownPositions = new Set(
|
||||||
if (eventCreatedAt <= watermark) {
|
(database.prepare('SELECT position FROM cassettes').all() as { position: number }[]).map(
|
||||||
return {
|
(r) => r.position
|
||||||
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 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'
|
'INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value'
|
||||||
)
|
)
|
||||||
|
|
||||||
const run = db.transaction(() => {
|
const appliedAt = Math.floor(Date.now() / 1000)
|
||||||
for (const [posKey, entry] of Object.entries(payload.positions)) {
|
let sawRecount = false
|
||||||
updateCassette.run(entry.denomination, entry.count, Number(posKey))
|
|
||||||
|
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')
|
bumpCassetteStateSeq()
|
||||||
// The operator just asserted real counts, which is what a recount is.
|
// A recount is an operator opening the bay and counting it, which is
|
||||||
// Whatever made the old numbers untrustworthy no longer applies.
|
// exactly what resolves an unverified count. Nothing else does: a refill
|
||||||
clearUncertain.run('countsUncertainSince', '')
|
// adds to a number still known to be wrong.
|
||||||
})
|
if (sawRecount) upsertMeta.run('countsUncertainSince', '')
|
||||||
|
})()
|
||||||
|
|
||||||
run()
|
|
||||||
console.log(
|
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]!
|
const row = rows[i]!
|
||||||
upsert.run(row.position ?? i + 1, row.denomination, row.count)
|
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 {
|
export function updateCassetteCountByPosition(position: number, delta: number): void {
|
||||||
if (!db) throw new Error('Database not initialized')
|
if (!db) throw new Error('Database not initialized')
|
||||||
|
const database = db
|
||||||
|
|
||||||
db.prepare('UPDATE cassettes SET count = MAX(0, count + ?) WHERE position = ?').run(
|
database.transaction(() => {
|
||||||
delta,
|
database
|
||||||
position
|
.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') {
|
if (t.type === 'cash_in') {
|
||||||
|
|
|
||||||
|
|
@ -1,11 +1,17 @@
|
||||||
/**
|
/**
|
||||||
* Operator-config consumer (aiolabs/lamassu-next#56).
|
* Operator-config consumer (aiolabs/lamassu-next#56, v2 per bitspire ADR-004).
|
||||||
*
|
*
|
||||||
* Subscribes to operator-published kind-30078 events carrying cassette
|
* Subscribes to operator-published kind-30078 events carrying cassette
|
||||||
* config updates, validates + applies them to state.db, and hot-reloads
|
* OPERATIONS — a refill, an empty, a recount, a denomination change —
|
||||||
* the HAL dispenser. Also publishes a one-shot ATM-state hello-event on
|
* applies the ones it has not already seen, and hot-reloads the HAL
|
||||||
* first boot so the operator dashboard (satmachineadmin) can auto-populate
|
* dispenser. It also publishes this machine's cassette state, which is
|
||||||
* `cassette_configs` rows for this machine.
|
* what populates the operator dashboard's bay rows.
|
||||||
|
*
|
||||||
|
* 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: a dashboard form loaded before a dispense silently
|
||||||
|
* discarded that dispense, and neither side could detect it. This machine now
|
||||||
|
* owns the count — it holds the notes — and the operator says what it did.
|
||||||
*
|
*
|
||||||
* Architecture (see ~/dev/coordination/log.md entries on 2026-05-30):
|
* Architecture (see ~/dev/coordination/log.md entries on 2026-05-30):
|
||||||
*
|
*
|
||||||
|
|
@ -35,6 +41,20 @@ import type {} from '@/types/electron'
|
||||||
|
|
||||||
const KIND_NIP78 = 30078
|
const KIND_NIP78 = 30078
|
||||||
|
|
||||||
|
/** The wire schema this machine speaks. Operations, not counts (ADR-004). */
|
||||||
|
const CASSETTE_SCHEMA_VERSION = 2
|
||||||
|
|
||||||
|
/** One operator-authored operation, as it arrives on the wire. */
|
||||||
|
type CassetteOp = {
|
||||||
|
id: string
|
||||||
|
at: number
|
||||||
|
type: 'refill' | 'empty' | 'recount' | 'set_denomination'
|
||||||
|
position: number
|
||||||
|
bills?: number
|
||||||
|
count?: number
|
||||||
|
denomination?: number
|
||||||
|
}
|
||||||
|
|
||||||
/** Accept operator events stamped up to this many seconds in the future. */
|
/** Accept operator events stamped up to this many seconds in the future. */
|
||||||
const MAX_FUTURE_SKEW_S = 60
|
const MAX_FUTURE_SKEW_S = 60
|
||||||
|
|
||||||
|
|
@ -165,17 +185,13 @@ async function handleOperatorConfigEvent(
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
// 2. Replay protection — drop stale events. NIP-78 replaceable events
|
// 2. There is deliberately no `created_at` watermark here any more.
|
||||||
// DO get re-delivered on reconnect/restart; without this check, the
|
// Under absolute counts it was the only replay defence, and it cost us:
|
||||||
// ATM would re-apply the same payload on every boot and clobber any
|
// an event re-delivered out of order was dropped whole, operations
|
||||||
// cash-out decrements that landed between operator publishes.
|
// included. Idempotency now rides on the operations themselves — the
|
||||||
const watermark = await api.getLastKnownConfigCreatedAt()
|
// operator mints an id per op and this machine records the ones it
|
||||||
if (event.created_at <= watermark) {
|
// applied — which is strictly stronger, because it survives an event
|
||||||
console.log(
|
// that mixes operations we have seen with ones we have not.
|
||||||
`[OperatorConfig] Stale event dropped (created_at=${event.created_at} <= watermark=${watermark})`
|
|
||||||
)
|
|
||||||
return
|
|
||||||
}
|
|
||||||
|
|
||||||
// 3. Clock-skew defense — reject events stamped too far in the future.
|
// 3. Clock-skew defense — reject events stamped too far in the future.
|
||||||
// Limits damage from a leaked operator nsec future-stamping a fake
|
// Limits damage from a leaked operator nsec future-stamping a fake
|
||||||
|
|
@ -189,7 +205,7 @@ async function handleOperatorConfigEvent(
|
||||||
}
|
}
|
||||||
|
|
||||||
// 4. Decrypt content (NIP-44 v2).
|
// 4. Decrypt content (NIP-44 v2).
|
||||||
let parsed: { positions: Record<string, { denomination: number; count: number }> }
|
let parsed: { schema_version?: number; ops?: unknown }
|
||||||
try {
|
try {
|
||||||
const plaintext = await cfg.signer.nip44Decrypt(event.pubkey, event.content)
|
const plaintext = await cfg.signer.nip44Decrypt(event.pubkey, event.content)
|
||||||
parsed = JSON.parse(plaintext) as typeof parsed
|
parsed = JSON.parse(plaintext) as typeof parsed
|
||||||
|
|
@ -197,22 +213,36 @@ async function handleOperatorConfigEvent(
|
||||||
console.error('[OperatorConfig] Decrypt/parse failed:', err)
|
console.error('[OperatorConfig] Decrypt/parse failed:', err)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if (!parsed || typeof parsed !== 'object' || !parsed.positions) {
|
if (!parsed || typeof parsed !== 'object' || !Array.isArray(parsed.ops)) {
|
||||||
console.error('[OperatorConfig] Payload missing `positions` field')
|
// A v1 operator publishing absolute counts lands here and is ignored.
|
||||||
|
// That direction fails safe: the machine keeps its own counts, which it
|
||||||
|
// is now the only writer of, and simply will not dispense notes it
|
||||||
|
// believes it lacks. The opposite — applying a count from a form loaded
|
||||||
|
// before a dispense — is what ADR-004 exists to stop.
|
||||||
|
console.error('[OperatorConfig] Payload missing `ops` array — dropped')
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
const ops = parsed.ops as CassetteOp[]
|
||||||
|
|
||||||
// 5. Atomic apply (cassettes + meta watermark) via IPC. The state-store
|
// 5. Apply the ones we have not seen, in one transaction with the sequence
|
||||||
// function re-validates watermark + position key-set equality +
|
// bump. No `created_at` watermark: each op carries an operator-minted id
|
||||||
// per-entry types inside the SQLite transaction. Duplicate
|
// and the machine records what it applied, so a re-delivered event is a
|
||||||
// denominations across positions are allowed — real machines load
|
// no-op on its own merits. The watermark would be strictly weaker and
|
||||||
// N cassettes of the same denomination for cash-out throughput.
|
// actively harmful — an event arriving out of order can still carry an
|
||||||
const result = await api.applyOperatorCassettesConfig(
|
// operation this machine has never seen.
|
||||||
{ positions: parsed.positions },
|
const result = await api.applyOperatorCassetteOps(ops)
|
||||||
event.created_at
|
for (const bad of result.rejected) {
|
||||||
)
|
console.warn(`[OperatorConfig] Op ${bad.id} rejected: ${bad.reason}`)
|
||||||
if (!result.applied) {
|
}
|
||||||
console.warn('[OperatorConfig] Apply rejected:', result.reason)
|
if (result.applied.length === 0) {
|
||||||
|
console.log(`[OperatorConfig] No new ops in event ${event.id.slice(0, 12)}…`)
|
||||||
|
// Still republish: the operator learns from our applied_ops echo that
|
||||||
|
// earlier operations landed, and an event carrying nothing new can be the
|
||||||
|
// first one we successfully answer after a relay outage.
|
||||||
|
const machineIdNoop = cfg.machineId ?? cfg.signer.pubkey
|
||||||
|
await publishCassettesState(cfg, api, machineIdNoop).catch((err) =>
|
||||||
|
console.warn('[OperatorConfig] post-apply cassettes-state republish failed:', err)
|
||||||
|
)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -220,8 +250,8 @@ async function handleOperatorConfigEvent(
|
||||||
// picks up the new per-position mapping. state.db is already updated;
|
// picks up the new per-position mapping. state.db is already updated;
|
||||||
// HAL re-init failure means the renderer's persistedInventory may be
|
// HAL re-init failure means the renderer's persistedInventory may be
|
||||||
// ahead of the HAL until next service restart — log loudly but don't
|
// ahead of the HAL until next service restart — log loudly but don't
|
||||||
// unwind the state.db apply (the operator wants their config landed;
|
// unwind the state.db apply (the operation happened physically; HAL
|
||||||
// HAL can catch up).
|
// can catch up).
|
||||||
const cassettesAfter = await api.loadCassettes()
|
const cassettesAfter = await api.loadCassettes()
|
||||||
const halResult = await api.halReloadCassettes(
|
const halResult = await api.halReloadCassettes(
|
||||||
cassettesAfter.map((c) => ({
|
cassettesAfter.map((c) => ({
|
||||||
|
|
@ -233,9 +263,7 @@ async function handleOperatorConfigEvent(
|
||||||
if (!halResult.ok) {
|
if (!halResult.ok) {
|
||||||
console.error('[OperatorConfig] HAL reload failed:', halResult.error)
|
console.error('[OperatorConfig] HAL reload failed:', halResult.error)
|
||||||
}
|
}
|
||||||
console.log(
|
console.log(`[OperatorConfig] Applied ops: ${result.applied.join(', ')}`)
|
||||||
`[OperatorConfig] Applied — created_at=${event.created_at}, positions=${Object.keys(parsed.positions).join(',')}`
|
|
||||||
)
|
|
||||||
|
|
||||||
// Republish our resulting cassette state so the operator's view reflects the
|
// Republish our resulting cassette state so the operator's view reflects the
|
||||||
// applied config (the "on cassette reload" case). Different d-tag from the
|
// applied config (the "on cassette reload" case). Different d-tag from the
|
||||||
|
|
@ -274,9 +302,22 @@ async function publishCassettesState(
|
||||||
// ignores this, so it needs no coordinated release. When set, the counts
|
// ignores this, so it needs no coordinated release. When set, the counts
|
||||||
// above are the machine's best guess, not a measurement.
|
// above are the machine's best guess, not a measurement.
|
||||||
const countsUncertainSince = await api.getCountsUncertainSince()
|
const countsUncertainSince = await api.getCountsUncertainSince()
|
||||||
const payload = countsUncertainSince
|
|
||||||
? { positions, counts_uncertain_since: countsUncertainSince }
|
// `applied_ops` is the acknowledgement leg. An addressable event gives its
|
||||||
: { positions }
|
// publisher no failure signal — the relay returns OK for an event it then
|
||||||
|
// discards — so echoing the ids back is the only way the operator can tell
|
||||||
|
// an operation that landed from one that was merely sent. `seq` lets a
|
||||||
|
// reader reject a regression without trusting a clock: `created_at` has
|
||||||
|
// second granularity and ties break on event id, so it cannot order two
|
||||||
|
// reports from the same second.
|
||||||
|
const [appliedOps, seq] = await Promise.all([api.getAppliedOpIds(), api.getCassetteStateSeq()])
|
||||||
|
const payload: Record<string, unknown> = {
|
||||||
|
schema_version: CASSETTE_SCHEMA_VERSION,
|
||||||
|
positions,
|
||||||
|
seq,
|
||||||
|
applied_ops: appliedOps,
|
||||||
|
}
|
||||||
|
if (countsUncertainSince) payload.counts_uncertain_since = countsUncertainSince
|
||||||
const ciphertext = await cfg.signer.nip44Encrypt(operatorPubkey, JSON.stringify(payload))
|
const ciphertext = await cfg.signer.nip44Encrypt(operatorPubkey, JSON.stringify(payload))
|
||||||
|
|
||||||
// Force the stamp strictly above our last one. Addressable events are ordered
|
// Force the stamp strictly above our last one. Addressable events are ordered
|
||||||
|
|
|
||||||
20
apps/machine/src/types/electron.d.ts
vendored
20
apps/machine/src/types/electron.d.ts
vendored
|
|
@ -183,10 +183,22 @@ declare global {
|
||||||
pay: CardSession['pay']
|
pay: CardSession['pay']
|
||||||
amountMsat: number
|
amountMsat: number
|
||||||
}) => Promise<{ ok: boolean; bolt11?: string; reason?: string }>
|
}) => Promise<{ ok: boolean; bolt11?: string; reason?: string }>
|
||||||
applyOperatorCassettesConfig: (
|
applyOperatorCassetteOps: (
|
||||||
payload: { positions: Record<string, { denomination: number; count: number }> },
|
ops: {
|
||||||
eventCreatedAt: number
|
id: string
|
||||||
) => Promise<{ applied: true } | { applied: false; reason: string }>
|
at: number
|
||||||
|
type: 'refill' | 'empty' | 'recount' | 'set_denomination'
|
||||||
|
position: number
|
||||||
|
bills?: number
|
||||||
|
count?: number
|
||||||
|
denomination?: number
|
||||||
|
}[]
|
||||||
|
) => Promise<{
|
||||||
|
applied: string[]
|
||||||
|
rejected: { id: string; reason: string }[]
|
||||||
|
}>
|
||||||
|
getAppliedOpIds: (limit?: number) => Promise<string[]>
|
||||||
|
getCassetteStateSeq: () => Promise<number>
|
||||||
getFeeConfig: () => Promise<{
|
getFeeConfig: () => Promise<{
|
||||||
cashInFeeFraction: number
|
cashInFeeFraction: number
|
||||||
cashOutFeeFraction: number
|
cashOutFeeFraction: number
|
||||||
|
|
|
||||||
|
|
@ -65,21 +65,23 @@ The state document carries `applied_ops`, so the dashboard can render each publi
|
||||||
as applied or pending. This supplies the feedback leg a replaceable event cannot, without
|
as applied or pending. This supplies the feedback leg a replaceable event cannot, without
|
||||||
needing the transport to report failures.
|
needing the transport to report failures.
|
||||||
|
|
||||||
### 4a. Until decisions 1 to 4 ship, the overwrite is warned about, not prevented.
|
### 4a. Superseded. Before decisions 1 to 4 shipped, the overwrite was warned about.
|
||||||
|
|
||||||
The dashboard's publish dialog already states the failure plainly — that the publish will
|
The dashboard's publish dialog stated the failure plainly — that the publish would overwrite
|
||||||
overwrite the ATM's tracked counts, that decrements since the last baseline will be lost, and
|
the ATM's tracked counts, that decrements since the last baseline would be lost, and that it
|
||||||
that it should follow a physical refill rather than a mid-day tweak. It also says v2
|
should follow a physical refill rather than a mid-day tweak.
|
||||||
reconciliation will replace it.
|
|
||||||
|
|
||||||
Recording this because it changes how the gap should be read. It is a known, deliberately
|
Kept here rather than deleted, because it is the calibration for how much a warning is worth.
|
||||||
accepted risk carrying a human-factors mitigation, not an oversight, and the product had
|
It was a known, deliberately accepted risk carrying a human-factors mitigation, not an
|
||||||
already reached the same conclusion these decisions formalise. It is worth keeping in mind
|
oversight, and the product had already reached the same conclusion these decisions formalise.
|
||||||
that a warning is the weakest control available: it depends on an operator reading a dialog
|
It was also the weakest control available: it depended on an operator reading a dialog at the
|
||||||
at the end of a refill round, and it cannot help at all when the stale value is the one
|
end of a refill round, and it could not help at all when the stale value was the one already
|
||||||
already in the form. Confirmed live on 2026-09-22 — a dispense moved a bay from 54 to 53
|
in the form. Confirmed live on 2026-09-22 — a dispense moved a bay from 54 to 53 while a form
|
||||||
while a form loaded at 54 stayed open, and nothing but that dialog stood between the operator
|
loaded at 54 stayed open, and nothing but that dialog stood between the operator and
|
||||||
and discarding the decrement.
|
discarding the decrement.
|
||||||
|
|
||||||
|
The dialog and the endpoint behind it are both gone. The operator dashboard no longer has a
|
||||||
|
field that accepts a count, which is a stronger guarantee than any wording could be.
|
||||||
|
|
||||||
### 5. Ordering is decided by `created_at`, never by arrival order, on both sides.
|
### 5. Ordering is decided by `created_at`, never by arrival order, on both sides.
|
||||||
|
|
||||||
|
|
@ -138,13 +140,29 @@ eight consecutive restarts, three heartbeat republishes carried strictly increas
|
||||||
back off the relay, and the machine, the relay and the operator dashboard agreed on the counts
|
back off the relay, and the machine, the relay and the operator dashboard agreed on the counts
|
||||||
with timestamps correlated to the second.
|
with timestamps correlated to the second.
|
||||||
|
|
||||||
Decisions 1 through 4 are the v2 operations wire and are not yet built. Decision 4a describes
|
Decisions 1 through 4 are the v2 operations wire and shipped in spirekeeper#46 and
|
||||||
what stands in for them meanwhile.
|
bitspire#106.
|
||||||
|
|
||||||
|
On the operator side there is no longer any endpoint that accepts a count: the absolute
|
||||||
|
publish, its CRUD write and its request model were removed rather than deprecated. The
|
||||||
|
dashboard records operations and renders each as applied or pending from the machine's
|
||||||
|
`applied_ops` echo. On the machine side, schema v13 adds a `cassette_ops` dedup ledger, the
|
||||||
|
`created_at` watermark on this path is retired in favour of per-op ids, and the state document
|
||||||
|
carries `schema_version`, `seq` and `applied_ops`.
|
||||||
|
|
||||||
|
The wire shapes are those given under decisions 1 and 4 above.
|
||||||
|
|
||||||
Cutover for v2 is strict, no compatibility code: spirekeeper deploys first, machines follow on
|
Cutover for v2 is strict, no compatibility code: spirekeeper deploys first, machines follow on
|
||||||
their nightly pull. During that window a not-yet-updated ATM ignores an ops payload, so an
|
their nightly pull. During that window a not-yet-updated ATM ignores an ops payload, so an
|
||||||
operator refill does not land until it updates — which fails safe, since the machine
|
operator refill does not land until it updates — which fails safe, since the machine
|
||||||
under-counts and will not dispense bills it believes it lacks.
|
under-counts and will not dispense bills it believes it lacks. In the other direction an
|
||||||
|
updated machine drops a v1 absolute-count payload on the missing `ops` array, which is the
|
||||||
|
same safe direction: the machine keeps the counts it is now the only writer of.
|
||||||
|
|
||||||
|
One gap stays open deliberately. An operation recorded while the relay is unreachable waits
|
||||||
|
for the operator's next action to be published, because only an operator action triggers a
|
||||||
|
publish. The window makes that self-healing once anything is published, but nothing on the
|
||||||
|
operator side republishes on its own. An operator-side heartbeat is the fix; it is not built.
|
||||||
|
|
||||||
## Alternatives considered
|
## Alternatives considered
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue