From fa36ee1058c2727aa7688ada049df77028ebd114 Mon Sep 17 00:00:00 2001 From: Aaron Qian Date: Tue, 22 Sep 2026 20:41:50 -0700 Subject: [PATCH] app: live page subscribes through the bus manager --- src/components/stream-tab.tsx | 62 +++----- src/lib/control.test.ts | 124 ++------------- src/lib/control.ts | 62 +------- src/lib/stream.test.ts | 2 +- src/lib/stream.ts | 2 +- src/lib/telemetry-poll.test.ts | 280 --------------------------------- src/lib/telemetry-poll.ts | 246 ----------------------------- src/lib/telemetry.test.ts | 86 ++++++++++ src/lib/telemetry.ts | 104 ++++++++++++ src/routes/live.tsx | 240 ++++++++++------------------ tests/e2e/control.spec.ts | 32 +++- tests/e2e/helpers.ts | 8 +- 12 files changed, 345 insertions(+), 903 deletions(-) delete mode 100644 src/lib/telemetry-poll.test.ts delete mode 100644 src/lib/telemetry-poll.ts create mode 100644 src/lib/telemetry.test.ts create mode 100644 src/lib/telemetry.ts diff --git a/src/components/stream-tab.tsx b/src/components/stream-tab.tsx index 180a78e..ae51b3f 100644 --- a/src/components/stream-tab.tsx +++ b/src/components/stream-tab.tsx @@ -1,5 +1,5 @@ import { CircleQuestionMark, Download, Zap } from "lucide-react"; -import { useEffect, useId, useMemo, useState } from "react"; +import { useId, useMemo, useState } from "react"; import type uPlot from "uplot"; import { Chart, type ChartOptions } from "@/components/uplot"; import { Button } from "@/components/ui/button"; @@ -15,6 +15,7 @@ import { Checkbox } from "@/components/ui/checkbox"; import { Input } from "@/components/ui/input"; import { Label } from "@/components/ui/label"; import { Popover, PopoverContent, PopoverTrigger } from "@/components/ui/popover"; +import { useBus, useReadOnce } from "@/lib/bus/hooks"; import { useChartTokens, type ChartTokens } from "@/lib/chart-theme"; import { hex16 } from "@/lib/format"; import { useSession } from "@/lib/session"; @@ -37,16 +38,10 @@ import { import { BIAS_REGISTERS, CONFIG_REGISTERS, - decodeSpan, - spanOver, + configFrom, type TelemetryConfig, -} from "@/lib/telemetry-poll"; -import { - biasesFromTable, - calibrationFromTable, - calibrationStatus, - senseFromTable, -} from "@/lib/units"; +} from "@/lib/telemetry"; +import { calibrationStatus } from "@/lib/units"; import { useUnitsPref } from "@/lib/use-pref"; const DEFAULT_COUNT = 120; @@ -57,6 +52,8 @@ const COUNT_MAX = 65535; const WINDOW_MAX_MS = 4_294_967; const DASH = [6, 4]; const COLORS: readonly (keyof ChartTokens)[] = ["series1", "series2", "series3", "ctx"]; +/** The conversion registers plus the tick rate a burst's time axis needs. */ +const STREAM_REGISTERS: readonly string[] = [...CONFIG_REGISTERS, ...BIAS_REGISTERS, "tick_hz"]; interface StreamConfig extends TelemetryConfig { tickHz: number; @@ -125,8 +122,9 @@ function saveCsv(text: string, name: string): void { } export function StreamTab({ id }: { id: number }) { - const { descriptor, descriptorError, run } = useSession(); - const [config, setConfig] = useState(); + const { descriptor, descriptorError } = useSession(); + const bus = useBus(); + const stream = useReadOnce(id, STREAM_REGISTERS, [descriptor]); const [fields, setFields] = useState(DEFAULT_FIELDS); const [count, setCount] = useState(String(DEFAULT_COUNT)); const [windowMs, setWindowMs] = useState(String(DEFAULT_WINDOW_MS)); @@ -139,33 +137,14 @@ export function StreamTab({ id }: { id: number }) { const windowId = useId(); const fieldsId = useId(); - useEffect(() => { - if (descriptor === undefined) return; - const all = descriptor.fields(); - const config = spanOver(all, [...CONFIG_REGISTERS, "tick_hz"]); - const biases = spanOver(all, BIAS_REGISTERS); - let live = true; - run(async (c) => { - const read = decodeSpan(config, await c.read(id, config.addr, config.count)); - const bias = decodeSpan(biases, await c.read(id, biases.addr, biases.count)); - return { - sense: senseFromTable(read), - cal: calibrationFromTable(read), - biases: biasesFromTable(bias), - tickHz: read("tick_hz"), - }; - }).then( - (c) => { - if (live) setConfig(c); - }, - (e: unknown) => { - if (live) setError(message(e)); - }, - ); - return () => { - live = false; - }; - }, [descriptor, id, run]); + const { snapshot } = stream; + const config = useMemo( + () => + snapshot === undefined + ? undefined + : { ...configFrom(snapshot.read), tickHz: snapshot.read("tick_hz") }, + [snapshot], + ); const calibrated = config !== undefined && calibrationStatus(config.cal).valid; const raw = !calibrated || unitsPref === "raw"; @@ -199,7 +178,8 @@ export function StreamTab({ id }: { id: number }) { return; } setPending(true); - run((c) => c.telBurst(id, descriptor, mask, samples, window * 1000)) + bus + .command((c) => c.telBurst(id, descriptor, mask, samples, window * 1000)) .then( (burst) => { const rows = decodeBurst(burst.frames, mask); @@ -215,7 +195,7 @@ export function StreamTab({ id }: { id: number }) { }); }; - const problem = error ?? descriptorError; + const problem = error ?? (descriptor === undefined ? undefined : stream.error) ?? descriptorError; const heading = `${fieldsId}-chart`; const title = (family: "position" | "electrical") => [...new Set(units.filter((u) => familyOf(u.key) === family).map((u) => u.unit))].join(", "); diff --git a/src/lib/control.test.ts b/src/lib/control.test.ts index 4df18b4..f9774dc 100644 --- a/src/lib/control.test.ts +++ b/src/lib/control.test.ts @@ -1,5 +1,6 @@ -import { readFileSync } from "node:fs"; -import { afterEach, beforeEach, expect, test, vi } from "vitest"; +import type { Field } from "@openservocore/client"; +import { expect, test } from "vitest"; +import descriptor from "../../../open-servo-core/descriptors/osc-servo/0.1.json"; import fixture from "../../tests/fixtures/stall-24mhz.json"; import { clampGoal, @@ -11,7 +12,6 @@ import { goalSpec, isMode, isWindow, - LatestWins, LIMIT_REGISTERS, limitsFromTable, MODES, @@ -21,15 +21,11 @@ import { type GoalContext, type Limits, } from "./control"; -import { spanOver, type FieldInfo, type Sample } from "./telemetry-poll"; +import { decodeSpan, planSpans } from "./bus/spans"; +import type { Sample } from "./telemetry"; import { ADC_MAX_COUNT, ampsPerCount, senseFromTable, type Calibration } from "./units"; -const { fields } = JSON.parse( - readFileSync( - new URL("../../../open-servo-core/descriptors/osc-servo/0.1.json", import.meta.url), - "utf8", - ), -) as { fields: (FieldInfo & { variants?: { name: string; value: number }[] })[] }; +const fields = descriptor.fields as Field[]; const sense = senseFromTable((name) => { const v = (fixture.meta.sense as Record)[name]; @@ -47,13 +43,13 @@ const swapped: Calibration = { ...cal, rawMin: 3800, rawMax: 200 }; const limits: Limits = { dutyMaxQ15: 30000, velocityLimitCps: 4000, currentLimitCounts: 280 }; const ctx: GoalContext = { cal, sense, limits, raw: false }; -test("the control and limit spans over the 0.1 descriptor fit one READ each", () => { - expect(spanOver(fields, CONTROL_REGISTERS)).toMatchObject({ addr: 384, count: 18 }); - expect(spanOver(fields, LIMIT_REGISTERS)).toMatchObject({ addr: 54, count: 20 }); +test("the control and limit registers plan one read each", () => { + expect(planSpans(fields, CONTROL_REGISTERS)).toEqual([{ addr: 384, count: 18 }]); + expect(planSpans(fields, LIMIT_REGISTERS)).toEqual([{ addr: 54, count: 20 }]); }); test("decodeControl reads the switch, the mode and every goal", () => { - const span = spanOver(fields, CONTROL_REGISTERS); + const span = { addr: 384, count: 18 }; const bytes = new Uint8Array(span.count); const view = new DataView(bytes.buffer); const at = (name: string) => { @@ -67,7 +63,7 @@ test("decodeControl reads the switch, the mode and every goal", () => { view.setInt32(at("goal_position"), 2048, true); view.setInt32(at("goal_velocity"), -600, true); view.setInt16(at("goal_current"), 150, true); - expect(decodeControl(span, bytes)).toEqual({ + expect(decodeControl(decodeSpan(fields, span, bytes))).toEqual({ torque: true, mode: 2, goals: { goal_duty: -1234, goal_position: 2048, goal_velocity: -600, goal_current: 150 }, @@ -167,101 +163,3 @@ test("isWindow accepts the listed windows only", () => { expect(isWindow(30)).toBe(true); expect(isWindow(5)).toBe(false); }); - -beforeEach(() => { - vi.useFakeTimers(); -}); - -afterEach(() => { - vi.useRealTimers(); -}); - -function deferredSender() { - const sent: number[] = []; - const resolvers: (() => void)[] = []; - const send = (v: number) => - new Promise((resolve) => { - sent.push(v); - resolvers.push(resolve); - }); - const settle = async () => { - resolvers.shift()?.(); - await vi.advanceTimersByTimeAsync(0); - }; - return { sent, send, settle }; -} - -test("the first value goes out at once and a burst collapses to its newest", async () => { - const { sent, send, settle } = deferredSender(); - const errors: unknown[] = []; - const w = new LatestWins(send, 200, (e) => errors.push(e)); - w.push(1); - expect(sent).toEqual([1]); - w.push(2); - w.push(3); - w.push(4); - expect(sent).toEqual([1]); - await settle(); - await vi.advanceTimersByTimeAsync(199); - expect(sent).toEqual([1]); - await vi.advanceTimersByTimeAsync(1); - expect(sent).toEqual([1, 4]); - await settle(); - await vi.advanceTimersByTimeAsync(200); - expect(sent).toEqual([1, 4]); - expect(errors).toEqual([]); -}); - -test("a value pushed within the gap after a settled send waits for the gap", async () => { - const { sent, send, settle } = deferredSender(); - const w = new LatestWins(send, 200, () => undefined); - w.push(1); - await settle(); - await vi.advanceTimersByTimeAsync(50); - w.push(2); - expect(sent).toEqual([1]); - await vi.advanceTimersByTimeAsync(150); - expect(sent).toEqual([1, 2]); -}); - -test("a value pushed after the gap has passed goes out at once", async () => { - const { sent, send, settle } = deferredSender(); - const w = new LatestWins(send, 200, () => undefined); - w.push(1); - await settle(); - await vi.advanceTimersByTimeAsync(200); - w.push(2); - expect(sent).toEqual([1, 2]); -}); - -test("a rejected send reports the error and the next value still goes out", async () => { - const errors: unknown[] = []; - const sent: number[] = []; - const w = new LatestWins( - (v) => { - sent.push(v); - return v === 1 ? Promise.reject(new Error("boom")) : Promise.resolve(); - }, - 200, - (e) => errors.push(e), - ); - w.push(1); - await vi.advanceTimersByTimeAsync(0); - expect(errors).toHaveLength(1); - w.push(2); - await vi.advanceTimersByTimeAsync(200); - expect(sent).toEqual([1, 2]); -}); - -test("stop drops the pending value and ignores later pushes", async () => { - const { sent, send, settle } = deferredSender(); - const w = new LatestWins(send, 200, () => undefined); - w.push(1); - w.push(2); - w.stop(); - await settle(); - await vi.advanceTimersByTimeAsync(500); - w.push(3); - await vi.advanceTimersByTimeAsync(500); - expect(sent).toEqual([1]); -}); diff --git a/src/lib/control.ts b/src/lib/control.ts index a06f4e5..6b46d20 100644 --- a/src/lib/control.ts +++ b/src/lib/control.ts @@ -1,8 +1,7 @@ // The Live page's control cluster, minus React: the servo's mode and goal -// registers, each goal's range and unit conversion, and the write policy -// behind a slider drag. +// registers, and each goal's range and unit conversion. -import { decodeSpan, type Sample, type Span } from "./telemetry-poll"; +import type { Sample } from "./telemetry"; import { ADC_MAX_COUNT, ampsPerCount, @@ -25,9 +24,6 @@ export function isWindow(value: number): value is WindowS { return WINDOWS_S.some((w) => w === value); } -/** Gap between goal writes while a slider drags: at most 5 per second. */ -export const GOAL_WRITE_GAP_MS = 200; - /** Duties are q15 fractions of full drive (firmware regions/control.rs). */ const Q15 = 2 ** 15; /** goal_duty is an i16, so full drive itself is one count out of reach. */ @@ -78,8 +74,7 @@ export interface Limits { currentLimitCounts: number; } -export function decodeControl(span: Span, bytes: Uint8Array): ControlState { - const read = decodeSpan(span, bytes); +export function decodeControl(read: ReadRegister): ControlState { return { torque: read("torque_enable") !== 0, mode: read("mode"), @@ -230,54 +225,3 @@ export interface GoalSpec extends GoalUnits { export function goalSpec(mode: ModeName, ctx: GoalContext): GoalSpec { return { ...goalUnits(mode, ctx), range: goalRange(mode, ctx.cal, ctx.limits) }; } - -/** - * Latest wins: a burst of values collapses to the newest, one send runs at a - * time, and the next starts no sooner than `gapMs` after the previous settled. - * The first value of a burst goes out at once. - */ -export class LatestWins { - private next: { value: T } | undefined; - private busy = false; - private timer: ReturnType | undefined; - private stopped = false; - - constructor( - private readonly send: (value: T) => Promise, - private readonly gapMs: number, - private readonly onError: (error: unknown) => void, - ) {} - - push(value: T): void { - if (this.stopped) return; - this.next = { value }; - if (!this.busy && this.timer === undefined) this.flush(); - } - - /** Drops what has not been sent; a send already running finishes. */ - stop(): void { - this.stopped = true; - this.next = undefined; - clearTimeout(this.timer); - this.timer = undefined; - } - - private flush(): void { - const next = this.next; - if (next === undefined) return; - this.next = undefined; - this.busy = true; - this.send(next.value) - .catch((e: unknown) => { - if (!this.stopped) this.onError(e); - }) - .finally(() => { - this.busy = false; - if (this.stopped) return; - this.timer = setTimeout(() => { - this.timer = undefined; - this.flush(); - }, this.gapMs); - }); - } -} diff --git a/src/lib/stream.test.ts b/src/lib/stream.test.ts index 704587b..e11afd3 100644 --- a/src/lib/stream.test.ts +++ b/src/lib/stream.test.ts @@ -14,7 +14,7 @@ import { toCsv, unitsFor, } from "./stream"; -import type { TelemetryConfig } from "./telemetry-poll"; +import type { TelemetryConfig } from "./telemetry"; import { calibrationFromTable, senseFromTable } from "./units"; /** Mirrors core tel.rs `sample(i)`: every field, window_valid on even i. */ diff --git a/src/lib/stream.ts b/src/lib/stream.ts index 15fb9f3..418c0d8 100644 --- a/src/lib/stream.ts +++ b/src/lib/stream.ts @@ -3,7 +3,7 @@ // that mirrors the ident host's StreamAssembler, and the CSV export. import { dutyPercent } from "./control"; -import type { TelemetryConfig } from "./telemetry-poll"; +import type { TelemetryConfig } from "./telemetry"; import { busV, currentMa, diff --git a/src/lib/telemetry-poll.test.ts b/src/lib/telemetry-poll.test.ts deleted file mode 100644 index 257f7ef..0000000 --- a/src/lib/telemetry-poll.test.ts +++ /dev/null @@ -1,280 +0,0 @@ -import { readFileSync } from "node:fs"; -import { afterEach, beforeEach, expect, test, vi } from "vitest"; -import { - BIAS_REGISTERS, - CONFIG_REGISTERS, - decodeSample, - decodeSpan, - READ_MAX, - registersOf, - SAMPLE_REGISTERS, - SampleRing, - spanOver, - startTelemetry, - type FieldInfo, - type Sample, - type TelemetryConfig, -} from "./telemetry-poll"; - -const descriptor = JSON.parse( - readFileSync( - new URL("../../../open-servo-core/descriptors/osc-servo/0.1.json", import.meta.url), - "utf8", - ), -) as { fields: FieldInfo[] }; -const { fields } = descriptor; - -const field = (name: string, addr: number, width: number, kind: FieldInfo["kind"]): FieldInfo => ({ - name, - addr, - width, - kind, -}); - -test("spanOver covers the named fields from the lowest address to the highest end", () => { - const span = spanOver( - [field("a", 10, 2, "uint"), field("b", 4, 4, "int"), field("c", 20, 1, "uint")], - ["a", "c"], - ); - expect(span).toMatchObject({ addr: 10, count: 11 }); - expect(span.fields.map((f) => f.name)).toEqual(["a", "c"]); -}); - -test("spanOver rejects a missing register and a span over one READ", () => { - expect(() => spanOver([field("a", 0, 2, "uint")], ["b"])).toThrow("descriptor has no b"); - expect(() => - spanOver([field("a", 0, 2, "uint"), field("z", READ_MAX, 1, "uint")], ["a", "z"]), - ).toThrow("do not fit one READ"); -}); - -test("the sample span over the 0.1 descriptor is goal_duty through ntc_raw", () => { - const span = spanOver(fields, SAMPLE_REGISTERS); - expect(span.addr).toBe(390); - expect(span.count).toBe(210); -}); - -test("the config and bias spans over the 0.1 descriptor fit one READ each", () => { - expect(spanOver(fields, CONFIG_REGISTERS)).toMatchObject({ addr: 128, count: 172 }); - expect(spanOver(fields, BIAS_REGISTERS)).toMatchObject({ addr: 594, count: 8 }); -}); - -test("registersOf lists what a binder reads, in order", () => { - expect(registersOf((read) => [read("x"), read("y")])).toEqual(["x", "y"]); -}); - -test("decodeSpan reads little-endian values by width and sign", () => { - const span = spanOver( - [ - field("u8", 100, 1, "uint"), - field("i8", 101, 1, "int"), - field("u16", 102, 2, "uint"), - field("i16", 104, 2, "int"), - field("u32", 106, 4, "uint"), - field("i32", 110, 4, "int"), - ], - ["u8", "i8", "u16", "i16", "u32", "i32"], - ); - const bytes = new Uint8Array([ - 0xff, 0xff, 0x34, 0x12, 0xfe, 0xff, 0x78, 0x56, 0x34, 0x12, 0xfd, 0xff, 0xff, 0xff, - ]); - const read = decodeSpan(span, 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(-3); - expect(() => read("nope")).toThrow("span has no nope"); - const flags = spanOver([field("on", 0, 1, "bool"), field("mode", 1, 1, "enum")], ["on", "mode"]); - const readFlags = decodeSpan(flags, new Uint8Array([1, 3])); - expect(readFlags("on")).toBe(1); - expect(readFlags("mode")).toBe(3); - expect(() => decodeSpan(span, bytes.subarray(1))).toThrow("expected 14"); -}); - -test("decodeSample scales omega_hat_cps out of Q16 and keeps the rest in counts", () => { - const span = spanOver(fields, SAMPLE_REGISTERS); - const bytes = new Uint8Array(span.count); - const view = new DataView(bytes.buffer); - const at = (name: string) => { - const f = fields.find((x) => x.name === name); - if (f === undefined) throw new Error(name); - return f.addr - span.addr; - }; - view.setInt16(at("goal_duty"), -1000, true); - view.setInt32(at("goal_position"), 2048, true); - view.setInt32(at("goal_velocity"), -500, true); - view.setInt16(at("goal_current"), 250, true); - view.setUint8(at("mode_active"), 3); - view.setInt32(at("omega_hat_cps"), -3 * 65536, true); - view.setInt16(at("duty_applied_q15"), -900, true); - view.setUint16(at("pos"), 1234, true); - view.setUint16(at("current"), 300, true); - view.setUint16(at("vmotor_a"), 800, true); - view.setUint16(at("vmotor_b"), 700, true); - view.setUint16(at("vbus_raw"), 3600, true); - view.setUint16(at("ntc_raw"), 2000, true); - expect(decodeSample(span, bytes, 1.5)).toEqual({ - t: 1.5, - pos: 1234, - goal: 2048, - goalVelocity: -500, - goalCurrent: 250, - goalDuty: -1000, - dutyApplied: -900, - modeActive: 3, - velocity: -3, - current: 300, - vbus: 3600, - vmotorA: 800, - vmotorB: 700, - ntc: 2000, - }); -}); - -const sampleAt = (t: number): Sample => ({ - t, - pos: 0, - goal: 0, - goalVelocity: 0, - goalCurrent: 0, - goalDuty: 0, - dutyApplied: 0, - modeActive: 0, - velocity: 0, - current: 0, - vbus: 0, - vmotorA: 0, - vmotorB: 0, - ntc: 0, -}); - -test("the ring keeps the window's worth of samples, oldest first", () => { - const ring = new SampleRing(30); - for (const t of [0, 10, 20, 30, 31]) ring.push(sampleAt(t)); - expect(ring.samples.map((s) => s.t)).toEqual([10, 20, 30, 31]); - expect(ring.samples).not.toBe(ring.samples); -}); - -beforeEach(() => { - vi.useFakeTimers(); -}); -afterEach(() => { - vi.useRealTimers(); -}); - -/** A fake servo whose registers carry their own address in each byte. */ -function fakeRead(): { - calls: [number, number][]; - read: (addr: number, count: number) => Promise; - pending: Set<(bytes: Uint8Array | undefined) => void>; -} { - const calls: [number, number][] = []; - const pending = new Set<(bytes: Uint8Array | undefined) => void>(); - return { - calls, - pending, - read: (addr, count) => { - calls.push([addr, count]); - return new Promise((resolve) => { - pending.add((bytes) => { - pending.clear(); - resolve(bytes ?? Uint8Array.from({ length: count }, (_, i) => (addr + i) & 0xff)); - }); - }); - }, - }; -} - -function settle(fake: ReturnType, bytes?: Uint8Array): void { - for (const resolve of fake.pending) resolve(bytes); -} - -test("startTelemetry reads config, then biases, then samples on every tick", async () => { - const fake = fakeRead(); - const configs: TelemetryConfig[] = []; - const samples: Sample[] = []; - let now = 0; - const stop = startTelemetry({ - fields, - read: fake.read, - now: () => now, - periodMs: 100, - onConfig: (c) => configs.push(c), - onSample: (s) => samples.push(s), - onError: (e: unknown) => { - throw e; - }, - }); - expect(fake.calls).toEqual([[128, 172]]); - settle(fake); - await vi.advanceTimersByTimeAsync(100); - expect(fake.calls).toEqual([ - [128, 172], - [594, 8], - ]); - settle(fake); - await vi.advanceTimersByTimeAsync(0); - expect(configs).toHaveLength(1); - expect(configs[0]?.sense.shuntMohm).toBe(0xf3f2); - expect(configs[0]?.biases.currentBiasCounts).toBe(0x5352); - now = 7; - await vi.advanceTimersByTimeAsync(100); - expect(fake.calls.at(-1)).toEqual([390, 210]); - settle(fake); - await vi.advanceTimersByTimeAsync(0); - expect(samples).toHaveLength(1); - expect(samples[0]?.t).toBe(7); - expect(samples[0]?.pos).toBe(0x4140); - stop(); - await vi.advanceTimersByTimeAsync(500); - expect(fake.calls).toHaveLength(3); -}); - -test("a tick whose read is still pending is skipped, not queued", async () => { - const fake = fakeRead(); - startTelemetry({ - fields, - read: fake.read, - now: () => 0, - periodMs: 100, - onConfig: () => undefined, - onSample: () => undefined, - onError: () => undefined, - }); - await vi.advanceTimersByTimeAsync(350); - expect(fake.calls).toHaveLength(4); - const samples: Sample[] = []; - const skipping = startTelemetry({ - fields, - read: () => Promise.resolve(undefined), - now: () => 0, - periodMs: 100, - onConfig: () => undefined, - onSample: (s) => samples.push(s), - onError: () => undefined, - }); - await vi.advanceTimersByTimeAsync(1000); - expect(samples).toHaveLength(0); - skipping(); -}); - -test("the first error stops the poll and is reported once", async () => { - const errors: unknown[] = []; - let calls = 0; - startTelemetry({ - fields, - read: () => { - calls++; - return Promise.reject(new Error("busy")); - }, - now: () => 0, - periodMs: 100, - onConfig: () => undefined, - onSample: () => undefined, - onError: (e) => errors.push(e), - }); - await vi.advanceTimersByTimeAsync(500); - expect(calls).toBe(1); - expect(errors).toHaveLength(1); -}); diff --git a/src/lib/telemetry-poll.ts b/src/lib/telemetry-poll.ts deleted file mode 100644 index 959c32e..0000000 --- a/src/lib/telemetry-poll.ts +++ /dev/null @@ -1,246 +0,0 @@ -// The Live page's telemetry poll: the conversion registers once, then one READ -// over the sample registers per tick. Everything but `startTelemetry` is pure. - -import type { Kind } from "@openservocore/client"; -import { - biasesFromTable, - calibrationFromTable, - senseFromTable, - type Biases, - type Calibration, - type ReadRegister, - type Sense, -} from "./units"; - -export const POLL_HZ = 10; -/** Largest READ reply one frame carries (protocol sec 3.1). */ -export const READ_MAX = 252; -/** omega_hat_cps is csQ16, (counts/s) x 2^16 (firmware regions/telemetry.rs). */ -const Q16 = 2 ** 16; - -/** The descriptor's `Field`, narrowed to what a span needs. */ -export interface FieldInfo { - name: string; - addr: number; - width: number; - kind: Kind; -} - -/** One contiguous READ covering `fields`. */ -export interface Span { - addr: number; - count: number; - fields: FieldInfo[]; -} - -/** - * One poll in device counts: `velocity` and `goalVelocity` are counts/s, the - * duties q15, `modeActive` the `mode` enum's discriminant, `t` seconds. - */ -export interface Sample { - t: number; - pos: number; - goal: number; - goalVelocity: number; - goalCurrent: number; - goalDuty: number; - dutyApplied: number; - modeActive: number; - velocity: number; - current: number; - vbus: number; - vmotorA: number; - vmotorB: number; - ntc: number; -} - -export interface TelemetryConfig { - sense: Sense; - cal: Calibration; - biases: Biases; -} - -export const SAMPLE_REGISTERS: readonly string[] = [ - "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", -]; - -/** The registers a units binder reads, so the list lives in one place. */ -export function registersOf(binder: (read: ReadRegister) => unknown): string[] { - const names: string[] = []; - binder((name) => { - names.push(name); - return 0; - }); - return names; -} - -export const CONFIG_REGISTERS: readonly string[] = [ - ...registersOf(senseFromTable), - ...registersOf(calibrationFromTable), -]; -export const BIAS_REGISTERS: readonly string[] = registersOf(biasesFromTable); - -export function spanOver(fields: readonly FieldInfo[], names: readonly string[]): Span { - const picked = names.map((name) => { - const field = fields.find((f) => f.name === name); - if (field === undefined) throw new Error(`descriptor has no ${name}`); - return field; - }); - const addr = Math.min(...picked.map((f) => f.addr)); - const count = Math.max(...picked.map((f) => f.addr + f.width)) - addr; - if (count > READ_MAX) throw new Error(`${count} bytes do not fit one READ`); - return { addr, count, fields: picked }; -} - -// Decodes from the span's own layout, never the descriptor object: a new -// selection frees that while a read is still in flight. -export function decodeSpan(span: Span, bytes: Uint8Array): ReadRegister { - if (bytes.length !== span.count) { - throw new Error(`read returned ${bytes.length} bytes, expected ${span.count}`); - } - const view = new DataView(bytes.buffer, bytes.byteOffset, bytes.byteLength); - return (name) => { - const field = span.fields.find((f) => f.name === name); - if (field === undefined) throw new Error(`span has no ${name}`); - return decodeField(view, field.addr - span.addr, field); - }; -} - -/** Bools and enums decode as their unsigned byte. */ -function decodeField(view: DataView, offset: number, { name, width, kind }: FieldInfo): number { - if (kind === "bytes") throw new Error(`${name} is ${kind}, not a number`); - const signed = kind === "int"; - switch (width) { - case 1: - return signed ? view.getInt8(offset) : view.getUint8(offset); - case 2: - return signed ? view.getInt16(offset, true) : view.getUint16(offset, true); - case 4: - return signed ? view.getInt32(offset, true) : view.getUint32(offset, true); - default: - throw new Error(`${name} is ${width} bytes wide`); - } -} - -export function decodeSample(span: Span, bytes: Uint8Array, t: number): Sample { - const read = decodeSpan(span, bytes); - return { - t, - pos: read("pos"), - goal: read("goal_position"), - goalVelocity: read("goal_velocity"), - goalCurrent: read("goal_current"), - goalDuty: read("goal_duty"), - dutyApplied: read("duty_applied_q15"), - modeActive: read("mode_active"), - velocity: read("omega_hat_cps") / Q16, - current: read("current"), - vbus: read("vbus_raw"), - vmotorA: read("vmotor_a"), - vmotorB: read("vmotor_b"), - ntc: read("ntc_raw"), - }; -} - -/** The last `windowS` seconds of samples, oldest first. */ -export class SampleRing { - private buf: Sample[] = []; - - constructor(private readonly windowS: number) {} - - push(sample: Sample): void { - this.buf.push(sample); - const cutoff = sample.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); - } - - /** A fresh array each call, so a render can key on it. */ - get samples(): Sample[] { - return this.buf.slice(); - } -} - -export interface TelemetryOptions { - fields: readonly FieldInfo[]; - /** Resolves undefined when the read was skipped (a command was in flight). */ - read: (addr: number, count: number) => Promise; - /** Seconds. */ - now: () => number; - periodMs: number; - onConfig: (config: TelemetryConfig) => void; - onSample: (sample: Sample) => void; - onError: (error: unknown) => void; -} - -type Phase = { name: "config" } | { name: "biases"; config: ReadRegister } | { name: "samples" }; - -/** - * Ticks every `periodMs`: the sense and calibration span, then the bias span, - * then the sample span until stopped. A tick whose read is skipped retries the - * same phase next tick. The first error stops the poll. - */ -export function startTelemetry(o: TelemetryOptions): () => void { - const spans = { - config: spanOver(o.fields, CONFIG_REGISTERS), - biases: spanOver(o.fields, BIAS_REGISTERS), - samples: spanOver(o.fields, SAMPLE_REGISTERS), - }; - let phase: Phase = { name: "config" }; - let stopped = false; - - async function tick(): Promise { - const span = spans[phase.name]; - const bytes = await o.read(span.addr, span.count); - if (stopped || bytes === undefined) return; - switch (phase.name) { - case "config": - phase = { name: "biases", config: decodeSpan(span, bytes) }; - break; - case "biases": { - const { config } = phase; - o.onConfig({ - sense: senseFromTable(config), - cal: calibrationFromTable(config), - biases: biasesFromTable(decodeSpan(span, bytes)), - }); - phase = { name: "samples" }; - break; - } - case "samples": - o.onSample(decodeSample(span, bytes, o.now())); - break; - } - } - - const stop = (): void => { - stopped = true; - clearInterval(timer); - }; - const guarded = (): void => { - tick().catch((e: unknown) => { - if (stopped) return; - stop(); - o.onError(e); - }); - }; - const timer = setInterval(guarded, o.periodMs); - guarded(); - return stop; -} diff --git a/src/lib/telemetry.test.ts b/src/lib/telemetry.test.ts new file mode 100644 index 0000000..2d305ba --- /dev/null +++ b/src/lib/telemetry.test.ts @@ -0,0 +1,86 @@ +import type { Field } from "@openservocore/client"; +import { expect, test } from "vitest"; +import descriptor from "../../../open-servo-core/descriptors/osc-servo/0.1.json"; +import { decodeSpan, planSpans } from "./bus/spans"; +import { + BIAS_REGISTERS, + CONFIG_REGISTERS, + configFrom, + decodeSample, + registersOf, + SAMPLE_REGISTERS, +} from "./telemetry"; + +const fields = descriptor.fields as Field[]; + +test("the sample registers plan one read, goal_duty through ntc_raw", () => { + expect(planSpans(fields, SAMPLE_REGISTERS)).toEqual([{ addr: 390, count: 210 }]); +}); + +test("the conversion registers plan one read each side of the table", () => { + expect(planSpans(fields, CONFIG_REGISTERS)).toEqual([{ addr: 128, count: 172 }]); + expect(planSpans(fields, BIAS_REGISTERS)).toEqual([{ addr: 594, count: 8 }]); +}); + +test("registersOf lists what a binder reads, in order", () => { + expect(registersOf((read) => [read("x"), read("y")])).toEqual(["x", "y"]); +}); + +test("decodeSample scales omega_hat_cps out of Q16 and keeps the rest in counts", () => { + const span = { addr: 390, count: 210 }; + const bytes = new Uint8Array(span.count); + const view = new DataView(bytes.buffer); + const at = (name: string) => { + const f = fields.find((x) => x.name === name); + if (f === undefined) throw new Error(name); + return f.addr - span.addr; + }; + view.setInt16(at("goal_duty"), -1000, true); + view.setInt32(at("goal_position"), 2048, true); + view.setInt32(at("goal_velocity"), -500, true); + view.setInt16(at("goal_current"), 250, true); + view.setUint8(at("mode_active"), 3); + view.setInt32(at("omega_hat_cps"), -3 * 65536, true); + view.setInt16(at("duty_applied_q15"), -900, true); + view.setUint16(at("pos"), 1234, true); + view.setUint16(at("current"), 300, true); + view.setUint16(at("vmotor_a"), 800, true); + view.setUint16(at("vmotor_b"), 700, true); + view.setUint16(at("vbus_raw"), 3600, true); + view.setUint16(at("ntc_raw"), 2000, true); + expect(decodeSample(decodeSpan(fields, span, bytes), 1.5)).toEqual({ + t: 1.5, + pos: 1234, + goal: 2048, + goalVelocity: -500, + goalCurrent: 250, + goalDuty: -1000, + dutyApplied: -900, + modeActive: 3, + velocity: -3, + current: 300, + vbus: 3600, + vmotorA: 800, + vmotorB: 700, + ntc: 2000, + }); +}); + +test("configFrom binds the sense, the calibration and the biases from one reader", () => { + const values: Record = { + shunt_r_mohm: 10, + gain_milli: 32000, + vdd_mv: 3300, + raw_min: 200, + raw_max: 3800, + angle_min_cdeg: 0, + angle_max_cdeg: 18000, + gear_ratio_centi: 100, + current_bias_counts: 2048, + vmotor_bias_counts: 1024, + }; + const config = configFrom((name) => values[name] ?? 1); + expect(config.sense.shuntMohm).toBe(10); + expect(config.cal.rawMax).toBe(3800); + expect(config.biases.currentBiasCounts).toBe(2048); +}); diff --git a/src/lib/telemetry.ts b/src/lib/telemetry.ts new file mode 100644 index 0000000..df6b6b4 --- /dev/null +++ b/src/lib/telemetry.ts @@ -0,0 +1,104 @@ +// The Live page's telemetry: the registers one fast subscription carries, the +// conversion registers behind them, and the decode from a snapshot's reader. + +import { + biasesFromTable, + calibrationFromTable, + senseFromTable, + type Biases, + type Calibration, + type ReadRegister, + type Sense, +} from "./units"; + +/** The nominal rate of the bus manager's fast class, which the header prints. */ +export const POLL_HZ = 10; +/** omega_hat_cps is csQ16, (counts/s) x 2^16 (firmware regions/telemetry.rs). */ +const Q16 = 2 ** 16; + +/** + * One sample in device counts: `velocity` and `goalVelocity` are counts/s, the + * duties q15, `modeActive` the `mode` enum's discriminant, `t` seconds. + */ +export interface Sample { + t: number; + pos: number; + goal: number; + goalVelocity: number; + goalCurrent: number; + goalDuty: number; + dutyApplied: number; + modeActive: number; + velocity: number; + current: number; + vbus: number; + vmotorA: number; + vmotorB: number; + ntc: number; +} + +export interface TelemetryConfig { + sense: Sense; + cal: Calibration; + biases: Biases; +} + +export const SAMPLE_REGISTERS: readonly string[] = [ + "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", +]; + +/** The registers a units binder reads, so the list lives in one place. */ +export function registersOf(binder: (read: ReadRegister) => unknown): string[] { + const names: string[] = []; + binder((name) => { + names.push(name); + return 0; + }); + return names; +} + +export const CONFIG_REGISTERS: readonly string[] = [ + ...registersOf(senseFromTable), + ...registersOf(calibrationFromTable), +]; +export const BIAS_REGISTERS: readonly string[] = registersOf(biasesFromTable); + +export function decodeSample(read: ReadRegister, t: number): Sample { + return { + t, + pos: read("pos"), + goal: read("goal_position"), + goalVelocity: read("goal_velocity"), + goalCurrent: read("goal_current"), + goalDuty: read("goal_duty"), + dutyApplied: read("duty_applied_q15"), + modeActive: read("mode_active"), + velocity: read("omega_hat_cps") / Q16, + current: read("current"), + vbus: read("vbus_raw"), + vmotorA: read("vmotor_a"), + vmotorB: read("vmotor_b"), + ntc: read("ntc_raw"), + }; +} + +/** The conversion constants, from one read over CONFIG_REGISTERS and BIAS_REGISTERS. */ +export function configFrom(read: ReadRegister): TelemetryConfig { + return { + sense: senseFromTable(read), + cal: calibrationFromTable(read), + biases: biasesFromTable(read), + }; +} diff --git a/src/routes/live.tsx b/src/routes/live.tsx index fd98e2b..27e0cad 100644 --- a/src/routes/live.tsx +++ b/src/routes/live.tsx @@ -1,4 +1,4 @@ -import type { Descriptor, Field } from "@openservocore/client"; +import type { Field, Value } from "@openservocore/client"; import { createFileRoute, Link } from "@tanstack/react-router"; import { Pause, Play, TriangleAlert } from "lucide-react"; import { useEffect, useId, useMemo, useRef, useState, type KeyboardEvent } from "react"; @@ -28,42 +28,39 @@ import { Slider } from "@/components/ui/slider"; import { Switch } from "@/components/ui/switch"; import { Tabs, TabsContent, TabsList, TabsTrigger } from "@/components/ui/tabs"; import { Tooltip, TooltipContent, TooltipTrigger } from "@/components/ui/tooltip"; +import { useBus, useReadOnce, useRegisters, useRing } from "@/lib/bus/hooks"; import { useChartTokens, type ChartTokens } from "@/lib/chart-theme"; import { clampGoal, CONTROL_REGISTERS, decodeControl, dutyPercent, - GOAL_WRITE_GAP_MS, goalOf, goalSpec, goalUnits, isWindow, - LatestWins, LIMIT_REGISTERS, limitsFromTable, modeLabel, modeName, WINDOWS_S, - type ControlState, type GoalRegister, type GoalUnits, - type Limits, type ModeName, type WindowS, } from "@/lib/control"; import { isUnits, type Units } from "@/lib/prefs"; -import { useSession, type Session } from "@/lib/session"; +import { useSession } from "@/lib/session"; import { - decodeSpan, + BIAS_REGISTERS, + CONFIG_REGISTERS, + configFrom, + decodeSample, POLL_HZ, - SampleRing, - spanOver, - startTelemetry, + SAMPLE_REGISTERS, type Sample, - type Span, type TelemetryConfig, -} from "@/lib/telemetry-poll"; +} from "@/lib/telemetry"; import { busV, calibrationStatus, @@ -101,6 +98,10 @@ interface PanelDef { const DASH = [6, 4]; /** The ring keeps the longest window, so a shorter one is a view over the same samples. */ const RING_S: WindowS = 30; +/** One read behind every conversion the panels and the goal units need. */ +const CONVERSION_REGISTERS: readonly string[] = [...CONFIG_REGISTERS, ...BIAS_REGISTERS]; +/** Inside the sample span, so the mode costs no exchange of its own. */ +const MODE_REGISTERS: readonly string[] = ["mode"]; const POSITION: SeriesDef = { key: "position", label: "Position", token: "series1" }; const VELOCITY: SeriesDef = { key: "velocity", label: "Velocity", token: "series2", right: true }; @@ -281,49 +282,25 @@ function LivePage() { } function Telemetry({ id }: { id: number }) { - const { descriptor, descriptorError, run } = useSession(); - const [config, setConfig] = useState(); - const [samples, setSamples] = useState([]); - /** The samples on screen while paused; the poll keeps filling the ring behind them. */ + const { descriptor, descriptorError } = useSession(); + const snapshots = useRing(id, SAMPLE_REGISTERS, "fast", RING_S); + const conversion = useReadOnce(id, CONVERSION_REGISTERS, [descriptor]); + const modeSnapshot = useRegisters(id, MODE_REGISTERS, "fast"); + /** The samples on screen while paused; the subscription keeps filling the ring behind them. */ const [frozen, setFrozen] = useState(); const [windowS, setWindowS] = useState(RING_S); - const [error, setError] = useState(); const [shown, setShown] = useState(initialShown); - /** The servo's mode, switch and goals as last read back. */ - const [control, setControl] = useState(); const [unitsPref, setUnitsPref] = useUnitsPref(); const tokens = useChartTokens(); - useEffect(() => { - if (descriptor === undefined) return; - const ring = new SampleRing(RING_S); - let pending = false; - return startTelemetry({ - fields: descriptor.fields(), - // A tick due while the last read is still queued is skipped, so a busy - // bus never piles reads up behind the writes. - read: async (addr, count) => { - if (pending) return undefined; - pending = true; - try { - return await run((c) => c.read(id, addr, count)); - } finally { - pending = false; - } - }, - now: () => performance.now() / 1000, - periodMs: 1000 / POLL_HZ, - onConfig: setConfig, - onSample: (s) => { - ring.push(s); - setSamples(ring.samples); - }, - onError: (e: unknown) => { - setError(message(e)); - }, - }); - }, [descriptor, id, run]); - + const config = useMemo( + () => (conversion.snapshot === undefined ? undefined : configFrom(conversion.snapshot.read)), + [conversion.snapshot], + ); + const samples = useMemo( + () => snapshots.filter((s) => !s.stale).map((s) => decodeSample(s.read, s.t)), + [snapshots], + ); const calibrated = config !== undefined && calibrationStatus(config.cal).valid; const raw = !calibrated || unitsPref === "raw"; const shownSamples = frozen ?? samples; @@ -336,9 +313,9 @@ function Telemetry({ id }: { id: number }) { // mode_active from it within one slow tick, and the simulated fleet never // runs the kernel that would. const mode = - modeField === undefined || control === undefined + modeField === undefined || modeSnapshot === undefined || modeSnapshot.stale ? undefined - : modeName(modeField.variants, control.mode); + : modeName(modeField.variants, modeSnapshot.read("mode")); const active = modeField === undefined || latestSample === undefined ? undefined @@ -355,8 +332,12 @@ function Telemetry({ id }: { id: number }) { () => (config === undefined ? [] : toRows(shownSamples, config, raw, goal, windowS)), [shownSamples, config, raw, goal, windowS], ); + const newest = snapshots.at(-1); const latest = rows.at(-1); - const problem = error ?? descriptorError; + const problem = + (descriptor === undefined ? undefined : conversion.error) ?? + (newest?.stale === true ? newest.error : undefined) ?? + descriptorError; const windowId = useId(); return ( @@ -427,18 +408,15 @@ function Telemetry({ id }: { id: number }) {