Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
62 changes: 21 additions & 41 deletions src/components/stream-tab.tsx
Original file line number Diff line number Diff line change
@@ -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";
Expand All @@ -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";
Expand All @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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<StreamConfig>();
const { descriptor, descriptorError } = useSession();
const bus = useBus();
const stream = useReadOnce(id, STREAM_REGISTERS, [descriptor]);
const [fields, setFields] = useState<readonly FieldKey[]>(DEFAULT_FIELDS);
const [count, setCount] = useState(String(DEFAULT_COUNT));
const [windowMs, setWindowMs] = useState(String(DEFAULT_WINDOW_MS));
Expand All @@ -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<StreamConfig | undefined>(
() =>
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";
Expand Down Expand Up @@ -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);
Expand All @@ -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(", ");
Expand Down
124 changes: 11 additions & 113 deletions src/lib/control.test.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -11,7 +12,6 @@ import {
goalSpec,
isMode,
isWindow,
LatestWins,
LIMIT_REGISTERS,
limitsFromTable,
MODES,
Expand All @@ -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<string, number>)[name];
Expand All @@ -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) => {
Expand All @@ -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 },
Expand Down Expand Up @@ -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<void>((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<number>(
(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]);
});
62 changes: 3 additions & 59 deletions src/lib/control.ts
Original file line number Diff line number Diff line change
@@ -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,
Expand All @@ -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. */
Expand Down Expand Up @@ -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"),
Expand Down Expand Up @@ -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<T> {
private next: { value: T } | undefined;
private busy = false;
private timer: ReturnType<typeof setTimeout> | undefined;
private stopped = false;

constructor(
private readonly send: (value: T) => Promise<void>,
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);
});
}
}
2 changes: 1 addition & 1 deletion src/lib/stream.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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. */
Expand Down
2 changes: 1 addition & 1 deletion src/lib/stream.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading