diff --git a/src/components/bus-debug.tsx b/src/components/bus-debug.tsx new file mode 100644 index 0000000..04abda7 --- /dev/null +++ b/src/components/bus-debug.tsx @@ -0,0 +1,68 @@ +import { useSyncExternalStore } from "react"; +import { Card } from "@/components/ui/card"; +import { useBusStats } from "@/lib/bus/hooks"; +import type { Quantiles } from "@/lib/bus/manager"; + +const never = () => () => undefined; + +/** `?debug` only: the numbers the hardware procedure records. */ +export function BusDebug() { + const wanted = useSyncExternalStore( + never, + () => new URLSearchParams(location.search).has("debug"), + () => false, + ); + return wanted ? : null; +} + +function n(value: number, digits = 0): string { + return value.toFixed(digits); +} + +function q(value: Quantiles): string { + return `${n(value.p50, 1)}/${n(value.p95, 1)}`; +} + +function Panel() { + const stats = useBusStats(); + const rows: [string, string][] = [ + ["exchanges", `${stats.exchanges}`], + ["per second", n(stats.perSecond, 1)], + ["bytes/s", n(stats.bytesPerSecond)], + ["wait p50/p95", q(stats.wait)], + ["run small", q(stats.run.small)], + ["run medium", q(stats.run.medium)], + ["run large", q(stats.run.large)], + ["lag p50/p95", q(stats.lag)], + ["timeouts", `${stats.timeouts}`], + ["stalled", `${stats.stalled}`], + ["errors", `${stats.errors}`], + ["coalesced", `${stats.coalesced}`], + ["grouped", `${stats.grouped}`], + ["utilisation", n(stats.utilisation, 2)], + ["period fast/slow", `${n(stats.effectivePeriodMs.fast)}/${n(stats.effectivePeriodMs.slow)}`], + ]; + return ( + +
+ {rows.map(([label, value]) => ( +
+
{label}
+
{value}
+
+ ))} +
+ {[...stats.perServo].map(([id, s]) => ( +

+ ID {id} fails {s.consecutiveFailures} + {s.probing ? " probing" : ""} +

+ ))} +
+ ); +} diff --git a/src/components/calibration-card.tsx b/src/components/calibration-card.tsx index b8b1a47..215f006 100644 --- a/src/components/calibration-card.tsx +++ b/src/components/calibration-card.tsx @@ -9,7 +9,7 @@ import { Popover, PopoverContent, PopoverTrigger } from "@/components/ui/popover import { Skeleton } from "@/components/ui/skeleton"; import { Tooltip, TooltipContent, TooltipTrigger } from "@/components/ui/tooltip"; import { CALIBRATION_REGISTERS, editReason, type CalibrationRegister } from "@/lib/calibration"; -import { decodeSpan, span } from "@/lib/card-poll"; +import { decodeSpan, span } from "@/lib/bus/spans"; import { hexAddr } from "@/lib/format"; import { useSession } from "@/lib/session"; import { calibrationFromTable, calibrationStatus, type Calibration } from "@/lib/units"; diff --git a/src/lib/bus/hooks.ts b/src/lib/bus/hooks.ts new file mode 100644 index 0000000..d9ec8f8 --- /dev/null +++ b/src/lib/bus/hooks.ts @@ -0,0 +1,144 @@ +import type { OscClient, Value } from "@openservocore/client"; +import { createContext, useCallback, useContext, useEffect, useMemo, useState } from "react"; +import { useSyncExternalStore } from "react"; +import type { BusManager, BusStats, Rate, Snapshot } from "./manager"; +import type { BusStore } from "./store"; + +export interface BusHost { + manager: BusManager; + store: BusStore; +} + +export const BusContext = createContext(undefined); + +export function useBusHost(): BusHost { + const host = useContext(BusContext); + if (host === undefined) throw new Error("useBus outside BusContext"); + return host; +} + +function message(e: unknown): string { + return e instanceof Error ? e.message : String(e); +} + +/** A register list stable across renders that pass an equal one. */ +function useNames(registers: readonly string[]): readonly string[] { + const key = registers.join(","); + return useMemo(() => (key === "" ? [] : key.split(",")), [key]); +} + +export function useRegisters( + id: number, + registers: readonly string[], + rate: Rate, +): Snapshot | undefined { + const { manager, store } = useBusHost(); + const names = useNames(registers); + const bound = useMemo( + () => ({ + cell: store.cell(undefined), + sub: { id, registers: names, rate }, + }), + [store, id, names, rate], + ); + useEffect( + () => + manager.subscribe(bound.sub, (s) => { + bound.cell.set(s); + }), + [manager, bound], + ); + return useSyncExternalStore(bound.cell.subscribe, bound.cell.get, bound.cell.get); +} + +export function useRing( + id: number, + registers: readonly string[], + rate: Rate, + windowS: number, +): readonly Snapshot[] { + const { manager, store } = useBusHost(); + const names = useNames(registers); + const bound = useMemo( + () => ({ ring: store.ring(windowS), sub: { id, registers: names, rate } }), + [store, windowS, id, names, rate], + ); + useEffect( + () => + manager.subscribe(bound.sub, (s) => { + bound.ring.push(s); + }), + [manager, bound], + ); + return useSyncExternalStore(bound.ring.subscribe, bound.ring.get, bound.ring.get); +} + +export interface ReadOnce { + snapshot: Snapshot | undefined; + error: string | undefined; + reload: () => void; +} + +export function useReadOnce( + id: number, + registers: readonly string[], + deps: readonly unknown[] = [], +): ReadOnce { + const { manager } = useBusHost(); + const names = useNames(registers); + const key = JSON.stringify(deps); + const [generation, setGeneration] = useState(0); + const [result, setResult] = useState>({ + snapshot: undefined, + error: undefined, + }); + useEffect(() => { + let live = true; + void manager.readOnce(id, names).then( + (snapshot) => { + if (live) setResult({ snapshot, error: undefined }); + }, + (e: unknown) => { + if (live) setResult({ snapshot: undefined, error: message(e) }); + }, + ); + return () => { + live = false; + }; + }, [manager, id, names, key, generation]); + const reload = useCallback(() => { + setGeneration((g) => g + 1); + }, []); + return { ...result, reload }; +} + +export interface Bus { + write: (id: number, register: string, value: Value) => Promise; + command: (fn: (client: OscClient) => Promise) => Promise; + exclusive: (fn: (client: OscClient) => Promise) => Promise; +} + +export function useBus(): Bus { + const { manager } = useBusHost(); + return useMemo( + () => ({ + write: (id, register, value) => manager.write(id, register, value), + command: (fn) => manager.command(fn), + exclusive: (fn) => manager.exclusive(fn), + }), + [manager], + ); +} + +export function useBusStats(): BusStats { + const { manager } = useBusHost(); + const [stats, setStats] = useState(() => manager.stats()); + useEffect( + () => + manager.onStats(() => { + setStats(manager.stats()); + }), + [manager], + ); + return stats; +} diff --git a/src/lib/bus/manager.test.ts b/src/lib/bus/manager.test.ts new file mode 100644 index 0000000..d94dcbe --- /dev/null +++ b/src/lib/bus/manager.test.ts @@ -0,0 +1,666 @@ +import type { Field, OscClient, Value } from "@openservocore/client"; +import { expect, test } from "vitest"; +import { BusManager, type BusClient, type Clock, type Layout, type Snapshot } from "./manager"; + +// A small control table: goal and torque at the front, the live sensors at +// 10, one far register that cannot share a 252-byte read with the rest. +const FIELDS: Field[] = [ + { name: "goal", addr: 0, width: 2, access: "rw", kind: "int", variants: [] }, + { name: "torque", addr: 2, width: 1, access: "rw", kind: "bool", variants: [] }, + { name: "mode", addr: 3, width: 1, access: "rw", kind: "enum", variants: [] }, + { name: "pos", addr: 10, width: 2, access: "ro", kind: "uint", variants: [] }, + { name: "current", addr: 12, width: 2, access: "ro", kind: "uint", variants: [] }, + { name: "temp", addr: 200, width: 2, access: "ro", kind: "uint", variants: [] }, + { name: "far", addr: 400, width: 2, access: "ro", kind: "uint", variants: [] }, +]; + +const TABLE_SIZE = 512; +/** The manager samples event-loop lag twice a second. */ +const LAG_PROBE_MS = 500; + +function fieldOf(name: string): Field { + const f = FIELDS.find((f) => f.name === name); + if (f === undefined) throw new Error(`no field ${name}`); + return f; +} + +function decodeValue(name: string, bytes: Uint8Array): Value { + const f = fieldOf(name); + const view = new DataView(bytes.buffer, bytes.byteOffset, bytes.byteLength); + const raw = f.width === 1 ? view.getUint8(0) : view.getUint16(0, true); + switch (f.kind) { + case "int": + return { kind: "int", value: f.width === 1 ? view.getInt8(0) : view.getInt16(0, true) }; + case "bool": + return { kind: "bool", value: raw !== 0 }; + case "enum": + return { kind: "enum", value: raw }; + case "bytes": + return { kind: "bytes", value: bytes.slice() }; + case "uint": + return { kind: "uint", value: raw }; + } +} + +function encodeValue(name: string, value: Value): Uint8Array { + const f = fieldOf(name); + if (value.kind === "bytes") return value.value; + const n = typeof value.value === "boolean" ? Number(value.value) : value.value; + if (f.width === 1) return new Uint8Array([n & 0xff]); + const out = new Uint8Array(2); + new DataView(out.buffer).setInt16(0, n, true); + return out; +} + +interface TestLayout extends Layout { + decodes: string[]; +} + +function makeLayout(): TestLayout { + const decodes: string[] = []; + return { + fields: FIELDS, + decodes, + decode: (name, bytes) => { + decodes.push(name); + return decodeValue(name, bytes); + }, + encode: encodeValue, + }; +} + +interface Timer { + id: number; + at: number; + fn: () => void; +} + +/** Milliseconds, advanced by hand. */ +class FakeClock implements Clock { + t = 0; + private timers: Timer[] = []; + private nextId = 0; + + now(): number { + return this.t / 1000; + } + + after(ms: number, fn: () => void): () => void { + const id = ++this.nextId; + this.timers.push({ id, at: this.t + ms, fn }); + return () => { + this.timers = this.timers.filter((x) => x.id !== id); + }; + } + + get armed(): number[] { + return this.timers.map((x) => x.at).sort((a, b) => a - b); + } + + /** Moves the clock without firing anything: the bus was starved. */ + stall(ms: number): void { + this.t += ms; + } + + async advance(ms: number): Promise { + const target = this.t + ms; + await flush(); + for (;;) { + const due = this.timers + .filter((x) => x.at <= target) + .sort((a, b) => a.at - b.at || a.id - b.id)[0]; + if (due === undefined) break; + this.timers = this.timers.filter((x) => x !== due); + this.t = Math.max(this.t, due.at); + due.fn(); + await flush(); + } + this.t = target; + await flush(); + } +} + +async function flush(): Promise { + for (let i = 0; i < 40; i++) await Promise.resolve(); +} + +interface Call { + kind: "read" | "gread" | "write"; + ids: number[]; + addr: number; + count: number; + at: number; + data?: number[]; +} + +class FakeClient { + readonly log: Call[] = []; + latency = 0; + overlaps = 0; + /** "read:1", "read:1:10" or "write:1" to the message that call answers with. */ + readonly failures = new Map(); + /** Slots that answer nothing on a group read. */ + readonly silent = new Set(); + gread?: (ids: number[], addr: number, count: number) => Promise<(Uint8Array | undefined)[]>; + private busy = false; + private readonly tables = new Map(); + + constructor( + private readonly clock: FakeClock, + grouped = false, + ) { + if (!grouped) return; + this.gread = (ids, addr, count) => + this.settle({ kind: "gread", ids, addr, count, at: clock.t }, `read:${ids[0] ?? 0}`, () => + ids.map((id) => (this.silent.has(id) ? undefined : this.slice(id, addr, count))), + ); + } + + table(id: number): Uint8Array { + let t = this.tables.get(id); + if (t === undefined) { + t = new Uint8Array(TABLE_SIZE); + for (let i = 0; i < TABLE_SIZE; i++) t[i] = (id * 16 + i) & 0xff; + this.tables.set(id, t); + } + return t; + } + + private slice(id: number, addr: number, count: number): Uint8Array { + return this.table(id).slice(addr, addr + count); + } + + read(id: number, addr: number, count: number): Promise { + return this.settle( + { kind: "read", ids: [id], addr, count, at: this.clock.t }, + `read:${id}`, + () => this.slice(id, addr, count), + ); + } + + write(id: number, addr: number, data: Uint8Array): Promise { + return this.settle( + { kind: "write", ids: [id], addr, count: data.length, at: this.clock.t, data: [...data] }, + `write:${id}`, + () => { + this.table(id).set(data, addr); + }, + ); + } + + private settle(call: Call, key: string, action: () => T): Promise { + if (this.busy) { + this.overlaps++; + return Promise.reject(new Error("busy")); + } + this.busy = true; + this.log.push(call); + const fail = this.failures.get(`${key}:${call.addr}`) ?? this.failures.get(key); + const finish = (): T => { + this.busy = false; + if (fail !== undefined) throw new Error(fail); + return action(); + }; + if (this.latency === 0) return Promise.resolve().then(finish); + return new Promise((resolve, reject) => { + this.clock.after(this.latency, () => { + try { + resolve(finish()); + } catch (e) { + reject(e instanceof Error ? e : new Error(String(e))); + } + }); + }); + } + + reads(): [number, number, number][] { + return this.log + .filter((c) => c.kind === "read" || c.kind === "gread") + .map((c) => [c.ids[0] ?? 0, c.addr, c.count]); + } +} + +interface Bus { + clock: FakeClock; + client: FakeClient; + manager: BusManager; + layout: TestLayout; +} + +function makeBus( + options: { gread?: boolean; layoutFrom?: number } & ConstructorParameters< + typeof BusManager + >[1] = {}, +): Bus { + const { gread = false, layoutFrom, ...rest } = options; + const clock = new FakeClock(); + const client = new FakeClient(clock, gread); + const layout = makeLayout(); + const manager = new BusManager(clock, rest); + manager.attach(client as unknown as OscClient & BusClient, () => + layoutFrom === undefined ? layout : undefined, + ); + return { clock, client, manager, layout }; +} + +function collect(): { seen: Snapshot[]; listener: (s: Snapshot) => void } { + const seen: Snapshot[] = []; + return { seen, listener: (s) => seen.push(s) }; +} + +test("three servos, fast and slow subscriptions, writes and commands interleaved: the client never sees an overlap", async () => { + const { clock, client, manager } = makeBus(); + client.latency = 7; + for (const id of [1, 2, 3]) { + manager.subscribe({ id, registers: ["pos", "current"], rate: "fast" }, () => undefined); + manager.subscribe({ id, registers: ["temp"], rate: "slow" }, () => undefined); + } + const promises: Promise[] = []; + for (let i = 0; i < 6; i++) { + promises.push(manager.write(1 + (i % 3), "goal", { kind: "int", value: i })); + promises.push(manager.command((c) => c.read(2, 10, 2))); + await clock.advance(23); + } + await clock.advance(500); + await Promise.all(promises); + expect(client.overlaps).toBe(0); + expect(client.log.length).toBeGreaterThan(12); +}); + +test("a write submitted mid-read starts when that read settles, ahead of every due read", async () => { + const { clock, client, manager } = makeBus(); + client.latency = 50; + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + manager.subscribe({ id: 2, registers: ["pos"], rate: "fast" }, () => undefined); + await clock.advance(0); + expect(client.log).toHaveLength(1); + const write = manager.write(1, "goal", { kind: "int", value: 42 }); + await clock.advance(200); + await write; + expect(client.log[1]?.kind).toBe("write"); +}); + +test("writes to one register while one is in flight collapse to the newest; both promises settle together", async () => { + const { clock, client, manager } = makeBus(); + client.latency = 20; + const first = manager.write(1, "goal", { kind: "int", value: 1 }); + await flush(); + const second = manager.write(1, "goal", { kind: "int", value: 2 }); + const third = manager.write(1, "goal", { kind: "int", value: 3 }); + await clock.advance(20); + await first; + const order: string[] = []; + void second.then(() => order.push("second")); + void third.then(() => order.push("third")); + await clock.advance(20); + await Promise.all([second, third]); + expect(order).toEqual(["second", "third"]); + const writes = client.log.filter((c) => c.kind === "write"); + expect(writes).toHaveLength(2); + expect(writes[1]?.data).toEqual([3, 0]); + expect(manager.stats().coalesced).toBe(1); +}); + +test("writes to different registers of one servo keep submission order", async () => { + const { clock, client, manager } = makeBus(); + client.latency = 5; + const first = manager.write(1, "goal", { kind: "int", value: 1 }); + await flush(); + const rest = Promise.all([ + manager.write(1, "torque", { kind: "bool", value: false }), + manager.write(1, "mode", { kind: "enum", value: 2 }), + manager.write(1, "goal", { kind: "int", value: 2 }), + ]); + await clock.advance(200); + await Promise.all([first, rest]); + const writes = client.log.filter((c) => c.kind === "write").map((c) => c.addr); + expect(writes).toEqual([0, 2, 3, 0]); +}); + +test("a command is never coalesced with a write or reordered past one", async () => { + const { clock, client, manager } = makeBus(); + client.latency = 5; + const first = manager.write(1, "goal", { kind: "int", value: 1 }); + await flush(); + const rest = Promise.all([ + manager.write(1, "goal", { kind: "int", value: 2 }), + manager.command((c) => c.read(1, 200, 2)), + manager.write(1, "goal", { kind: "int", value: 3 }), + ]); + await clock.advance(200); + await Promise.all([first, rest]); + const seen = client.log.map((c) => `${c.kind}:${c.addr}`); + expect(seen.slice(0, 3)).toEqual(["write:0", "write:0", "read:200"]); + expect(client.log[1]?.data).toEqual([3, 0]); +}); + +test("a write dirties its field; the refresh read runs before the next live read and reaches every subscriber covering the field", async () => { + const { clock, client, manager } = makeBus(); + const live = collect(); + const dirty = collect(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, live.listener); + manager.subscribe({ id: 1, registers: ["goal"], rate: "slow" }, dirty.listener); + await clock.advance(1); + const before = client.reads().length; + await manager.write(1, "goal", { kind: "int", value: 77 }); + await flush(); + const after = client.reads().slice(before); + expect(after[0]).toEqual([1, 0, 2]); + expect(dirty.seen.at(-1)?.read("goal")).toBe(77); +}); + +test("a dirty field inside a live span due within a fast period is served by that live read: one exchange", async () => { + const { clock, client, manager } = makeBus(); + manager.subscribe({ id: 1, registers: ["goal", "pos"], rate: "fast" }, () => undefined); + await clock.advance(1); + const before = client.reads().length; + await manager.write(1, "goal", { kind: "int", value: 5 }); + await flush(); + expect(client.reads().slice(before)).toEqual([[1, 0, 12]]); +}); + +test("registers of one servo and rate merge into the fewest spans under 252 bytes, gaps read through", async () => { + const { clock, client, manager } = makeBus(); + manager.subscribe( + { id: 1, registers: ["goal", "pos", "temp", "far"], rate: "fast" }, + () => undefined, + ); + await clock.advance(1); + expect(client.reads()).toEqual([ + [1, 0, 202], + [1, 400, 2], + ]); +}); + +test("a slow field inside a fast span costs no exchange of its own", async () => { + const { clock, client, manager } = makeBus(); + manager.subscribe({ id: 1, registers: ["goal", "temp"], rate: "fast" }, () => undefined); + manager.subscribe({ id: 1, registers: ["pos"], rate: "slow" }, () => undefined); + await clock.advance(1000); + expect(new Set(client.reads().map(([, addr]) => addr))).toEqual(new Set([0])); +}); + +test("subscribers on one span receive the same Snapshot object; values are decoded once", async () => { + const { clock, client, manager, layout } = makeBus(); + const a = collect(); + const b = collect(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, a.listener); + manager.subscribe({ id: 1, registers: ["pos", "current"], rate: "fast" }, b.listener); + await clock.advance(1); + expect(client.reads()).toEqual([[1, 10, 4]]); + expect(a.seen).toHaveLength(1); + expect(a.seen[0]).toBe(b.seen[0]); + expect(layout.decodes.filter((n) => n === "pos")).toHaveLength(1); +}); + +test("a 1 s stall of the bus is followed by one read per span, not ten", async () => { + const { clock, client, manager } = makeBus(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + manager.subscribe({ id: 1, registers: ["far"], rate: "fast" }, () => undefined); + await clock.advance(1); + const before = client.reads().length; + clock.stall(1000); + await clock.advance(0); + expect(client.reads().length - before).toBe(2); +}); + +test("with latency above the fast period the slow class still gets its turn (earliest due first)", async () => { + const { clock, client, manager } = makeBus(); + client.latency = 150; + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + manager.subscribe({ id: 1, registers: ["far"], rate: "slow" }, () => undefined); + await clock.advance(3000); + const addrs = client.reads().map(([, addr]) => addr); + expect(addrs.filter((a) => a === 400).length).toBeGreaterThan(0); + expect(addrs.filter((a) => a === 10).length).toBeGreaterThan(1); +}); + +test("utilisation over target lengthens the fast period; two cycles under target shorten it back", async () => { + const { clock, client, manager } = makeBus({ fastMs: 100, target: 0.6 }); + client.latency = 70; + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + await clock.advance(100); + expect(manager.stats().utilisation).toBeCloseTo(0.7, 3); + expect(manager.stats().effectivePeriodMs.fast).toBeCloseTo(100 * (0.7 / 0.6), 3); + client.latency = 0; + await clock.advance(1000); + expect(manager.stats().effectivePeriodMs.fast).toBe(100); +}); + +test("nothing due: the clock is armed for the earliest due and no read is issued meanwhile", async () => { + const { clock, client, manager } = makeBus(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + await clock.advance(0); + expect(client.reads()).toHaveLength(1); + expect(clock.armed).toContain(100); + await clock.advance(99); + expect(client.reads()).toHaveLength(1); + await clock.advance(1); + expect(client.reads()).toHaveLength(2); +}); + +test("the same span on three servos goes out as one gread; a silent slot marks only its servo stale", async () => { + const { clock, client, manager } = makeBus({ gread: true }); + client.silent.add(2); + const seen = new Map(); + for (const id of [1, 2, 3]) { + manager.subscribe({ id, registers: ["pos"], rate: "fast" }, (s) => seen.set(id, s)); + } + await clock.advance(1); + expect(client.log.filter((c) => c.kind === "gread")).toHaveLength(1); + expect(client.log[0]?.ids).toEqual([1, 2, 3]); + expect(seen.get(1)?.stale).toBe(false); + expect(seen.get(2)?.stale).toBe(true); + expect(seen.get(3)?.stale).toBe(false); + expect(manager.stats().grouped).toBe(1); +}); + +test("without `gread` the same subscriptions produce per-servo reads and identical snapshots", async () => { + const grouped = makeBus({ gread: true }); + const single = makeBus(); + const from = async (bus: Bus): Promise<[number, number][]> => { + const out: [number, number][] = []; + for (const id of [1, 2, 3]) { + bus.manager.subscribe({ id, registers: ["pos"], rate: "fast" }, (s) => { + out.push([id, s.read("pos")]); + }); + } + await bus.clock.advance(1); + return out.sort((a, b) => a[0] - b[0]); + }; + const a = await from(grouped); + const b = await from(single); + expect(single.client.log.filter((c) => c.kind === "read")).toHaveLength(3); + expect(grouped.client.log.filter((c) => c.kind === "gread")).toHaveLength(1); + expect(b).toEqual(a); +}); + +test("a failed read marks the snapshot stale, keeps the subscription, three failures switch the servo to the probe cadence, one success restores it", async () => { + const { clock, client, manager } = makeBus(); + const seen = collect(); + const straddle = collect(); + client.failures.set("read:1", "no reply"); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, seen.listener); + // pos and far cannot share a 252-byte read, so this one straddles two spans. + manager.subscribe({ id: 1, registers: ["pos", "far"], rate: "fast" }, straddle.listener); + await clock.advance(250); + expect(seen.seen.at(-1)?.stale).toBe(true); + expect(seen.seen.at(-1)?.error).toBe("no reply"); + expect(straddle.seen.at(-1)?.stale).toBe(true); + expect(straddle.seen.at(-1)?.error).toBe("no reply"); + expect(manager.stats().perServo.get(1)?.consecutiveFailures).toBeGreaterThanOrEqual(3); + expect(manager.stats().perServo.get(1)?.probing).toBe(true); + client.failures.clear(); + const before = client.reads().length; + await clock.advance(1000); + expect(client.reads()[before]).toEqual([1, 0, 4]); + expect(manager.stats().perServo.get(1)?.probing).toBe(false); + await clock.advance(200); + expect(client.reads().slice(before + 1)).toContainEqual([1, 10, 2]); + expect(seen.seen.at(-1)?.stale).toBe(false); + expect(straddle.seen.at(-1)?.stale).toBe(false); +}); + +test("a subscriber mounting after a span's read is served from the cache without an exchange; the next read arrives on the span's own due", async () => { + const { clock, client, manager } = makeBus(); + manager.subscribe({ id: 1, registers: ["pos", "current"], rate: "slow" }, () => undefined); + await clock.advance(0); + expect(client.reads()).toEqual([[1, 10, 4]]); + await clock.advance(400); + const late = collect(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "slow" }, late.listener); + await clock.advance(0); + const table = client.table(1); + expect(late.seen).toHaveLength(1); + expect(late.seen[0]?.read("pos")).toBe((table[10] ?? 0) | ((table[11] ?? 0) << 8)); + expect(late.seen[0]?.stale).toBe(false); + // Dated by the read it came from, not by the moment it was served. + expect(late.seen[0]?.seq).toBe(1); + expect(client.reads()).toHaveLength(1); + await clock.advance(700); + expect(client.reads()).toHaveLength(2); + expect(late.seen).toHaveLength(2); +}); + +test("a timeout inserts the pacing pause before the next exchange", async () => { + const { clock, client, manager } = makeBus(); + client.failures.set("read:1", "read timeout"); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + manager.subscribe({ id: 2, registers: ["pos"], rate: "fast" }, () => undefined); + await clock.advance(0); + expect(client.reads()).toHaveLength(1); + await clock.advance(1); + expect(client.reads()).toHaveLength(2); + expect(manager.stats().timeouts).toBe(1); +}); + +test("unsubscribe during an in-flight read: the cache updates, no listener fires", async () => { + const { clock, client, manager } = makeBus(); + client.latency = 30; + const seen = collect(); + const stop = manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, seen.listener); + await clock.advance(0); + stop(); + await clock.advance(30); + expect(seen.seen).toHaveLength(0); + const after = collect(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, after.listener); + manager.detach("unplugged"); + expect(after.seen.at(-1)?.stale).toBe(true); + const table = client.table(1); + expect(after.seen.at(-1)?.read("pos")).toBe((table[10] ?? 0) | ((table[11] ?? 0) << 8)); +}); + +test('detach rejects pending writes with "not connected", lets the in-flight settle, marks snapshots stale; attach resumes the same subscriptions', async () => { + const { clock, client, manager } = makeBus(); + client.latency = 25; + const seen = collect(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, seen.listener); + await clock.advance(0); + const pending = manager.write(1, "goal", { kind: "int", value: 9 }); + manager.detach("unplugged"); + await expect(pending).rejects.toThrow("not connected"); + expect(seen.seen.at(-1)?.stale).toBe(true); + await clock.advance(30); + const before = client.log.length; + manager.attach(client as unknown as OscClient & BusClient, () => makeLayout()); + await clock.advance(30); + expect(client.log.length).toBeGreaterThan(before); + expect(seen.seen.at(-1)?.stale).toBe(false); +}); + +test("exclusive runs after the in-flight exchange, freezes the lanes, and due times restart without a burst", async () => { + const { clock, client, manager } = makeBus(); + client.latency = 40; + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + manager.subscribe({ id: 2, registers: ["pos"], rate: "fast" }, () => undefined); + await clock.advance(0); + const during: number[] = []; + const job = manager.exclusive(async (c) => { + during.push(client.log.length); + await c.read(3, 0, 2); + during.push(client.log.length); + return "done"; + }); + await clock.advance(40); + await clock.advance(40); + expect(await job).toBe("done"); + // The scan's own read is the only exchange the turn allowed. + expect(during[1]).toBe((during[0] ?? 0) + 1); + const after = client.reads().length; + await clock.advance(99); + expect(client.reads()).toHaveLength(after); + await clock.advance(400); + expect(client.reads().length).toBeGreaterThan(after); +}); + +test("roster drops a vanished id from scheduling; its subscriber sees stale", async () => { + const { clock, client, manager } = makeBus(); + const gone = collect(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + manager.subscribe({ id: 2, registers: ["pos"], rate: "fast" }, gone.listener); + await clock.advance(1); + manager.roster([1]); + expect(gone.seen.at(-1)?.stale).toBe(true); + const before = client.reads().filter(([id]) => id === 2).length; + await clock.advance(1000); + expect(client.reads().filter(([id]) => id === 2)).toHaveLength(before); +}); + +test("a validation rejection rejects the write and dirties nothing", async () => { + const { clock, client, manager } = makeBus(); + client.failures.set("write:1", "validation"); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + await clock.advance(1); + const before = client.reads().length; + await expect(manager.write(1, "goal", { kind: "int", value: 1 })).rejects.toThrow("validation"); + await flush(); + expect(client.reads()).toHaveLength(before); +}); + +test("stats record wait, run, lag, stalled, coalesced and grouped counts and the effective periods", async () => { + const { clock, client, manager } = makeBus({ gread: true }); + client.latency = 12; + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, () => undefined); + manager.subscribe({ id: 2, registers: ["pos"], rate: "fast" }, () => undefined); + // The lag probe measures setTimeout(0) drift: stalling the clock in the + // instant the probe arms makes the drift exactly the stall. + clock.after(LAG_PROBE_MS, () => { + clock.stall(40); + }); + await clock.advance(0); + void manager.write(1, "goal", { kind: "int", value: 1 }); + void manager.write(1, "goal", { kind: "int", value: 2 }); + await clock.advance(60); + const busy = manager.stats(); + expect(busy.wait.p95).toBeGreaterThan(0); + expect(busy.run.small.p50).toBeGreaterThan(0); + expect(busy.coalesced).toBe(1); + expect(busy.grouped).toBeGreaterThan(0); + client.failures.set("read:1", "recv guard expired"); + await clock.advance(2000); + const stats = manager.stats(); + expect(stats.exchanges).toBeGreaterThan(busy.exchanges); + expect(stats.lag.p95).toBeGreaterThanOrEqual(40); + expect(stats.stalled).toBeGreaterThan(0); + expect(stats.effectivePeriodMs.slow).toBe(1000); + expect(stats.recent.length).toBe(Math.min(256, stats.exchanges)); +}); + +test("a subscription before its layout loads is planned on `layoutChanged`", async () => { + const clock = new FakeClock(); + const client = new FakeClient(clock); + const manager = new BusManager(clock); + const held: { layout: Layout | undefined } = { layout: undefined }; + manager.attach(client as unknown as OscClient & BusClient, () => held.layout); + const seen = collect(); + manager.subscribe({ id: 1, registers: ["pos"], rate: "fast" }, seen.listener); + await clock.advance(500); + expect(client.log).toHaveLength(0); + await expect(manager.readOnce(1, ["pos"])).rejects.toThrow("no layout for ID 1"); + held.layout = makeLayout(); + manager.layoutChanged(1); + await clock.advance(1); + expect(client.reads()).toEqual([[1, 10, 2]]); + expect(seen.seen).toHaveLength(1); +}); diff --git a/src/lib/bus/manager.ts b/src/lib/bus/manager.ts new file mode 100644 index 0000000..68dd359 --- /dev/null +++ b/src/lib/bus/manager.ts @@ -0,0 +1,953 @@ +// One scheduler owns the adapter. Pages declare what they need and this +// decides what goes on the wire and when; the client reference never leaves +// the private field below, so an overlapping call cannot be written. + +import type { Field, OscClient, Value } from "@openservocore/client"; +import type { ReadRegister } from "../units"; +import { + decodeValues, + field, + planSpans, + readerOver, + within, + type Layout, + type Span, +} from "./spans"; +import { + classify, + StatsRecorder, + type BusStats, + type Exchange, + type Kind, + type Outcome, +} from "./stats"; + +export type { BusStats, Exchange, Quantiles } from "./stats"; +export type { Layout, Span } from "./spans"; + +export type Rate = "fast" | "slow"; + +export interface Subscription { + id: number; + registers: readonly string[]; + rate: Rate; +} + +export interface Snapshot { + id: number; + /** Exchange sequence number, monotonic per manager. */ + seq: number; + /** clock.now() at the exchange's start, seconds. */ + t: number; + /** Every field the read covered, decoded once. */ + values: ReadonlyMap; + /** Numeric accessor over `values`; throws for a bytes field or a name outside the read. */ + read: ReadRegister; + /** The last read of this span failed, or the servo is silent, or the client is detached. */ + stale: boolean; + error?: string; +} + +/** The slice of OscClient the scheduler drives; `gread` is optional. */ +export interface BusClient { + read(id: number, addr: number, count: number): Promise; + write(id: number, addr: number, data: Uint8Array): Promise; + gread?(ids: number[], addr: number, count: number): Promise<(Uint8Array | undefined)[]>; +} + +export interface Clock { + /** Seconds. */ + now(): number; + /** Runs `fn` once after `ms`; returns the cancel. */ + after(ms: number, fn: () => void): () => void; +} + +export const systemClock: Clock = { + now: () => performance.now() / 1000, + after: (ms, fn) => { + const timer = setTimeout(fn, ms); + return () => { + clearTimeout(timer); + }; + }, +}; + +const FAST_MS = 100; +const SLOW_MS = 1000; +const TARGET = 0.6; +const FAST_CAP_MS = 1000; +const SLOW_CAP_MS = 4000; +/** Consecutive failures that drop a servo to the probe cadence. */ +const PROBE_FAILURES = 3; +const PROBE_MS = 1000; +const PROBE: Span = { addr: 0, count: 4 }; +/** Fault pacing after a timeout (protocol sec 8). */ +const PACE_MS = 1; +const LAG_PROBE_MS = 500; +const STATS_MS = 500; + +interface Sub { + sub: Subscription; + listener: (s: Snapshot) => void; + fields: Field[] | undefined; +} + +interface PlannedSpan extends Span { + rate: Rate; + /** clock time in milliseconds at which this span wants its next read. */ + due: number; +} + +/** One register's last reading: the short-lived cache line a span leaves behind. */ +interface Cached { + value: Value; + t: number; + seq: number; +} + +interface ServoState { + spans: PlannedSpan[]; + probe: PlannedSpan; + cache: Map; + failures: number; + probing: boolean; +} + +interface Settle { + resolve: (value: T) => void; + reject: (error: unknown) => void; +} + +interface WriteItem { + kind: "write"; + id: number; + register: string; + value: Value; + settle: Settle[]; + queuedAt: number; +} + +interface CommandItem { + kind: "command"; + fn: (client: OscClient) => Promise; + settle: Settle; + queuedAt: number; +} + +type ControlItem = WriteItem | CommandItem; + +interface ExclusiveItem { + fn: (client: OscClient) => Promise; + settle: Settle; + queuedAt: number; +} + +interface ReadResult { + values: Map; + seq: number; + t: number; + error?: string; +} + +interface RefreshItem { + id: number; + addr: number; + count: number; + queuedAt: number; + done?: (result: ReadResult) => void; +} + +interface ReadTarget { + id: number; + span: PlannedSpan | undefined; +} + +interface ReadJob { + targets: ReadTarget[]; + addr: number; + count: number; + queuedAt: number; + done?: (result: ReadResult) => void; +} + +type Job = + | { lane: "exclusive"; item: ExclusiveItem } + | { lane: "control"; item: ControlItem } + | { lane: "refresh" | "live"; read: ReadJob }; + +function message(e: unknown): string { + return e instanceof Error ? e.message : String(e); +} + +export class BusManager { + private readonly clock: Clock; + private readonly fastMs: number; + private readonly slowMs: number; + private readonly target: number; + + private client: (OscClient & BusClient) | undefined; + private layoutOf: ((id: number) => Layout | undefined) | undefined; + private rosterIds: ReadonlySet | undefined; + + private readonly subs = new Set(); + private readonly fresh: Sub[] = []; + private readonly servos = new Map(); + private readonly control: ControlItem[] = []; + private readonly refresh: RefreshItem[] = []; + private readonly exclusives: ExclusiveItem[] = []; + + private readonly recorder = new StatsRecorder(); + private readonly statsListeners = new Set<() => void>(); + private stopTimer: (() => void) | undefined; + private stopProbe: (() => void) | undefined; + private stopStats: (() => void) | undefined; + + private inFlight = false; + private dispatching = false; + private seq = 0; + private utilisation = 0; + private under = 0; + private effective: { fast: number; slow: number }; + + constructor(clock: Clock, options: { fastMs?: number; slowMs?: number; target?: number } = {}) { + this.clock = clock; + this.fastMs = options.fastMs ?? FAST_MS; + this.slowMs = options.slowMs ?? SLOW_MS; + this.target = options.target ?? TARGET; + this.effective = { fast: this.fastMs, slow: this.slowMs }; + } + + // Connection + + attach(client: OscClient & BusClient, layout: (id: number) => Layout | undefined): void { + this.client = client; + this.layoutOf = layout; + const now = this.nowMs(); + for (const id of this.servos.keys()) { + const state = this.state(id); + state.failures = 0; + state.probing = false; + // A new client knows nothing: the old readings must not be served as fresh. + state.cache.clear(); + this.replan(id); + for (const span of state.spans) span.due = now; + } + this.startLagProbe(); + this.schedulePump(); + } + + detach(reason: string): void { + this.client = undefined; + this.layoutOf = undefined; + this.arm(undefined); + this.stopProbe?.(); + this.stopProbe = undefined; + const pending = this.control.splice(0); + const refresh = this.refresh.splice(0); + const exclusives = this.exclusives.splice(0); + for (const item of pending) this.rejectControl(item, new Error("not connected")); + for (const item of refresh) + item.done?.({ values: new Map(), seq: this.seq, t: 0, error: reason }); + for (const item of exclusives) item.settle.reject(new Error("not connected")); + for (const entry of this.subs) entry.listener(this.fromCache(entry, true, reason)); + } + + layoutChanged(id: number): void { + this.replan(id); + this.schedulePump(); + } + + /** After a scan: subscriptions on ids not listed stop scheduling and report `stale`. */ + roster(ids: readonly number[]): void { + this.rosterIds = new Set(ids); + for (const id of this.servos.keys()) this.replan(id); + for (const entry of this.subs) { + if (this.scheduled(entry.sub.id)) continue; + entry.listener(this.fromCache(entry, true, "not on the bus")); + } + this.schedulePump(); + } + + // Declarations + + subscribe(sub: Subscription, listener: (s: Snapshot) => void): () => void { + const entry: Sub = { sub, listener, fields: undefined }; + this.subs.add(entry); + this.fresh.push(entry); + this.replan(sub.id); + this.schedulePump(); + return () => { + this.subs.delete(entry); + this.replan(sub.id); + }; + } + + readOnce(id: number, registers: readonly string[]): Promise { + const layout = this.layoutOf?.(id); + if (layout === undefined) return Promise.reject(new Error(`no layout for ID ${id}`)); + let spans: Span[]; + try { + spans = planSpans(layout.fields, registers); + } catch (e) { + return Promise.reject(e instanceof Error ? e : new Error(message(e))); + } + const now = this.nowMs(); + return new Promise((resolve, reject) => { + const values = new Map(); + let left = spans.length; + let failed = false; + let seq = this.seq; + let t = this.clock.now(); + if (left === 0) { + resolve({ id, seq, t, values, read: readerOver(values), stale: false }); + return; + } + for (const s of spans) { + this.refresh.push({ + id, + addr: s.addr, + count: s.count, + queuedAt: now, + done: (result) => { + if (failed) return; + if (result.error !== undefined) { + failed = true; + reject(new Error(result.error)); + return; + } + for (const [name, value] of result.values) values.set(name, value); + seq = result.seq; + t = result.t; + if (--left > 0) return; + resolve({ id, seq, t, values, read: readerOver(values), stale: false }); + }, + }); + } + this.pump(); + }); + } + + write(id: number, register: string, value: Value): Promise { + return new Promise((resolve, reject) => { + if (this.client === undefined) { + reject(new Error("not connected")); + return; + } + const pending = this.control.find( + (i): i is WriteItem => i.kind === "write" && i.id === id && i.register === register, + ); + if (pending !== undefined) { + pending.value = value; + pending.settle.push({ resolve, reject }); + this.recorder.coalesce(); + return; + } + this.control.push({ + kind: "write", + id, + register, + value, + settle: [{ resolve, reject }], + queuedAt: this.nowMs(), + }); + this.pump(); + }); + } + + command(fn: (client: OscClient) => Promise): Promise { + return this.enqueue(fn, (item) => { + this.control.push({ kind: "command", ...item }); + }); + } + + exclusive(fn: (client: OscClient) => Promise): Promise { + return this.enqueue(fn, (item) => { + this.exclusives.push(item); + }); + } + + private enqueue( + fn: (client: OscClient) => Promise, + push: (item: ExclusiveItem) => void, + ): Promise { + return new Promise((resolve, reject) => { + if (this.client === undefined) { + reject(new Error("not connected")); + return; + } + push({ + fn, + settle: { + resolve: (value) => { + resolve(value as T); + }, + reject, + }, + queuedAt: this.nowMs(), + }); + this.pump(); + }); + } + + // Statistics + + stats(): BusStats { + const perServo = new Map(); + for (const [id, state] of this.servos) { + perServo.set(id, { consecutiveFailures: state.failures, probing: state.probing }); + } + return this.recorder.build({ + utilisation: this.utilisation, + effectivePeriodMs: { ...this.effective }, + perServo, + }); + } + + onStats(listener: () => void): () => void { + this.statsListeners.add(listener); + if (this.statsListeners.size === 1) this.tickStats(); + return () => { + this.statsListeners.delete(listener); + if (this.statsListeners.size > 0) return; + this.stopStats?.(); + this.stopStats = undefined; + }; + } + + private tickStats(): void { + this.stopStats = this.clock.after(STATS_MS, () => { + for (const listener of this.statsListeners) listener(); + if (this.statsListeners.size > 0) this.tickStats(); + }); + } + + private startLagProbe(): void { + this.stopProbe?.(); + const tick = (): void => { + const at = this.clock.now(); + this.clock.after(0, () => { + this.recorder.lag((this.clock.now() - at) * 1000); + }); + this.stopProbe = this.clock.after(LAG_PROBE_MS, tick); + }; + tick(); + } + + // Planning + + private nowMs(): number { + return this.clock.now() * 1000; + } + + private period(rate: Rate): number { + return rate === "fast" ? this.effective.fast : this.effective.slow; + } + + private scheduled(id: number): boolean { + return this.rosterIds === undefined || this.rosterIds.has(id); + } + + private state(id: number): ServoState { + let state = this.servos.get(id); + if (state === undefined) { + state = { + spans: [], + probe: { ...PROBE, rate: "slow", due: this.nowMs() }, + cache: new Map(), + failures: 0, + probing: false, + }; + this.servos.set(id, state); + } + return state; + } + + private spansOf(state: ServoState): PlannedSpan[] { + return state.probing ? [state.probe] : state.spans; + } + + /** The servo's spans, merged per rate class; a slow field inside a fast span gets none. */ + private replan(id: number): void { + const state = this.state(id); + const layout = this.layoutOf?.(id); + const mine = [...this.subs].filter((e) => e.sub.id === id); + if (mine.length === 0 || layout === undefined) { + state.spans = []; + return; + } + const names = (rate: Rate) => + mine.filter((e) => e.sub.rate === rate).flatMap((e) => e.sub.registers); + const fast = planSpans(layout.fields, names("fast")); + const slow = planSpans(layout.fields, names("slow"), fast); + const now = this.nowMs(); + const keep = (s: Span, rate: Rate): PlannedSpan => { + const prior = state.spans.find( + (p) => p.addr === s.addr && p.count === s.count && p.rate === rate, + ); + return { ...s, rate, due: prior?.due ?? now }; + }; + state.spans = [...fast.map((s) => keep(s, "fast")), ...slow.map((s) => keep(s, "slow"))]; + for (const entry of mine) + entry.fields = [...new Set(entry.sub.registers)].map((n) => field(layout.fields, n)); + this.restretch(); + } + + /** Effective periods follow the utilisation the subscriptions ask for (sec 3.6). */ + private restretch(): void { + let u = 0; + for (const [id, state] of this.servos) { + if (!this.scheduled(id)) continue; + for (const s of this.spansOf(state)) { + u += this.recorder.estimate(s.count) / (s.rate === "fast" ? this.fastMs : this.slowMs); + } + } + this.utilisation = u; + if (u > this.target) { + this.under = 0; + const fast = Math.min((this.fastMs * u) / this.target, FAST_CAP_MS); + this.effective = { + fast, + slow: + fast >= FAST_CAP_MS + ? Math.min((this.slowMs * u) / this.target, SLOW_CAP_MS) + : this.slowMs, + }; + return; + } + if (++this.under >= 2) this.effective = { fast: this.fastMs, slow: this.slowMs }; + } + + // Dispatch + + /** + * Plans settle before the wire does: subscriptions mounted in one tick + * merge into one span instead of each firing a read of its own. + */ + private schedulePump(): void { + this.arm( + this.clock.after(0, () => { + this.stopTimer = undefined; + this.pump(); + }), + ); + } + + private pump(): void { + this.serveFresh(); + if (this.inFlight || this.client === undefined) return; + const job = this.pick(); + if (job === undefined) { + this.armNext(); + return; + } + this.arm(undefined); + this.inFlight = true; + void this.execute(job); + } + + private pick(): Job | undefined { + const exclusive = this.exclusives.shift(); + if (exclusive !== undefined) return { lane: "exclusive", item: exclusive }; + const control = this.control.shift(); + if (control !== undefined) return { lane: "control", item: control }; + const now = this.nowMs(); + const item = this.refresh[0]; + if (item !== undefined) { + this.refresh.shift(); + const live = this.liveCovering(item, now); + if (live !== undefined) { + live.due = now + this.period(live.rate); + return { + lane: "live", + read: { + targets: [{ id: item.id, span: live }], + addr: live.addr, + count: live.count, + queuedAt: item.queuedAt, + done: item.done, + }, + }; + } + return { + lane: "refresh", + read: { + targets: [{ id: item.id, span: undefined }], + addr: item.addr, + count: item.count, + queuedAt: item.queuedAt, + done: item.done, + }, + }; + } + return this.pickLive(now); + } + + /** Rule 4: a live span that already covers the dirty field and is due soon serves both. */ + private liveCovering(item: RefreshItem, now: number): PlannedSpan | undefined { + const state = this.servos.get(item.id); + if (state === undefined || !this.scheduled(item.id)) return undefined; + return this.spansOf(state).find( + (s) => + s.addr <= item.addr && + s.addr + s.count >= item.addr + item.count && + s.due <= now + this.effective.fast, + ); + } + + private pickLive(now: number): Job | undefined { + let best: { id: number; span: PlannedSpan } | undefined; + for (const [id, state] of this.servos) { + if (!this.scheduled(id)) continue; + for (const span of this.spansOf(state)) { + if (span.due > now) continue; + if (best === undefined || span.due < best.span.due) best = { id, span }; + } + } + if (best === undefined) return undefined; + const chosen = best; + const targets: ReadTarget[] = [{ id: chosen.id, span: chosen.span }]; + if (this.client?.gread !== undefined) { + for (const [id, state] of this.servos) { + if (id === chosen.id || !this.scheduled(id)) continue; + const span = this.spansOf(state).find( + (s) => + s.addr === chosen.span.addr && + s.count === chosen.span.count && + s.due <= now + this.period(s.rate), + ); + if (span !== undefined) targets.push({ id, span }); + } + } + for (const t of targets) { + if (t.span !== undefined) t.span.due = now + this.period(t.span.rate); + } + return { + lane: "live", + read: { targets, addr: chosen.span.addr, count: chosen.span.count, queuedAt: now }, + }; + } + + private async execute(job: Job): Promise { + const client = this.client; + if (client === undefined) { + this.inFlight = false; + return; + } + // Development assertion: two client calls can never overlap (sec 4). + if (this.dispatching) throw new Error("bus dispatch re-entered"); + this.dispatching = true; + const seq = ++this.seq; + const startedAt = this.nowMs(); + const t = this.clock.now(); + let outcome: Outcome = "ok"; + let kind: Kind = "command"; + let bytes = 0; + let id: number | undefined; + let addr: number | undefined; + let queuedAt = startedAt; + try { + switch (job.lane) { + case "exclusive": { + queuedAt = job.item.queuedAt; + try { + job.item.settle.resolve(await job.item.fn(client)); + } catch (e) { + outcome = classify(e); + job.item.settle.reject(e); + } + this.restart(); + break; + } + case "control": { + queuedAt = job.item.queuedAt; + if (job.item.kind === "command") { + try { + job.item.settle.resolve(await job.item.fn(client)); + } catch (e) { + outcome = classify(e); + job.item.settle.reject(e); + } + break; + } + kind = "write"; + id = job.item.id; + const item = job.item; + let f: Field; + let data: Uint8Array; + try { + const layout = this.layoutOf?.(item.id); + if (layout === undefined) throw new Error(`no layout for ID ${item.id}`); + f = field(layout.fields, item.register); + data = layout.encode(item.register, item.value); + } catch (e) { + outcome = "error"; + this.rejectControl(item, e); + break; + } + addr = f.addr; + bytes = data.length; + try { + await client.write(item.id, f.addr, data); + for (const s of item.settle) s.resolve(); + this.markDirty(item.id, f); + } catch (e) { + outcome = classify(e); + this.rejectControl(item, e); + } + break; + } + case "refresh": + case "live": { + const { read } = job; + queuedAt = read.queuedAt; + kind = read.targets.length > 1 ? "gread" : "read"; + addr = read.addr; + id = read.targets.length === 1 ? read.targets[0]?.id : undefined; + bytes = read.count * read.targets.length; + outcome = await this.runRead(client, read, seq, t); + break; + } + } + } finally { + this.dispatching = false; + } + const settledAt = this.nowMs(); + const exchange: Exchange = { + seq, + lane: job.lane, + kind, + id, + addr, + bytes, + queuedAt, + startedAt, + settledAt, + outcome, + }; + this.recorder.record(exchange); + this.restretch(); + this.inFlight = false; + if (outcome === "timeout") { + this.clock.after(PACE_MS, () => { + this.pump(); + }); + return; + } + this.pump(); + } + + private async runRead( + client: OscClient & BusClient, + read: ReadJob, + seq: number, + t: number, + ): Promise { + const gread = client.gread?.bind(client); + if (read.targets.length > 1 && gread !== undefined) { + this.recorder.group(); + let parts: (Uint8Array | undefined)[]; + try { + parts = await gread( + read.targets.map((x) => x.id), + read.addr, + read.count, + ); + } catch (e) { + for (const target of read.targets) + this.fail(target.id, read, seq, t, message(e), read.done); + return classify(e); + } + let outcome: Outcome = "ok"; + read.targets.forEach((target, i) => { + const bytes = parts[i]; + if (bytes === undefined) { + outcome = "timeout"; + this.fail(target.id, read, seq, t, "silent", read.done); + return; + } + this.deliver(target.id, read, bytes, seq, t, read.done); + }); + return outcome; + } + const target = read.targets[0]; + if (target === undefined) return "ok"; + try { + const bytes = await client.read(target.id, read.addr, read.count); + this.deliver(target.id, read, bytes, seq, t, read.done); + return "ok"; + } catch (e) { + this.fail(target.id, read, seq, t, message(e), read.done); + return classify(e); + } + } + + // Fan-out + + private deliver( + id: number, + span: Span, + bytes: Uint8Array, + seq: number, + t: number, + done: ((r: ReadResult) => void) | undefined, + ): void { + const layout = this.layoutOf?.(id); + if (layout === undefined) return; + const values = decodeValues(layout, span, bytes); + const state = this.state(id); + for (const [name, value] of values) state.cache.set(name, { value, t, seq }); + if (state.probing || state.failures > 0) { + state.probing = false; + state.failures = 0; + this.replan(id); + } + const snapshot: Snapshot = { id, seq, t, values, read: readerOver(values), stale: false }; + this.fanOut(id, span, snapshot); + done?.({ values, seq, t }); + } + + private fail( + id: number, + span: Span, + seq: number, + t: number, + error: string, + done: ((r: ReadResult) => void) | undefined, + ): void { + const state = this.state(id); + state.failures++; + if (state.failures >= PROBE_FAILURES && !state.probing) { + state.probing = true; + state.probe.due = this.nowMs() + PROBE_MS; + } + // Every subscriber with a field in the failed span hears, straddlers too; + // what the cache still holds of their own registers rides along. + for (const entry of this.subs) { + if (entry.sub.id !== id || entry.fields === undefined) continue; + if (!entry.fields.some((f) => within(span, f))) continue; + entry.listener(this.fromCache(entry, true, error)); + } + done?.({ values: new Map(), seq, t, error }); + } + + /** + * Every subscriber the read covers gets the same Snapshot object. One whose + * registers straddle two spans gets them assembled from the servo's cache, + * once every one of them has landed. + */ + private fanOut(id: number, span: Span, snapshot: Snapshot): void { + for (const entry of this.subs) { + if (entry.sub.id !== id || entry.fields === undefined) continue; + if (entry.fields.every((f) => within(span, f))) { + entry.listener(snapshot); + continue; + } + if (!entry.fields.some((f) => within(span, f))) continue; + const assembled = this.fromCache(entry, false); + if (assembled.values.size === entry.fields.length) entry.listener(assembled); + } + } + + /** + * The subscription's own registers out of the servo's cache, dated by the + * oldest of them: a snapshot served this way is never fresher than its + * stalest member. + */ + private fromCache(entry: Sub, stale: boolean, error?: string): Snapshot { + const cache = this.servos.get(entry.sub.id)?.cache; + const values = new Map(); + let t = this.clock.now(); + let seq = this.seq; + for (const name of entry.fields?.map((f) => f.name) ?? entry.sub.registers) { + const hit = cache?.get(name); + if (hit === undefined) continue; + values.set(name, hit.value); + t = Math.min(t, hit.t); + seq = Math.min(seq, hit.seq); + } + return { + id: entry.sub.id, + seq, + t, + values, + read: readerOver(values), + stale, + error, + }; + } + + /** A subscriber that mounts mid-period sees the cache at once, not the next read. */ + private serveFresh(): void { + for (const entry of this.fresh.splice(0)) { + if (!this.subs.has(entry) || entry.fields === undefined) continue; + const snapshot = this.fromCache(entry, this.unsure(entry.sub.id)); + if (snapshot.values.size !== entry.fields.length) continue; + entry.listener(snapshot); + } + } + + /** The servo is not being read on its spans right now. */ + private unsure(id: number): boolean { + if (this.client === undefined || !this.scheduled(id)) return true; + const state = this.servos.get(id); + return state !== undefined && (state.probing || state.failures > 0); + } + + // Bookkeeping + + private rejectControl(item: ControlItem, error: unknown): void { + if (item.kind === "write") for (const s of item.settle) s.reject(error); + else item.settle.reject(error); + } + + /** A write acks: re-read the smallest subscribed span covering it, else the field alone. */ + private markDirty(id: number, f: Field): void { + const state = this.servos.get(id); + const covering = state?.spans + .filter((s) => within(s, f)) + .reduce( + (best, s) => (best === undefined || s.count < best.count ? s : best), + undefined, + ); + const span = covering ?? { addr: f.addr, count: f.width }; + const already = this.refresh.some( + (r) => r.id === id && r.addr === span.addr && r.count === span.count && r.done === undefined, + ); + if (already) return; + this.refresh.push({ id, addr: span.addr, count: span.count, queuedAt: this.nowMs() }); + } + + /** After an exclusive turn every span restarts its period, so no burst follows. */ + private restart(): void { + const now = this.nowMs(); + for (const state of this.servos.values()) { + for (const span of [...state.spans, state.probe]) span.due = now + this.period(span.rate); + } + } + + private arm(cancel: (() => void) | undefined): void { + this.stopTimer?.(); + this.stopTimer = cancel; + } + + private armNext(): void { + let earliest: number | undefined; + for (const [id, state] of this.servos) { + if (!this.scheduled(id)) continue; + for (const span of this.spansOf(state)) { + if (earliest === undefined || span.due < earliest) earliest = span.due; + } + } + if (earliest === undefined) { + this.arm(undefined); + return; + } + const wait = Math.max(0, earliest - this.nowMs()); + this.arm( + this.clock.after(wait, () => { + this.stopTimer = undefined; + this.pump(); + }), + ); + } +} diff --git a/src/lib/bus/spans.test.ts b/src/lib/bus/spans.test.ts new file mode 100644 index 0000000..ea57194 --- /dev/null +++ b/src/lib/bus/spans.test.ts @@ -0,0 +1,238 @@ +import type { Field, Value } from "@openservocore/client"; +import { expect, test } from "vitest"; +import descriptor from "../../../../open-servo-core/descriptors/osc-servo/0.1.json"; +import { + CONSTANT_REGISTERS, + constantsFrom, + decodeSpan, + decodeValues, + faultText, + healthFrom, + HEALTH_REGISTERS, + liveFrom, + LIVE_REGISTERS, + plan, + planSpans, + readerOver, + READ_MAX, + span, + within, + type Layout, +} from "./spans"; + +const fields = descriptor.fields as Field[]; + +test("the live span is one read from pos through vmotor_bias_counts", () => { + expect(span(fields, LIVE_REGISTERS)).toEqual({ addr: 576, count: 26 }); +}); + +test("the constants span is one read from raw_min through vmotor_bias_nom_counts", () => { + expect(span(fields, CONSTANT_REGISTERS)).toEqual({ addr: 128, count: 172 }); +}); + +test("plan fails loudly on a descriptor missing a register", () => { + expect(() => plan(fields.filter((f) => f.name !== "ntc_raw"))).toThrow("ntc_raw"); +}); + +const sample: Field[] = [ + { name: "u8", addr: 10, width: 1, access: "ro", kind: "uint", variants: [] }, + { name: "i8", addr: 11, width: 1, access: "ro", kind: "int", variants: [] }, + { name: "u16", addr: 12, width: 2, access: "ro", kind: "uint", variants: [] }, + { name: "i16", addr: 14, width: 2, access: "ro", kind: "int", variants: [] }, + { name: "u32", addr: 16, width: 4, access: "ro", kind: "uint", variants: [] }, + { name: "i32", addr: 20, width: 4, access: "ro", kind: "int", variants: [] }, + { name: "blob", addr: 24, width: 3, access: "ro", kind: "bytes", variants: [] }, + { name: "beyond", addr: 27, width: 1, access: "ro", kind: "uint", variants: [] }, +]; + +test("decodeSpan reads little-endian widths and signs at the span offset", () => { + const bytes = new Uint8Array([ + 0xff, 0xff, 0x34, 0x12, 0xfe, 0xff, 0x78, 0x56, 0x34, 0x12, 0xff, 0xff, 0xff, 0xff, 1, 2, 3, + ]); + const read = decodeSpan(sample, { addr: 10, count: bytes.length }, bytes); + expect(read("u8")).toBe(255); + expect(read("i8")).toBe(-1); + expect(read("u16")).toBe(0x1234); + expect(read("i16")).toBe(-2); + expect(read("u32")).toBe(0x12345678); + expect(read("i32")).toBe(-1); + expect(() => read("blob")).toThrow("not a number"); + expect(() => read("beyond")).toThrow("outside"); +}); + +test("decodeSpan honours a subarray's own offset", () => { + const backing = new Uint8Array([9, 9, 0x34, 0x12]); + const read = decodeSpan(sample, { addr: 12, count: 2 }, backing.subarray(2)); + expect(read("u16")).toBe(0x1234); +}); + +test("spans merge in address order, gaps read through, a new span past READ_MAX", () => { + const wide: Field[] = [ + { name: "a", addr: 0, width: 2, access: "ro", kind: "uint", variants: [] }, + { name: "b", addr: 100, width: 2, access: "ro", kind: "uint", variants: [] }, + { name: "c", addr: READ_MAX, width: 1, access: "ro", kind: "uint", variants: [] }, + { name: "d", addr: READ_MAX + 4, width: 1, access: "ro", kind: "uint", variants: [] }, + ]; + expect(planSpans(wide, ["d", "a", "c", "b"])).toEqual([ + { addr: 0, count: 102 }, + { addr: READ_MAX, count: 5 }, + ]); +}); + +test("a field already inside a served span gets no span of its own", () => { + const served = planSpans(sample, ["u8", "u16"]); + expect(served).toEqual([{ addr: 10, count: 4 }]); + expect(planSpans(sample, ["i8", "i32"], served)).toEqual([{ addr: 20, count: 4 }]); + const inside = sample.filter((f) => served.some((s) => within(s, f))).map((f) => f.name); + expect(inside).toEqual(["u8", "i8", "u16"]); +}); + +test("planSpans rejects a register the descriptor does not carry", () => { + expect(() => planSpans(sample, ["nope"])).toThrow("descriptor has no nope"); +}); + +test("the 0.1 sample span is goal_duty through ntc_raw in one read", () => { + expect( + planSpans(fields, [ + "goal_duty", + "goal_position", + "goal_velocity", + "goal_current", + "mode_active", + "omega_hat_cps", + "duty_applied_q15", + "pos", + "current", + "vmotor_a", + "vmotor_b", + "vbus_raw", + "ntc_raw", + ]), + ).toEqual([{ addr: 390, count: 210 }]); +}); + +/** A descriptor codec over `sample`, plain numbers only. */ +const layout: Layout = { + fields: sample, + decode: (name, bytes) => { + const f = sample.find((f) => f.name === name); + if (f === undefined) throw new Error(name); + if (f.kind === "bytes") return { kind: "bytes", value: bytes.slice() }; + const view = new DataView(bytes.buffer, bytes.byteOffset, bytes.byteLength); + const value = + f.width === 1 + ? f.kind === "int" + ? view.getInt8(0) + : view.getUint8(0) + : f.width === 2 + ? f.kind === "int" + ? view.getInt16(0, true) + : view.getUint16(0, true) + : f.kind === "int" + ? view.getInt32(0, true) + : view.getUint32(0, true); + return { kind: f.kind === "int" ? "int" : "uint", value }; + }, + encode: () => new Uint8Array(), +}; + +test("decodeValues covers every field wholly inside the read, and no more", () => { + const bytes = Uint8Array.from({ length: 6 }, (_, i) => i + 1); + const values = decodeValues(layout, { addr: 10, count: 6 }, bytes); + expect([...values.keys()]).toEqual(["u8", "i8", "u16", "i16"]); + expect(values.get("u16")).toEqual({ kind: "uint", value: 0x0403 } satisfies Value); +}); + +test("decodeValues stops at a short reply instead of reading past it", () => { + const values = decodeValues(layout, { addr: 10, count: 6 }, new Uint8Array([1, 2, 3])); + expect([...values.keys()]).toEqual(["u8", "i8"]); +}); + +test("readerOver throws for a name outside the read and for a bytes field", () => { + const read = readerOver( + new Map([ + ["u8", { kind: "uint", value: 7 }], + ["on", { kind: "bool", value: true }], + ["blob", { kind: "bytes", value: new Uint8Array(1) }], + ]), + ); + expect(read("u8")).toBe(7); + expect(read("on")).toBe(1); + expect(() => read("blob")).toThrow("not a number"); + expect(() => read("missing")).toThrow("outside the read"); +}); + +test("the live, health and constants layouts decode from a whole-table image", () => { + const image = new Uint8Array(descriptor.table_size); + const view = new DataView(image.buffer); + const put = (name: string, value: number) => { + const f = fields.find((f) => f.name === name); + if (f === undefined) throw new Error(name); + if (f.width === 4) view.setUint32(f.addr, value, true); + else if (f.width === 2) { + if (f.kind === "int") view.setInt16(f.addr, value, true); + else view.setUint16(f.addr, value, true); + } else view.setUint8(f.addr, value); + }; + put("raw_min", 118); + put("raw_max", 3990); + put("angle_min_cdeg", -9500); + put("angle_max_cdeg", 9500); + put("gear_ratio_centi", 100); + put("shunt_r_mohm", 10); + put("gain_milli", 32000); + put("vdd_mv", 3300); + put("ntc_beta", 3950); + put("pos", 2054); + put("current", 2100); + put("vbus_raw", 3072); + put("ntc_raw", 1500); + put("current_bias_counts", 2048); + put("vmotor_bias_counts", 700); + put("fault_flags", 0x04); + put("status_flags", 0x01); + put("trim_steps", 3); + put("crc_fail_count", 9); + put("framing_drop_count", 4); + const read = (names: readonly string[]) => { + const s = span(fields, names); + return decodeSpan(fields, s, image.subarray(s.addr, s.addr + s.count)); + }; + const constants = constantsFrom(read(CONSTANT_REGISTERS)); + expect(constants.calibration).toEqual({ + rawMin: 118, + rawMax: 3990, + angleMinCdeg: -9500, + angleMaxCdeg: 9500, + gearRatioCenti: 100, + }); + expect(constants.sense.gainMilli).toBe(32000); + expect(constants.calibrated.valid).toBe(true); + expect(liveFrom(read(LIVE_REGISTERS))).toEqual({ + pos: 2054, + current: 2100, + vbusRaw: 3072, + ntcRaw: 1500, + biases: { currentBiasCounts: 2048, vmotorBiasCounts: 700 }, + }); + expect(healthFrom(read(HEALTH_REGISTERS))).toEqual({ + faultFlags: 0x04, + configDirty: true, + trimSteps: 3, + crcFailCount: 9, + framingDropCount: 4, + }); +}); + +test("the fleet card reads the live sensors and the health block as one span", () => { + expect(planSpans(fields, [...LIVE_REGISTERS, ...HEALTH_REGISTERS])).toEqual([ + { addr: 512, count: 90 }, + ]); +}); + +test("faultText names the latched bits, lowest first", () => { + expect(faultText(0)).toBeUndefined(); + expect(faultText(1 << 2)).toBe("Stall detected"); + expect(faultText((1 << 0) | (1 << 5))).toBe("Overcurrent, Undervoltage"); + expect(faultText(1 << 7)).toBe("Fault 0x80"); +}); diff --git a/src/lib/card-poll.ts b/src/lib/bus/spans.ts similarity index 53% rename from src/lib/card-poll.ts rename to src/lib/bus/spans.ts index 0dafae6..7bdd486 100644 --- a/src/lib/card-poll.ts +++ b/src/lib/bus/spans.ts @@ -1,4 +1,7 @@ -import type { Field, Health, OscClient } from "@openservocore/client"; +// Register names to the reads that carry them, and the bytes back to values. +// Pure: no client, no clock, no React. + +import type { Descriptor, Field, Health, Value } from "@openservocore/client"; import { biasesFromTable, calibrationFromTable, @@ -9,51 +12,37 @@ import { type CalibrationStatus, type ReadRegister, type Sense, -} from "./units"; +} from "../units"; -export const POLL_MS = 1000; +/** A READ payload is at most 252 bytes (protocol sec 5.1). */ +export const READ_MAX = 252; -/** CALIB pot_lut, kinematics, sense and sense_ext: read once per servo. */ -export const CONSTANT_REGISTERS: readonly string[] = [ - "raw_min", - "raw_max", - "angle_min_cdeg", - "angle_max_cdeg", - "gear_ratio_centi", - "shunt_r_mohm", - "gain_milli", - "vdd_mv", - "vmotor_div_top", - "vmotor_div_bot", - "vbus_div_top_ohm", - "vbus_div_bot_ohm", - "ntc_pullup_ohm", - "ntc_r25_ohm", - "ntc_beta", - "vmotor_bias_nom_counts", -]; - -/** TELEMETRY sensors the cards show, plus the biases their conversion needs. */ -export const LIVE_REGISTERS: readonly string[] = [ - "pos", - "current", - "vbus_raw", - "ntc_raw", - "current_bias_counts", - "vmotor_bias_counts", -]; +/** `status_flags` bit 0: modified since the last save (protocol sec 9.4). */ +const STATUS_FLAG_CONFIG_DIRTY = 1 << 0; export interface Span { addr: number; count: number; } -function field(fields: readonly Field[], name: string): Field { +/** The slice of a descriptor the scheduler plans and decodes through. */ +export interface Layout { + encode: Descriptor["encode"]; + decode: Descriptor["decode"]; + fields: readonly Field[]; +} + +export function field(fields: readonly Field[], name: string): Field { const f = fields.find((f) => f.name === name); if (f === undefined) throw new Error(`descriptor has no ${name}`); return f; } +/** The field lies wholly inside the span. */ +export function within(span: Span, f: Field): boolean { + return f.addr >= span.addr && f.addr + f.width <= span.addr + span.count; +} + /** The one contiguous read covering every named register. */ export function span(fields: readonly Field[], names: readonly string[]): Span { let lo = Infinity; @@ -66,6 +55,56 @@ export function span(fields: readonly Field[], names: readonly string[]): Span { return { addr: lo, count: hi - lo }; } +/** + * The fewest reads covering `names`, fields in address order, each joined to + * the span before it unless that would pass READ_MAX. Gaps are read through: + * a byte costs a few microseconds, a second exchange a millisecond or more. + * Fields already inside a `served` span get no span of their own. + */ +export function planSpans( + fields: readonly Field[], + names: readonly string[], + served: readonly Span[] = [], +): Span[] { + const picked = [...new Set(names)] + .map((name) => field(fields, name)) + .filter((f) => !served.some((s) => within(s, f))) + .sort((a, b) => a.addr - b.addr); + const spans: Span[] = []; + for (const f of picked) { + const last = spans.at(-1); + const end = f.addr + f.width; + if (last !== undefined && end - last.addr <= READ_MAX) { + last.count = Math.max(last.count, end - last.addr); + } else { + spans.push({ addr: f.addr, count: f.width }); + } + } + return spans; +} + +/** Every field the read covered, decoded once through the descriptor's codec. */ +export function decodeValues(layout: Layout, span: Span, bytes: Uint8Array): Map { + const end = span.addr + Math.min(span.count, bytes.length); + const values = new Map(); + for (const f of layout.fields) { + if (f.addr < span.addr || f.addr + f.width > end) continue; + const at = f.addr - span.addr; + values.set(f.name, layout.decode(f.name, bytes.subarray(at, at + f.width))); + } + return values; +} + +/** A numeric reader over decoded values; bools read as 0 or 1. */ +export function readerOver(values: ReadonlyMap): ReadRegister { + return (name) => { + const v = values.get(name); + if (v === undefined) throw new Error(`${name} outside the read`); + if (v.kind === "bytes") throw new Error(`${name} is not a number`); + return typeof v.value === "boolean" ? Number(v.value) : v.value; + }; +} + /** A register reader over the bytes one `span` read returned. */ export function decodeSpan(fields: readonly Field[], span: Span, bytes: Uint8Array): ReadRegister { const view = new DataView(bytes.buffer, bytes.byteOffset, bytes.byteLength); @@ -87,6 +126,45 @@ export function decodeSpan(fields: readonly Field[], span: Span, bytes: Uint8Arr }; } +/** CALIB pot_lut, kinematics, sense and sense_ext: read once per servo. */ +export const CONSTANT_REGISTERS: readonly string[] = [ + "raw_min", + "raw_max", + "angle_min_cdeg", + "angle_max_cdeg", + "gear_ratio_centi", + "shunt_r_mohm", + "gain_milli", + "vdd_mv", + "vmotor_div_top", + "vmotor_div_bot", + "vbus_div_top_ohm", + "vbus_div_bot_ohm", + "ntc_pullup_ohm", + "ntc_r25_ohm", + "ntc_beta", + "vmotor_bias_nom_counts", +]; + +/** TELEMETRY sensors the cards show, plus the biases their conversion needs. */ +export const LIVE_REGISTERS: readonly string[] = [ + "pos", + "current", + "vbus_raw", + "ntc_raw", + "current_bias_counts", + "vmotor_bias_counts", +]; + +/** TELEMETRY-COMMON front, the fields `OscClient.health` reads as one block. */ +export const HEALTH_REGISTERS: readonly string[] = [ + "fault_flags", + "status_flags", + "trim_steps", + "crc_fail_count", + "framing_drop_count", +]; + export interface Plan { fields: Field[]; constants: Span; @@ -130,28 +208,22 @@ export function liveFrom(read: ReadRegister): Live { }; } +export function healthFrom(read: ReadRegister): Health { + return { + faultFlags: read("fault_flags"), + configDirty: (read("status_flags") & STATUS_FLAG_CONFIG_DIRTY) !== 0, + trimSteps: read("trim_steps"), + crcFailCount: read("crc_fail_count"), + framingDropCount: read("framing_drop_count"), + }; +} + export interface CardValues { constants: Constants; live: Live; health: Health; } -/** One card's values: the constants read only when not already known. */ -export async function readCard( - client: OscClient, - id: number, - plan: Plan, - constants: Constants | undefined, -): Promise { - if (constants === undefined) { - const bytes = await client.read(id, plan.constants.addr, plan.constants.count); - constants = constantsFrom(decodeSpan(plan.fields, plan.constants, bytes)); - } - const bytes = await client.read(id, plan.live.addr, plan.live.count); - const live = liveFrom(decodeSpan(plan.fields, plan.live, bytes)); - return { constants, live, health: await client.health(id) }; -} - /** Kernel fault latch bits (firmware kernel/faults.rs), lowest bit first. */ const FAULTS: readonly string[] = [ "Overcurrent", diff --git a/src/lib/bus/stats.ts b/src/lib/bus/stats.ts new file mode 100644 index 0000000..1ef6766 --- /dev/null +++ b/src/lib/bus/stats.ts @@ -0,0 +1,171 @@ +// Every exchange is timed at three points; these are the numbers the ?debug +// panel reads and the hardware procedure records. + +export type Lane = "control" | "refresh" | "live" | "exclusive"; +export type Kind = "read" | "gread" | "write" | "command"; +export type Outcome = "ok" | "timeout" | "stalled" | "error"; + +export interface Exchange { + seq: number; + lane: Lane; + kind: Kind; + id: number | undefined; + addr: number | undefined; + bytes: number; + queuedAt: number; + startedAt: number; + settledAt: number; + outcome: Outcome; +} + +export interface Quantiles { + p50: number; + p95: number; +} + +export interface BusStats { + exchanges: number; + perSecond: number; + bytesPerSecond: number; + /** Milliseconds; p50 and p95 over the last 256 exchanges. */ + wait: Quantiles; + run: { small: Quantiles; medium: Quantiles; large: Quantiles }; + /** setTimeout(0) drift sampled twice a second: the event-loop lag probe. */ + lag: Quantiles; + timeouts: number; + stalled: number; + errors: number; + coalesced: number; + grouped: number; + utilisation: number; + effectivePeriodMs: { fast: number; slow: number }; + perServo: ReadonlyMap; + recent: readonly Exchange[]; +} + +export type SizeClass = "small" | "medium" | "large"; + +const SMALL_MAX = 32; +const MEDIUM_MAX = 128; +/** Exchanges kept for the quantiles. */ +const WINDOW = 256; +/** Lag samples kept, about two minutes at the probe's cadence. */ +const LAG_WINDOW = 256; +/** Weight of the newest run in the per-class mean. */ +const EWMA_ALPHA = 0.2; + +export function sizeClass(bytes: number): SizeClass { + if (bytes < SMALL_MAX) return "small"; + if (bytes < MEDIUM_MAX) return "medium"; + return "large"; +} + +const ZERO: Quantiles = { p50: 0, p95: 0 }; + +function quantiles(values: number[]): Quantiles { + if (values.length === 0) return ZERO; + const sorted = [...values].sort((a, b) => a - b); + const at = (q: number) => sorted[Math.min(sorted.length - 1, Math.floor(q * sorted.length))] ?? 0; + return { p50: at(0.5), p95: at(0.95) }; +} + +export interface StatsExtra { + utilisation: number; + effectivePeriodMs: { fast: number; slow: number }; + perServo: ReadonlyMap; +} + +/** The exchange log and the counters over it. */ +export class StatsRecorder { + private readonly ring: Exchange[] = []; + private readonly lags: number[] = []; + private readonly mean: Record = { small: 0, medium: 0, large: 0 }; + private exchanges = 0; + private timeouts = 0; + private stalled = 0; + private errors = 0; + private coalesced = 0; + private grouped = 0; + + record(e: Exchange): void { + this.exchanges++; + this.ring.push(e); + if (this.ring.length > WINDOW) this.ring.shift(); + switch (e.outcome) { + case "timeout": + this.timeouts++; + break; + case "stalled": + this.stalled++; + break; + case "error": + this.errors++; + break; + case "ok": + break; + } + if (e.kind !== "read" && e.kind !== "gread") return; + const c = sizeClass(e.bytes); + const run = e.settledAt - e.startedAt; + this.mean[c] = this.mean[c] === 0 ? run : this.mean[c] * (1 - EWMA_ALPHA) + run * EWMA_ALPHA; + } + + coalesce(): void { + this.coalesced++; + } + + group(): void { + this.grouped++; + } + + lag(ms: number): void { + this.lags.push(ms); + if (this.lags.length > LAG_WINDOW) this.lags.shift(); + } + + /** Milliseconds one read of `bytes` is expected to take. */ + estimate(bytes: number): number { + return this.mean[sizeClass(bytes)]; + } + + build(extra: StatsExtra): BusStats { + const first = this.ring[0]; + const last = this.ring.at(-1); + const seconds = + first === undefined || last === undefined ? 0 : (last.settledAt - first.startedAt) / 1000; + const bytes = this.ring.reduce((sum, e) => sum + e.bytes, 0); + const runs = (c: SizeClass) => + quantiles( + this.ring + .filter((e) => (e.kind === "read" || e.kind === "gread") && sizeClass(e.bytes) === c) + .map((e) => e.settledAt - e.startedAt), + ); + return { + exchanges: this.exchanges, + perSecond: seconds > 0 ? this.ring.length / seconds : 0, + bytesPerSecond: seconds > 0 ? bytes / seconds : 0, + wait: quantiles(this.ring.map((e) => e.startedAt - e.queuedAt)), + run: { small: runs("small"), medium: runs("medium"), large: runs("large") }, + lag: quantiles(this.lags), + timeouts: this.timeouts, + stalled: this.stalled, + errors: this.errors, + coalesced: this.coalesced, + grouped: this.grouped, + utilisation: extra.utilisation, + effectivePeriodMs: extra.effectivePeriodMs, + perServo: extra.perServo, + recent: this.ring.slice(), + }; + } +} + +/** A link or servo failure, told apart by the message the client rejected with. */ +export function classify(error: unknown): Exclude { + const text = (error instanceof Error ? error.message : String(error)).toLowerCase(); + if (text.includes("guard") || text.includes("stall")) return "stalled"; + if (text.includes("timeout") || text.includes("timed out") || text.includes("silent")) { + return "timeout"; + } + return "error"; +} diff --git a/src/lib/bus/store.test.ts b/src/lib/bus/store.test.ts new file mode 100644 index 0000000..ff36da1 --- /dev/null +++ b/src/lib/bus/store.test.ts @@ -0,0 +1,83 @@ +import { expect, test } from "vitest"; +import type { Snapshot } from "./manager"; +import { BusStore } from "./store"; + +function frames(): { store: BusStore; tick: () => void } { + const queue: (() => void)[] = []; + const store = new BusStore((fn) => queue.push(fn)); + return { + store, + tick: () => { + const pending = queue.splice(0); + for (const fn of pending) fn(); + }, + }; +} + +function snapshot(t: number): Snapshot { + const values = new Map([["pos", { kind: "uint" as const, value: t }]]); + return { id: 1, seq: t, t, values, read: () => t, stale: false }; +} + +test("a cell notifies once per frame however many reads land", () => { + const { store, tick } = frames(); + const cell = store.cell(0); + let notified = 0; + cell.subscribe(() => notified++); + for (const v of [1, 2, 3, 4]) cell.set(v); + expect(cell.get()).toBe(0); + expect(notified).toBe(0); + tick(); + expect(cell.get()).toBe(4); + expect(notified).toBe(1); +}); + +test("only changed cells notify", () => { + const { store, tick } = frames(); + const moved = store.cell("a"); + const still = store.cell("x"); + let movedSeen = 0; + let stillSeen = 0; + moved.subscribe(() => movedSeen++); + still.subscribe(() => stillSeen++); + moved.set("b"); + still.set("x"); + tick(); + expect(movedSeen).toBe(1); + expect(stillSeen).toBe(0); + expect(still.get()).toBe("x"); +}); + +test("one frame is scheduled for the whole store", () => { + let scheduled = 0; + const queue: (() => void)[] = []; + const store = new BusStore((fn) => { + scheduled++; + queue.push(fn); + }); + const a = store.cell(0); + const b = store.cell(0); + a.set(1); + b.set(1); + a.set(2); + expect(scheduled).toBe(1); + for (const fn of queue.splice(0)) fn(); + a.set(3); + expect(scheduled).toBe(2); +}); + +test("the ring publishes a fresh array per frame and trims to the window", () => { + const { store, tick } = frames(); + const ring = store.ring(30); + let seen = 0; + ring.subscribe(() => seen++); + for (const t of [0, 10, 20]) ring.push(snapshot(t)); + tick(); + expect(seen).toBe(1); + const first = ring.get(); + expect(first.map((s) => s.t)).toEqual([0, 10, 20]); + ring.push(snapshot(31)); + tick(); + expect(ring.get().map((s) => s.t)).toEqual([10, 20, 31]); + expect(ring.get()).not.toBe(first); +}); diff --git a/src/lib/bus/store.ts b/src/lib/bus/store.ts new file mode 100644 index 0000000..f3fe3c4 --- /dev/null +++ b/src/lib/bus/store.ts @@ -0,0 +1,124 @@ +// The React binding for one manager. Reads settle whenever the bus frees; +// this collects them and publishes once per animation frame, so however many +// land inside a frame React renders once and only where a value moved. + +import type { Snapshot } from "./manager"; + +export type Frame = (fn: () => void) => void; + +interface Publishable { + publish: () => void; +} + +export class Cell { + private value: T; + private staged: T; + private readonly listeners = new Set<() => void>(); + + constructor( + private readonly store: BusStore, + initial: T, + ) { + this.value = initial; + this.staged = initial; + } + + readonly get = (): T => this.value; + + readonly subscribe = (listener: () => void): (() => void) => { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + }; + + set(next: T): void { + this.staged = next; + this.store.stage(this); + } + + publish = (): void => { + if (Object.is(this.value, this.staged)) return; + this.value = this.staged; + for (const listener of this.listeners) listener(); + }; +} + +/** The last `windowS` seconds of snapshots, oldest first, one array per frame. */ +export class RingCell { + private buf: Snapshot[] = []; + private value: readonly Snapshot[] = []; + private readonly listeners = new Set<() => void>(); + + constructor( + private readonly store: BusStore, + private readonly windowS: number, + ) {} + + readonly get = (): readonly Snapshot[] => this.value; + + readonly subscribe = (listener: () => void): (() => void) => { + this.listeners.add(listener); + return () => { + this.listeners.delete(listener); + }; + }; + + push(snapshot: Snapshot): void { + this.buf.push(snapshot); + const cutoff = snapshot.t - this.windowS; + let drop = 0; + for (const s of this.buf) { + if (s.t >= cutoff) break; + drop++; + } + if (drop > 0) this.buf.splice(0, drop); + this.store.stage(this); + } + + publish = (): void => { + this.value = this.buf.slice(); + for (const listener of this.listeners) listener(); + }; +} + +export class BusStore { + private readonly pending = new Set(); + private scheduled = false; + + constructor(private readonly frame: Frame) {} + + cell(initial: T): Cell { + return new Cell(this, initial); + } + + ring(windowS: number): RingCell { + return new RingCell(this, windowS); + } + + /** @internal */ + stage(cell: Publishable): void { + this.pending.add(cell); + if (this.scheduled) return; + this.scheduled = true; + this.frame(() => { + this.scheduled = false; + this.flush(); + }); + } + + private flush(): void { + const cells = [...this.pending]; + this.pending.clear(); + for (const cell of cells) cell.publish(); + } +} + +export const animationFrame: Frame = + typeof requestAnimationFrame === "function" + ? (fn) => { + requestAnimationFrame(fn); + } + : (fn) => { + setTimeout(fn, 16); + }; diff --git a/src/lib/card-poll.test.ts b/src/lib/card-poll.test.ts deleted file mode 100644 index ecfaed4..0000000 --- a/src/lib/card-poll.test.ts +++ /dev/null @@ -1,121 +0,0 @@ -import type { Field } from "@openservocore/client"; -import { expect, test } from "vitest"; -import descriptor from "../../../open-servo-core/descriptors/osc-servo/0.1.json"; -import { - CONSTANT_REGISTERS, - constantsFrom, - decodeSpan, - faultText, - LIVE_REGISTERS, - liveFrom, - plan, - span, -} from "./card-poll"; - -const fields = descriptor.fields as Field[]; - -test("the live span is one read from pos through vmotor_bias_counts", () => { - expect(span(fields, LIVE_REGISTERS)).toEqual({ addr: 576, count: 26 }); -}); - -test("the constants span is one read from raw_min through vmotor_bias_nom_counts", () => { - expect(span(fields, CONSTANT_REGISTERS)).toEqual({ addr: 128, count: 172 }); -}); - -test("plan fails loudly on a descriptor missing a register", () => { - expect(() => plan(fields.filter((f) => f.name !== "ntc_raw"))).toThrow("ntc_raw"); -}); - -const sample: Field[] = [ - { name: "u8", addr: 10, width: 1, access: "ro", kind: "uint", variants: [] }, - { name: "i8", addr: 11, width: 1, access: "ro", kind: "int", variants: [] }, - { name: "u16", addr: 12, width: 2, access: "ro", kind: "uint", variants: [] }, - { name: "i16", addr: 14, width: 2, access: "ro", kind: "int", variants: [] }, - { name: "u32", addr: 16, width: 4, access: "ro", kind: "uint", variants: [] }, - { name: "i32", addr: 20, width: 4, access: "ro", kind: "int", variants: [] }, - { name: "blob", addr: 24, width: 3, access: "ro", kind: "bytes", variants: [] }, - { name: "beyond", addr: 27, width: 1, access: "ro", kind: "uint", variants: [] }, -]; - -test("decodeSpan reads little-endian widths and signs at the span offset", () => { - const bytes = new Uint8Array([ - 0xff, 0xff, 0x34, 0x12, 0xfe, 0xff, 0x78, 0x56, 0x34, 0x12, 0xff, 0xff, 0xff, 0xff, 1, 2, 3, - ]); - const read = decodeSpan(sample, { addr: 10, count: bytes.length }, bytes); - expect(read("u8")).toBe(255); - expect(read("i8")).toBe(-1); - expect(read("u16")).toBe(0x1234); - expect(read("i16")).toBe(-2); - expect(read("u32")).toBe(0x12345678); - expect(read("i32")).toBe(-1); - expect(() => read("blob")).toThrow("not a number"); - expect(() => read("beyond")).toThrow("outside"); -}); - -test("decodeSpan honours a subarray's own offset", () => { - const backing = new Uint8Array([9, 9, 0x34, 0x12]); - const read = decodeSpan(sample, { addr: 12, count: 2 }, backing.subarray(2)); - expect(read("u16")).toBe(0x1234); -}); - -test("the live and constants layouts decode from a whole-table image", () => { - const image = new Uint8Array(descriptor.table_size); - const view = new DataView(image.buffer); - const put = (name: string, value: number) => { - const f = fields.find((f) => f.name === name); - if (f === undefined) throw new Error(name); - if (f.width === 2) { - if (f.kind === "int") view.setInt16(f.addr, value, true); - else view.setUint16(f.addr, value, true); - } else view.setUint8(f.addr, value); - }; - put("raw_min", 118); - put("raw_max", 3990); - put("angle_min_cdeg", -9500); - put("angle_max_cdeg", 9500); - put("gear_ratio_centi", 100); - put("shunt_r_mohm", 10); - put("gain_milli", 32000); - put("vdd_mv", 3300); - put("ntc_beta", 3950); - put("pos", 2054); - put("current", 2100); - put("vbus_raw", 3072); - put("ntc_raw", 1500); - put("current_bias_counts", 2048); - put("vmotor_bias_counts", 700); - const p = plan(fields); - const constants = constantsFrom( - decodeSpan( - fields, - p.constants, - image.subarray(p.constants.addr, p.constants.addr + p.constants.count), - ), - ); - expect(constants.calibration).toEqual({ - rawMin: 118, - rawMax: 3990, - angleMinCdeg: -9500, - angleMaxCdeg: 9500, - gearRatioCenti: 100, - }); - expect(constants.sense.gainMilli).toBe(32000); - expect(constants.calibrated.valid).toBe(true); - const live = liveFrom( - decodeSpan(fields, p.live, image.subarray(p.live.addr, p.live.addr + p.live.count)), - ); - expect(live).toEqual({ - pos: 2054, - current: 2100, - vbusRaw: 3072, - ntcRaw: 1500, - biases: { currentBiasCounts: 2048, vmotorBiasCounts: 700 }, - }); -}); - -test("faultText names the latched bits, lowest first", () => { - expect(faultText(0)).toBeUndefined(); - expect(faultText(1 << 2)).toBe("Stall detected"); - expect(faultText((1 << 0) | (1 << 5))).toBe("Overcurrent, Undervoltage"); - expect(faultText(1 << 7)).toBe("Fault 0x80"); -}); diff --git a/src/lib/command-queue.test.ts b/src/lib/command-queue.test.ts deleted file mode 100644 index 2534ed6..0000000 --- a/src/lib/command-queue.test.ts +++ /dev/null @@ -1,74 +0,0 @@ -import { expect, test } from "vitest"; -import { CommandQueue } from "./command-queue"; - -interface Fake { - log: string[]; - busy: boolean; - command(name: string, gate?: Promise): Promise; -} - -function fake(): Fake { - return { - log: [], - busy: false, - async command(name, gate) { - if (this.busy) throw new Error("busy"); - this.busy = true; - this.log.push(`start ${name}`); - await gate; - this.log.push(`end ${name}`); - this.busy = false; - return name; - }, - }; -} - -function deferred(): { promise: Promise; resolve: () => void } { - let resolve: () => void = () => undefined; - const promise = new Promise((r) => { - resolve = r; - }); - return { promise, resolve }; -} - -test("jobs run one at a time in submission order", async () => { - const client = fake(); - const queue = new CommandQueue(() => client); - const gate = deferred(); - const a = queue.run((c) => c.command("a", gate.promise)); - const b = queue.run((c) => c.command("b")); - const c = queue.run((c) => c.command("c")); - await Promise.resolve(); - expect(client.log).toEqual(["start a"]); - gate.resolve(); - expect(await Promise.all([a, b, c])).toEqual(["a", "b", "c"]); - expect(client.log).toEqual(["start a", "end a", "start b", "end b", "start c", "end c"]); -}); - -test("a rejected job does not block the next", async () => { - const client = fake(); - const queue = new CommandQueue(() => client); - const failed = queue.run(() => Promise.reject(new Error("range"))); - const next = queue.run((c) => c.command("next")); - await expect(failed).rejects.toThrow("range"); - expect(await next).toBe("next"); -}); - -test("a job sees the client open when its turn comes, not when it was queued", async () => { - let client: Fake | undefined = fake(); - const queue = new CommandQueue(() => client); - const gate = deferred(); - const first = queue.run((c) => c.command("first", gate.promise)); - const second = queue.run((c) => c.command("second")); - await Promise.resolve(); - expect(client.log).toEqual(["start first"]); - client = undefined; - gate.resolve(); - expect(await first).toBe("first"); - await expect(second).rejects.toThrow("not connected"); -}); - -test("run rejects while nothing is open", async () => { - const queue = new CommandQueue(() => undefined); - await expect(queue.run((c) => c.command("x"))).rejects.toThrow("not connected"); -}); diff --git a/src/lib/command-queue.ts b/src/lib/command-queue.ts deleted file mode 100644 index 1946283..0000000 --- a/src/lib/command-queue.ts +++ /dev/null @@ -1,21 +0,0 @@ -/** - * One promise chain for every adapter command. The adapter takes one command - * at a time and the wasm client rejects an overlapping call as "busy", so - * every caller queues here and each job sees the state the previous one left. - */ -export class CommandQueue { - private chain: Promise = Promise.resolve(); - - constructor(private readonly client: () => C | undefined) {} - - /** Runs `fn` after every queued job; rejects if no client is open when its turn comes. */ - run(fn: (client: C) => Promise): Promise { - const job = this.chain.then(() => { - const client = this.client(); - if (client === undefined) throw new Error("not connected"); - return fn(client); - }); - this.chain = job.catch(() => undefined); - return job; - } -} diff --git a/src/lib/session.tsx b/src/lib/session.tsx index 76cad8a..81fa7bb 100644 --- a/src/lib/session.tsx +++ b/src/lib/session.tsx @@ -17,8 +17,20 @@ import { type ReactNode, } from "react"; import { openClient, simRequested } from "./backend"; -import { faultText, plan, POLL_MS, readCard, type CardValues, type Plan } from "./card-poll"; -import { CommandQueue } from "./command-queue"; +import { BusContext, type BusHost } from "./bus/hooks"; +import { BusManager, systemClock, type Snapshot } from "./bus/manager"; +import { + CONSTANT_REGISTERS, + constantsFrom, + faultText, + healthFrom, + HEALTH_REGISTERS, + liveFrom, + LIVE_REGISTERS, + type CardValues, + type Layout, +} from "./bus/spans"; +import { animationFrame, BusStore } from "./bus/store"; import { fetchDescriptor } from "./descriptor"; import { idle, @@ -46,7 +58,7 @@ export interface Session { descriptor: Descriptor | undefined; descriptorError: string | undefined; descriptorFor: (servo: Servo) => Descriptor | undefined; - /** Per uid, the values the cards poll about once a second. */ + /** Per uid, the values the cards subscribe to about once a second. */ values: ReadonlyMap; connect: () => Promise; disconnect: () => Promise; @@ -55,37 +67,37 @@ export interface Session { setBaud: (rate: BaudRate) => Promise; select: (id: number | undefined) => void; /** - * The only way to talk to the adapter from a page: every command (scan, - * rails, this and other pages' reads) waits its turn on one queue, so - * nothing overlaps and the client's "busy" never surfaces. + * The bus manager's control lane under the name the pages still use: a + * command waits for at most the exchange in flight, and the client's + * "busy" can never surface. */ run: (fn: (client: OscClient) => Promise) => Promise; - /** Re-reads a servo's CALIB constants on the next poll, after a calibration write. */ + /** Re-reads a servo's CALIB constants, after a calibration write. */ refreshConstants: (uid: string) => void; } const SessionContext = createContext(undefined); -/** One descriptor and the card reads planned over it, shared by every servo of that model and firmware. */ -interface Layout { +/** One descriptor, shared by every servo of that model and firmware. */ +interface Model { descriptor: Descriptor; - plan: Plan; + layout: Layout; } -interface Snapshot { +interface Snap { state: SessionState; client: OscClient | undefined; simulated: boolean; - layouts: ReadonlyMap; + models: ReadonlyMap; layoutErrors: ReadonlyMap; values: ReadonlyMap; } -const initial: Snapshot = { +const initial: Snap = { state: idle, client: undefined, simulated: false, - layouts: new Map(), + models: new Map(), layoutErrors: new Map(), values: new Map(), }; @@ -95,11 +107,19 @@ function message(e: unknown): string { } /** Descriptors are published per model and major.minor. */ -function layoutKey(ping: Ping): string { +function modelKey(ping: Ping): string { const [major, minor] = unpackVersion(ping.fw); return `${ping.model}/${major}.${minor}`; } +function busLayout(descriptor: Descriptor): Layout { + return { + encode: descriptor.encode.bind(descriptor), + decode: descriptor.decode.bind(descriptor), + fields: descriptor.fields(), + }; +} + // The adapter takes one command at a time, so the pings run in sequence. async function pingAll(client: OscClient, found: Found[]): Promise { const count = new Map(); @@ -119,12 +139,14 @@ async function release(client: OscClient): Promise { } } -/** Owns the wasm client and its commands; `reduce` owns every state change. */ +/** Owns the wasm client through the bus manager; `reduce` owns every state change. */ class Controller { private snap = initial; - private readonly queue = new CommandQueue(() => this.snap.client); - private pollGen = 0; - private readonly stale = new Set(); + private readonly bus = new BusManager(systemClock); + readonly host: BusHost = { manager: this.bus, store: new BusStore(animationFrame) }; + private readonly cards = new Map>(); + private readonly subscriptions = new Map void>(); + private readonly reading = new Set(); private readonly fetching = new Set(); private readonly listeners = new Set<() => void>(); @@ -135,9 +157,9 @@ class Controller { }; }; - readonly snapshot = (): Snapshot => this.snap; + readonly snapshot = (): Snap => this.snap; - private set(patch: Partial): void { + private set(patch: Partial): void { this.snap = { ...this.snap, ...patch }; for (const listener of this.listeners) listener(); } @@ -147,39 +169,52 @@ class Controller { } // Stable so a card's effect can depend on it. - readonly run = (fn: (client: OscClient) => Promise): Promise => this.queue.run(fn); + readonly run = (fn: (client: OscClient) => Promise): Promise => this.bus.command(fn); refreshConstants(uid: string): void { - this.stale.add(uid); + const servo = this.snap.state.servos.find((s) => s.uid === uid); + if (servo === undefined) return; + const card = this.cards.get(uid); + if (card !== undefined) card.constants = undefined; + this.readConstants(servo); } - layoutFor(servo: Servo): Layout | undefined { - return servo.ping === undefined ? undefined : this.snap.layouts.get(layoutKey(servo.ping)); + modelFor(servo: Servo): Model | undefined { + return servo.ping === undefined ? undefined : this.snap.models.get(modelKey(servo.ping)); } layoutErrorFor(servo: Servo): string | undefined { - return servo.ping === undefined ? undefined : this.snap.layoutErrors.get(layoutKey(servo.ping)); + return servo.ping === undefined ? undefined : this.snap.layoutErrors.get(modelKey(servo.ping)); } - private loadLayouts(servos: Servo[]): void { + private readonly layoutById = (id: number): Layout | undefined => { + const servo = this.snap.state.servos.find((s) => s.id === id); + return servo === undefined ? undefined : this.modelFor(servo)?.layout; + }; + + private loadModels(servos: Servo[]): void { for (const servo of servos) { const { ping } = servo; if (ping === undefined) continue; - const key = layoutKey(ping); - if (this.snap.layouts.has(key) || this.fetching.has(key)) continue; + const key = modelKey(ping); + if (this.snap.models.has(key) || this.fetching.has(key)) continue; this.fetching.add(key); fetchDescriptor(ping.model, ping.fw) .then( (descriptor) => { - if (this.snap.layouts.has(key)) { + if (this.snap.models.has(key)) { descriptor.free(); return; } - const layouts = new Map(this.snap.layouts); - layouts.set(key, { descriptor, plan: plan(descriptor.fields()) }); + const models = new Map(this.snap.models); + models.set(key, { descriptor, layout: busLayout(descriptor) }); const layoutErrors = new Map(this.snap.layoutErrors); layoutErrors.delete(key); - this.set({ layouts, layoutErrors }); + this.set({ models, layoutErrors }); + for (const s of this.snap.state.servos) { + if (s.ping !== undefined && modelKey(s.ping) === key) this.bus.layoutChanged(s.id); + } + this.syncCards(); }, (e: unknown) => { const layoutErrors = new Map(this.snap.layoutErrors); @@ -205,76 +240,131 @@ class Controller { return; } this.set({ client, simulated: simRequested() }); + this.bus.attach(client, this.layoutById); await this.scan(client); } /** With `rate`, migrates the fleet first (servos, then the host, then the reunion sweep). */ private async scan(client: OscClient, rate?: BaudRate): Promise { this.dispatch({ type: "scan" }); - try { - await this.run(async () => { - if (rate !== undefined) { - const ids = [...new Set(this.snap.state.servos.map((s) => s.id))]; - const roster = await client.setBaud(ids, rate); - this.dispatch({ type: "migrated", roster }); + this.clearCards(); + // The whole scan is one exclusive turn: the lanes stay frozen, so on a + // failure the client is released with nothing else reaching for it. + const failure = await this.bus + .exclusive(async (): Promise => { + try { + if (rate !== undefined) { + const ids = [...new Set(this.snap.state.servos.map((s) => s.id))]; + const roster = await client.setBaud(ids, rate); + this.dispatch({ type: "migrated", roster }); + } + const baud = await client.findBusBaud(); + const servos = await pingAll(client, await client.discover()); + const rails = await client.rails(); + if (this.snap.client !== client) return undefined; + this.dispatch({ type: "found", servos, baud, rails }); + this.bus.roster([...new Set(servos.map((s) => s.id))]); + this.loadModels(servos); + return undefined; + } catch (e) { + this.bus.detach(message(e)); + return e; } - const baud = await client.findBusBaud(); - const servos = await pingAll(client, await client.discover()); - const rails = await client.rails(); - if (this.snap.client !== client) return; - this.set({ values: new Map() }); - this.dispatch({ type: "found", servos, baud, rails }); - this.loadLayouts(servos); - }); - if (this.snap.client === client) this.startPoll(client); - } catch (e) { - if (this.snap.client !== client) return; - this.set({ client: undefined, simulated: false, values: new Map() }); - this.dispatch({ type: "fail", error: message(e) }); - await release(client); + }) + .catch((e: unknown) => e); + if (failure === undefined) { + this.syncCards(); + return; } + if (this.snap.client !== client) return; + this.set({ client: undefined, simulated: false }); + this.clearCards(); + this.dispatch({ type: "fail", error: message(failure) }); + await release(client); } - // Fixed cadence; a tick still queued or running when the next is due is - // skipped, so a slow bus never piles reads up behind a scan. - private startPoll(client: OscClient): void { - const gen = ++this.pollGen; - const live = () => - gen === this.pollGen && this.snap.client === client && this.snap.state.status === "ready"; - let ticking = false; - const tick = () => { - if (!live()) { - clearInterval(timer); - return; + // The fleet cards: one slow subscription per servo over the live sensors and + // the health block, plus the CALIB constants read once. + private syncCards(): void { + const wanted = new Set(); + if (this.snap.state.status === "ready") { + for (const servo of this.snap.state.servos) { + if (this.modelFor(servo) === undefined) continue; + wanted.add(servo.uid); + if (this.subscriptions.has(servo.uid)) continue; + this.cards.set(servo.uid, {}); + this.subscriptions.set( + servo.uid, + this.bus.subscribe( + { id: servo.id, registers: [...LIVE_REGISTERS, ...HEALTH_REGISTERS], rate: "slow" }, + (snapshot) => { + this.onCard(servo.uid, snapshot); + }, + ), + ); + this.readConstants(servo); } - if (ticking) return; - ticking = true; - this.run((c) => this.readCards(c, live)) - .catch(() => undefined) - .finally(() => { - ticking = false; - }); - }; - const timer = setInterval(tick, POLL_MS); - tick(); + } + for (const [uid, stop] of this.subscriptions) { + if (wanted.has(uid)) continue; + stop(); + this.subscriptions.delete(uid); + this.cards.delete(uid); + } + this.publishCards(); } - private async readCards(client: OscClient, live: () => boolean): Promise { - for (const servo of this.snap.state.servos) { - if (!live()) return; - const layout = this.layoutFor(servo); - if (layout === undefined) continue; - const values = new Map(this.snap.values); - try { - const prior = this.stale.delete(servo.uid) - ? undefined - : this.snap.values.get(servo.uid)?.constants; - values.set(servo.uid, await readCard(client, servo.id, layout.plan, prior)); - } catch { - values.delete(servo.uid); + private clearCards(): void { + for (const stop of this.subscriptions.values()) stop(); + this.subscriptions.clear(); + this.cards.clear(); + this.set({ values: new Map() }); + } + + private readConstants(servo: Servo): void { + if (this.reading.has(servo.uid)) return; + this.reading.add(servo.uid); + void this.bus + .readOnce(servo.id, CONSTANT_REGISTERS) + .then( + (snapshot) => { + const card = this.cards.get(servo.uid); + if (card === undefined) return; + card.constants = constantsFrom(snapshot.read); + this.publishCards(); + }, + () => undefined, + ) + .finally(() => this.reading.delete(servo.uid)); + } + + private onCard(uid: string, snapshot: Snapshot): void { + const card = this.cards.get(uid); + if (card === undefined) return; + if (snapshot.stale) { + card.live = undefined; + card.health = undefined; + } else { + card.live = liveFrom(snapshot.read); + card.health = healthFrom(snapshot.read); + // A constants read that failed, or one a calibration write invalidated, + // is retried on the next card snapshot. + if (card.constants === undefined) { + const servo = this.snap.state.servos.find((s) => s.uid === uid); + if (servo !== undefined) this.readConstants(servo); } - if (live()) this.set({ values }); } + this.publishCards(); + } + + private publishCards(): void { + const values = new Map(); + for (const [uid, card] of this.cards) { + const { constants, live, health } = card; + if (constants === undefined || live === undefined || health === undefined) continue; + values.set(uid, { constants, live, health }); + } + this.set({ values }); } async discover(): Promise { @@ -293,7 +383,7 @@ class Controller { const { client } = this.snap; if (client === undefined || this.snap.state.status !== "ready") return; // Each toggle merges over the state the previous one acked. - const rails = await this.run((c) => { + const rails = await this.bus.command((c) => { const cur = this.snap.state.rails ?? { v3v3: false, v5: false }; return c.setRails(patch.v3v3 ?? cur.v3v3, patch.v5 ?? cur.v5); }); @@ -303,12 +393,17 @@ class Controller { async disconnect(): Promise { const { client } = this.snap; if (client === undefined) return; - // Queued behind any command still on the old client; `run` would already - // see no client, so the release rides the chain directly. - const done = this.run(() => Promise.resolve()).catch(() => undefined); - this.set({ client: undefined, simulated: false, values: new Map() }); + this.set({ client: undefined, simulated: false }); + this.clearCards(); this.dispatch({ type: "disconnect" }); - await done; + // Detaching from inside an exclusive turn: the exchange in flight has + // settled and the lanes are frozen, so the client is free to release. + await this.bus + .exclusive(() => { + this.bus.detach("disconnected"); + return Promise.resolve(); + }) + .catch(() => undefined); await release(client); } @@ -316,7 +411,7 @@ class Controller { if (this.snap.state.status !== "ready") return; this.dispatch({ type: "select", id }); const servo = this.snap.state.servos.find((s) => s.id === id); - if (servo !== undefined) this.loadLayouts([servo]); + if (servo !== undefined) this.loadModels([servo]); } } @@ -350,9 +445,9 @@ export function SessionProvider({ children }: { children: ReactNode }) { servos: snap.state.servos.map((s) => withHealth(s, snap.values.get(s.uid))), selected: snap.state.selected, missing: snap.state.missing, - descriptor: selected === undefined ? undefined : ctl.layoutFor(selected)?.descriptor, + descriptor: selected === undefined ? undefined : ctl.modelFor(selected)?.descriptor, descriptorError: selected === undefined ? undefined : ctl.layoutErrorFor(selected), - descriptorFor: (servo) => ctl.layoutFor(servo)?.descriptor, + descriptorFor: (servo) => ctl.modelFor(servo)?.descriptor, values: snap.values, connect: () => ctl.connect(), disconnect: async () => { @@ -369,7 +464,11 @@ export function SessionProvider({ children }: { children: ReactNode }) { ctl.refreshConstants(uid); }, }; - return {children}; + return ( + + {children} + + ); } export function useSession(): Session { diff --git a/src/routes/__root.tsx b/src/routes/__root.tsx index d3ebd63..b7aa2e1 100644 --- a/src/routes/__root.tsx +++ b/src/routes/__root.tsx @@ -1,6 +1,7 @@ import { createRootRoute, HeadContent, Outlet, Scripts } from "@tanstack/react-router"; import { useEffect, useSyncExternalStore, type ReactNode } from "react"; import { AppSidebar } from "@/components/app-sidebar"; +import { BusDebug } from "@/components/bus-debug"; import { UnsupportedBrowser } from "@/components/unsupported-browser"; import { SidebarInset, SidebarProvider, SidebarTrigger } from "@/components/ui/sidebar"; import { TooltipProvider } from "@/components/ui/tooltip"; @@ -41,6 +42,7 @@ function RootComponent() {
Use a larger window
+ ) : ( diff --git a/src/routes/index.tsx b/src/routes/index.tsx index 7c2b724..fd2ee35 100644 --- a/src/routes/index.tsx +++ b/src/routes/index.tsx @@ -13,7 +13,7 @@ import { CardTitle, } from "@/components/ui/card"; import { Skeleton } from "@/components/ui/skeleton"; -import { faultText, type CardValues } from "@/lib/card-poll"; +import { faultText, type CardValues } from "@/lib/bus/spans"; import { formatQuantity, formatVersion, hex16 } from "@/lib/format"; import { useSession, type Servo } from "@/lib/session"; import { busV, currentMa, DISPLAY, positionDeg, temperatureC } from "@/lib/units"; diff --git a/tests/e2e/bus.spec.ts b/tests/e2e/bus.spec.ts new file mode 100644 index 0000000..0621a56 --- /dev/null +++ b/tests/e2e/bus.spec.ts @@ -0,0 +1,28 @@ +import { expect, test } from "@playwright/test"; + +const READOUT = /^-?\d+(\.\d+)? deg$/; + +test("the debug panel reports a clean bus while the dashboard polls", async ({ page }) => { + await page.goto("/?sim=1,2&debug"); + await expect(page.getByRole("button", { name: /^Connected/ })).toBeVisible(); + const panel = page.getByRole("region", { name: "Bus statistics" }); + await expect(panel).toBeVisible(); + for (const id of [1, 2]) { + const card = page.getByRole("link", { name: new RegExp(`^ID ${id}\\b`) }); + await expect(card.getByText(READOUT)).toBeVisible(); + } + await page.waitForTimeout(2000); + await expect(panel.locator('dt:text-is("stalled") + dd')).toHaveText("0"); + await expect(panel.locator('dt:text-is("errors") + dd')).toHaveText("0"); + await expect(panel.locator('dt:text-is("exchanges") + dd')).not.toHaveText("0"); + for (const id of [1, 2]) { + const card = page.getByRole("link", { name: new RegExp(`^ID ${id}\\b`) }); + await expect(card.getByText(READOUT)).toBeVisible(); + } +}); + +test("without ?debug the panel stays out of the way", async ({ page }) => { + await page.goto("/?sim=1"); + await expect(page.getByRole("button", { name: /^Connected/ })).toBeVisible(); + await expect(page.getByRole("region", { name: "Bus statistics" })).toHaveCount(0); +});