Execute oxfmt (#1166)

Ran
pnpm oxfmt .

This should get the web repo aligned so that we can better enforce oxfmt going forward.
This commit is contained in:
Austin
2026-06-16 20:30:08 -04:00
committed by GitHub
parent c5690d1edb
commit d2cb51d489
378 changed files with 5548 additions and 1974 deletions
+22 -5
View File
@@ -6,7 +6,10 @@ export type {
MeshClientOptions,
} from "./src/core/client/MeshClient.ts";
export { MeshRegistry } from "./src/core/registry/MeshRegistry.ts";
export type { ConnectionId, RegistryEntry } from "./src/core/registry/MeshRegistry.ts";
export type {
ConnectionId,
RegistryEntry,
} from "./src/core/registry/MeshRegistry.ts";
// Constants & errors
export { Constants } from "./src/core/constants/index.ts";
@@ -17,7 +20,11 @@ export {
} from "./src/core/errors/MeshError.ts";
// Shared signal primitives
export { createStore, SignalMap, toReadonly } from "./src/core/signals/createStore.ts";
export {
createStore,
SignalMap,
toReadonly,
} from "./src/core/signals/createStore.ts";
export type { ReadonlySignal } from "./src/core/signals/createStore.ts";
// Logging
@@ -38,7 +45,11 @@ export { Queue } from "./src/core/queue/Queue.ts";
export { Xmodem } from "./src/core/xmodem/Xmodem.ts";
// Transport interface
export type { DeviceOutput, HttpRetryConfig, Transport } from "./src/core/transport/Transport.ts";
export type {
DeviceOutput,
HttpRetryConfig,
Transport,
} from "./src/core/transport/Transport.ts";
export { DeviceStatusEnum } from "./src/core/transport/Transport.ts";
// Commonly-used runtime enums exported directly for ergonomic access.
@@ -109,7 +120,10 @@ export type {
RadioConfigSection,
} from "./src/features/config/index.ts";
export { InMemoryTelemetryRepository, TelemetryClient } from "./src/features/telemetry/index.ts";
export {
InMemoryTelemetryRepository,
TelemetryClient,
} from "./src/features/telemetry/index.ts";
export type {
TelemetryClientOptions,
TelemetryKind,
@@ -125,7 +139,10 @@ export { TraceRouteClient } from "./src/features/traceroute/index.ts";
export type { TraceRoute } from "./src/features/traceroute/index.ts";
export { FilesClient } from "./src/features/files/index.ts";
export type { FileTransfer, TransferStatus } from "./src/features/files/index.ts";
export type {
FileTransfer,
TransferStatus,
} from "./src/features/files/index.ts";
// Phase-A legacy shims (removed in Phase C)
export { MeshDevice } from "./src/shim/legacyMeshDevice.ts";
+47 -32
View File
@@ -2,6 +2,22 @@
"name": "@meshtastic/sdk",
"version": "1.0.0",
"description": "Domain-driven SDK for Meshtastic devices. Feature slices (device/chat/nodes/channels/config/telemetry/position/traceroute/files) with signals-backed reactive state. Replaces @meshtastic/core.",
"license": "GPL-3.0-only",
"repository": {
"type": "git",
"url": "git+https://github.com/meshtastic/web.git",
"directory": "packages/sdk"
},
"files": [
"package.json",
"README.md",
"LICENSE",
"dist"
],
"type": "module",
"main": "./dist/mod.js",
"module": "./dist/mod.js",
"types": "./dist/mod.d.ts",
"exports": {
".": {
"types": "./dist/mod.d.ts",
@@ -20,37 +36,6 @@
"default": "./src/core/testing/index.ts"
}
},
"type": "module",
"main": "./dist/mod.js",
"module": "./dist/mod.js",
"types": "./dist/mod.d.ts",
"license": "GPL-3.0-only",
"repository": {
"type": "git",
"url": "git+https://github.com/meshtastic/web.git",
"directory": "packages/sdk"
},
"tsdown": {
"entry": {
"mod": "mod.ts",
"transport": "src/core/transport/index.ts",
"protobuf": "src/core/protobuf/index.ts",
"testing": "src/core/testing/index.ts"
},
"platform": "browser",
"target": "esnext",
"dts": true,
"format": ["esm"],
"splitting": false,
"sourcemap": false,
"minify": false,
"treeshake": true,
"report": false,
"clean": true
},
"jsrInclude": ["mod.ts", "src", "README.md", "LICENSE"],
"jsrExclude": ["src/**/*.test.ts", "tests"],
"files": ["package.json", "README.md", "LICENSE", "dist"],
"scripts": {
"preinstall": "npx only-allow pnpm",
"prepack": "cp ../../LICENSE ./LICENSE",
@@ -69,5 +54,35 @@
"crc": "npm:crc@^4.3.2",
"ste-simple-events": "^3.0.11",
"tslog": "^4.9.3"
}
},
"tsdown": {
"clean": true,
"dts": true,
"entry": {
"mod": "mod.ts",
"protobuf": "src/core/protobuf/index.ts",
"testing": "src/core/testing/index.ts",
"transport": "src/core/transport/index.ts"
},
"format": [
"esm"
],
"minify": false,
"platform": "browser",
"report": false,
"sourcemap": false,
"splitting": false,
"target": "esnext",
"treeshake": true
},
"jsrExclude": [
"src/**/*.test.ts",
"tests"
],
"jsrInclude": [
"mod.ts",
"src",
"README.md",
"LICENSE"
]
}
+5 -2
View File
@@ -80,7 +80,8 @@ function renameEntries() {
}
// Match: import { i as Transport, n as DeviceStatusEnum, ... } from "./Foo-hash.js";
const CHUNK_IMPORT_RE = /^import\s*\{([^}]+)\}\s*from\s*"\.\/([A-Za-z]+-[A-Za-z0-9_-]+)\.js";\s*$/m;
const CHUNK_IMPORT_RE =
/^import\s*\{([^}]+)\}\s*from\s*"\.\/([A-Za-z]+-[A-Za-z0-9_-]+)\.js";\s*$/m;
function inlineChunks(entryPath) {
if (!existsSync(entryPath)) return;
@@ -150,7 +151,9 @@ function deleteOrphanedChunks() {
for (const entry of KNOWN_ENTRIES) {
const p = join(distDir, `${entry}.d.ts`);
if (!existsSync(p)) continue;
if (readFileSync(p, "utf8").includes(`./${f.replace(/\.d\.ts$/, ".js")}`)) {
if (
readFileSync(p, "utf8").includes(`./${f.replace(/\.d\.ts$/, ".js")}`)
) {
referenced = true;
break;
}
@@ -18,7 +18,14 @@ describe("MeshClient.progress", () => {
expect(client.progress.value.phase).toBe("configuring");
expect(client.progress.value).toEqual({
phase: "configuring",
received: { config: 0, modules: 0, channels: 0, nodes: 0, myInfo: false, metadata: false },
received: {
config: 0,
modules: 0,
channels: 0,
nodes: 0,
myInfo: false,
metadata: false,
},
});
});
@@ -27,14 +34,25 @@ describe("MeshClient.progress", () => {
const client = new MeshClient({ transport });
void client.configure();
client.events.onConfigPacket.dispatch(create(Protobuf.Config.ConfigSchema, {}));
client.events.onConfigPacket.dispatch(create(Protobuf.Config.ConfigSchema, {}));
client.events.onChannelPacket.dispatch(create(Protobuf.Channel.ChannelSchema, {}));
client.events.onNodeInfoPacket.dispatch(create(Protobuf.Mesh.NodeInfoSchema, {}));
client.events.onMyNodeInfo.dispatch(create(Protobuf.Mesh.MyNodeInfoSchema, {}));
client.events.onConfigPacket.dispatch(
create(Protobuf.Config.ConfigSchema, {}),
);
client.events.onConfigPacket.dispatch(
create(Protobuf.Config.ConfigSchema, {}),
);
client.events.onChannelPacket.dispatch(
create(Protobuf.Channel.ChannelSchema, {}),
);
client.events.onNodeInfoPacket.dispatch(
create(Protobuf.Mesh.NodeInfoSchema, {}),
);
client.events.onMyNodeInfo.dispatch(
create(Protobuf.Mesh.MyNodeInfoSchema, {}),
);
const cur = client.progress.value;
if (cur.phase !== "configuring") throw new Error("expected configuring phase");
if (cur.phase !== "configuring")
throw new Error("expected configuring phase");
expect(cur.received.config).toBe(2);
expect(cur.received.channels).toBe(1);
expect(cur.received.nodes).toBe(1);
@@ -46,7 +64,9 @@ describe("MeshClient.progress", () => {
const { transport } = createFakeTransport();
const client = new MeshClient({ transport });
void client.configure();
client.events.onConfigPacket.dispatch(create(Protobuf.Config.ConfigSchema, {}));
client.events.onConfigPacket.dispatch(
create(Protobuf.Config.ConfigSchema, {}),
);
client.events.onConfigComplete.dispatch(0);
const cur = client.progress.value;
@@ -59,7 +79,9 @@ describe("MeshClient.progress", () => {
const { transport } = createFakeTransport();
const client = new MeshClient({ transport });
// not calling configure() — phase is idle
client.events.onConfigPacket.dispatch(create(Protobuf.Config.ConfigSchema, {}));
client.events.onConfigPacket.dispatch(
create(Protobuf.Config.ConfigSchema, {}),
);
expect(client.progress.value.phase).toBe("idle");
});
@@ -67,12 +89,21 @@ describe("MeshClient.progress", () => {
const { transport } = createFakeTransport();
const client = new MeshClient({ transport });
void client.configure();
client.events.onConfigPacket.dispatch(create(Protobuf.Config.ConfigSchema, {}));
client.events.onConfigPacket.dispatch(
create(Protobuf.Config.ConfigSchema, {}),
);
client.events.onConfigComplete.dispatch(0);
void client.configure();
expect(client.progress.value).toEqual({
phase: "configuring",
received: { config: 0, modules: 0, channels: 0, nodes: 0, myInfo: false, metadata: false },
received: {
config: 0,
modules: 0,
channels: 0,
nodes: 0,
myInfo: false,
metadata: false,
},
});
});
});
+39 -14
View File
@@ -11,7 +11,12 @@ import { Queue } from "../queue/Queue.ts";
import { createStore, type ReadonlySignal } from "../signals/createStore.ts";
import type { Transport } from "../transport/Transport.ts";
import { DeviceStatusEnum } from "../transport/Transport.ts";
import { ChannelNumber, type Destination, Emitter, type PacketMetadata } from "../types.ts";
import {
ChannelNumber,
type Destination,
Emitter,
type PacketMetadata,
} from "../types.ts";
import { Xmodem } from "../xmodem/Xmodem.ts";
import { ChatClient } from "../../features/chat/index.ts";
import { ChannelsClient } from "../../features/channels/index.ts";
@@ -101,7 +106,9 @@ export class MeshClient {
* sends `configCompleteId`.
*/
public readonly progress: ReadonlySignal<ConnectionProgress>;
private readonly progressStore = createStore<ConnectionProgress>({ phase: "idle" });
private readonly progressStore = createStore<ConnectionProgress>({
phase: "idle",
});
private _heartbeatIntervalId: ReturnType<typeof setInterval> | undefined;
@@ -153,7 +160,10 @@ export class MeshClient {
}
public configure(): Promise<number> {
this.log.debug(Emitter[Emitter.Configure], "⚙️ Requesting device configuration");
this.log.debug(
Emitter[Emitter.Configure],
"⚙️ Requesting device configuration",
);
this.updateDeviceStatus(DeviceStatusEnum.DeviceConfiguring);
this.progressStore.write.value = {
phase: "configuring",
@@ -164,12 +174,14 @@ export class MeshClient {
payloadVariant: { case: "wantConfigId", value: this.configId },
});
return this.sendRaw(toBinary(Protobuf.Mesh.ToRadioSchema, toRadio)).catch((e) => {
if (this.device.status.value === DeviceStatusEnum.DeviceDisconnected) {
throw new Error("Device connection lost");
}
throw e;
});
return this.sendRaw(toBinary(Protobuf.Mesh.ToRadioSchema, toRadio)).catch(
(e) => {
if (this.device.status.value === DeviceStatusEnum.DeviceDisconnected) {
throw new Error("Device connection lost");
}
throw e;
},
);
}
/**
@@ -186,7 +198,10 @@ export class MeshClient {
if (cur.phase !== "configuring") return;
const next: ConnectionProgressCounters = {
...cur.received,
[field]: typeof cur.received[field] === "number" ? cur.received[field] + 1 : true,
[field]:
typeof cur.received[field] === "number"
? cur.received[field] + 1
: true,
} as ConnectionProgressCounters;
this.progressStore.write.value = { phase: "configuring", received: next };
};
@@ -200,7 +215,8 @@ export class MeshClient {
this.events.onConfigComplete.subscribe(() => {
const cur = this.progressStore.read.value;
const received = cur.phase === "configuring" ? cur.received : { ...EMPTY_COUNTERS };
const received =
cur.phase === "configuring" ? cur.received : { ...EMPTY_COUNTERS };
this.progressStore.write.value = { phase: "configured", received };
});
}
@@ -219,7 +235,10 @@ export class MeshClient {
}
this._heartbeatIntervalId = setInterval(() => {
this.heartbeat().catch((err: Error) => {
this.log.error(Emitter[Emitter.Ping], `⚠️ Unable to send heartbeat: ${err.message}`);
this.log.error(
Emitter[Emitter.Ping],
`⚠️ Unable to send heartbeat: ${err.message}`,
);
});
}, interval);
}
@@ -286,10 +305,16 @@ export class MeshClient {
meshPacket.rxTime = Math.trunc(Date.now() / 1000);
this.events.onMeshPacket.dispatch(meshPacket);
}
return await this.sendRaw(toBinary(Protobuf.Mesh.ToRadioSchema, toRadioMessage), meshPacket.id);
return await this.sendRaw(
toBinary(Protobuf.Mesh.ToRadioSchema, toRadioMessage),
meshPacket.id,
);
}
public async sendRaw(toRadio: Uint8Array, id: number = generatePacketId()): Promise<number> {
public async sendRaw(
toRadio: Uint8Array,
id: number = generatePacketId(),
): Promise<number> {
if (toRadio.length > 512) {
throw new PacketTooLargeError(toRadio.length);
}
+71 -27
View File
@@ -11,66 +11,110 @@ import type { LogEventPacket, PacketMetadata } from "../types.ts";
*/
export class EventBus {
public readonly onLogEvent = new SimpleEventDispatcher<LogEventPacket>();
public readonly onFromRadio = new SimpleEventDispatcher<Protobuf.Mesh.FromRadio>();
public readonly onMeshPacket = new SimpleEventDispatcher<Protobuf.Mesh.MeshPacket>();
public readonly onMyNodeInfo = new SimpleEventDispatcher<Protobuf.Mesh.MyNodeInfo>();
public readonly onNodeInfoPacket = new SimpleEventDispatcher<Protobuf.Mesh.NodeInfo>();
public readonly onChannelPacket = new SimpleEventDispatcher<Protobuf.Channel.Channel>();
public readonly onConfigPacket = new SimpleEventDispatcher<Protobuf.Config.Config>();
public readonly onFromRadio =
new SimpleEventDispatcher<Protobuf.Mesh.FromRadio>();
public readonly onMeshPacket =
new SimpleEventDispatcher<Protobuf.Mesh.MeshPacket>();
public readonly onMyNodeInfo =
new SimpleEventDispatcher<Protobuf.Mesh.MyNodeInfo>();
public readonly onNodeInfoPacket =
new SimpleEventDispatcher<Protobuf.Mesh.NodeInfo>();
public readonly onChannelPacket =
new SimpleEventDispatcher<Protobuf.Channel.Channel>();
public readonly onConfigPacket =
new SimpleEventDispatcher<Protobuf.Config.Config>();
public readonly onModuleConfigPacket =
new SimpleEventDispatcher<Protobuf.ModuleConfig.ModuleConfig>();
public readonly onAtakPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onMessagePacket = new SimpleEventDispatcher<PacketMetadata<string>>();
public readonly onAtakPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onMessagePacket = new SimpleEventDispatcher<
PacketMetadata<string>
>();
public readonly onRemoteHardwarePacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.RemoteHardware.HardwareMessage>
>();
public readonly onPositionPacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.Mesh.Position>
>();
public readonly onUserPacket = new SimpleEventDispatcher<PacketMetadata<Protobuf.Mesh.User>>();
public readonly onUserPacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.Mesh.User>
>();
public readonly onRoutingPacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.Mesh.Routing>
>();
public readonly onDeviceMetadataPacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.Mesh.DeviceMetadata>
>();
public readonly onCannedMessageModulePacket = new SimpleEventDispatcher<PacketMetadata<string>>();
public readonly onCannedMessageModulePacket = new SimpleEventDispatcher<
PacketMetadata<string>
>();
public readonly onWaypointPacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.Mesh.Waypoint>
>();
public readonly onAudioPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onDetectionSensorPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onPingPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onIpTunnelPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onAudioPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onDetectionSensorPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onPingPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onIpTunnelPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onPaxcounterPacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.PaxCount.Paxcount>
>();
public readonly onSerialPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onStoreForwardPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onRangeTestPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onSerialPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onStoreForwardPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onRangeTestPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onTelemetryPacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.Telemetry.Telemetry>
>();
public readonly onZpsPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onSimulatorPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onZpsPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onSimulatorPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onTraceRoutePacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.Mesh.RouteDiscovery>
>();
public readonly onNeighborInfoPacket = new SimpleEventDispatcher<
PacketMetadata<Protobuf.Mesh.NeighborInfo>
>();
public readonly onAtakPluginPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onMapReportPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onPrivatePacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onAtakForwarderPacket = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
public readonly onAtakPluginPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onMapReportPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onPrivatePacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onAtakForwarderPacket = new SimpleEventDispatcher<
PacketMetadata<Uint8Array>
>();
public readonly onClientNotificationPacket =
new SimpleEventDispatcher<Protobuf.Mesh.ClientNotification>();
public readonly onDeviceStatus = new SimpleEventDispatcher<DeviceStatusEnum>();
public readonly onLogRecord = new SimpleEventDispatcher<Protobuf.Mesh.LogRecord>();
public readonly onDeviceStatus =
new SimpleEventDispatcher<DeviceStatusEnum>();
public readonly onLogRecord =
new SimpleEventDispatcher<Protobuf.Mesh.LogRecord>();
public readonly onMeshHeartbeat = new SimpleEventDispatcher<Date>();
public readonly onDeviceDebugLog = new SimpleEventDispatcher<Uint8Array>();
public readonly onPendingSettingsChange = new SimpleEventDispatcher<boolean>();
public readonly onQueueStatus = new SimpleEventDispatcher<Protobuf.Mesh.QueueStatus>();
public readonly onPendingSettingsChange =
new SimpleEventDispatcher<boolean>();
public readonly onQueueStatus =
new SimpleEventDispatcher<Protobuf.Mesh.QueueStatus>();
public readonly onConfigComplete = new SimpleEventDispatcher<number>();
public readonly onRebooted = new SimpleEventDispatcher<void>();
}
+13 -6
View File
@@ -1,6 +1,7 @@
import { Logger } from "tslog";
const prettyLogTemplate = "{{hh}}:{{MM}}:{{ss}}:{{ms}}\t{{logLevelName}}\t[{{name}}]\t";
const prettyLogTemplate =
"{{hh}}:{{MM}}:{{ss}}:{{ms}}\t{{logLevelName}}\t[{{name}}]\t";
/**
* Minimum log level. Default is `info` (3). Lifts to `debug` (2) when:
@@ -15,8 +16,11 @@ function defaultMinLevel(): number {
try {
if (
typeof globalThis !== "undefined" &&
typeof (globalThis as { localStorage?: Storage }).localStorage !== "undefined" &&
(globalThis as { localStorage: Storage }).localStorage.getItem("mesh-debug") === "1"
typeof (globalThis as { localStorage?: Storage }).localStorage !==
"undefined" &&
(globalThis as { localStorage: Storage }).localStorage.getItem(
"mesh-debug",
) === "1"
) {
return 2;
}
@@ -25,15 +29,18 @@ function defaultMinLevel(): number {
}
if (
typeof globalThis !== "undefined" &&
(globalThis as { process?: { env?: Record<string, string | undefined> } }).process?.env
?.MESH_DEBUG === "1"
(globalThis as { process?: { env?: Record<string, string | undefined> } })
.process?.env?.MESH_DEBUG === "1"
) {
return 2;
}
return 3;
}
export function createLogger(name: string, options?: { minLevel?: number }): Logger<unknown> {
export function createLogger(
name: string,
options?: { minLevel?: number },
): Logger<unknown> {
return new Logger({
name,
prettyLogTemplate,
@@ -48,9 +48,16 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
case "packet": {
let decodedMessage: Protobuf.Mesh.FromRadio;
try {
decodedMessage = fromBinary(Protobuf.Mesh.FromRadioSchema, chunk.data);
decodedMessage = fromBinary(
Protobuf.Mesh.FromRadioSchema,
chunk.data,
);
} catch (e) {
sink.log.error(Emitter[Emitter.HandleFromRadio], "⚠️ Received undecodable packet", e);
sink.log.error(
Emitter[Emitter.HandleFromRadio],
"⚠️ Received undecodable packet",
e,
);
break;
}
sink.events.onFromRadio.dispatch(decodedMessage);
@@ -69,7 +76,9 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
break;
}
case "myInfo": {
sink.events.onMyNodeInfo.dispatch(decodedMessage.payloadVariant.value);
sink.events.onMyNodeInfo.dispatch(
decodedMessage.payloadVariant.value,
);
sink.log.info(
Emitter[Emitter.HandleFromRadio],
"📱 Received Node info for this device",
@@ -81,7 +90,9 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
Emitter[Emitter.HandleFromRadio],
`📱 Received Node Info packet for node: ${decodedMessage.payloadVariant.value.num}`,
);
sink.events.onNodeInfoPacket.dispatch(decodedMessage.payloadVariant.value);
sink.events.onNodeInfoPacket.dispatch(
decodedMessage.payloadVariant.value,
);
if (decodedMessage.payloadVariant.value.position) {
sink.events.onPositionPacket.dispatch({
@@ -120,12 +131,19 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
"⚠️ Received Config packet of variant: UNK",
);
}
sink.events.onConfigPacket.dispatch(decodedMessage.payloadVariant.value);
sink.events.onConfigPacket.dispatch(
decodedMessage.payloadVariant.value,
);
break;
}
case "logRecord": {
sink.log.trace(Emitter[Emitter.HandleFromRadio], "Received onLogRecord");
sink.events.onLogRecord.dispatch(decodedMessage.payloadVariant.value);
sink.log.trace(
Emitter[Emitter.HandleFromRadio],
"Received onLogRecord",
);
sink.events.onLogRecord.dispatch(
decodedMessage.payloadVariant.value,
);
break;
}
case "configCompleteId": {
@@ -133,7 +151,9 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
Emitter[Emitter.HandleFromRadio],
`⚙️ Received config complete id: ${decodedMessage.payloadVariant.value}`,
);
sink.events.onConfigComplete.dispatch(decodedMessage.payloadVariant.value);
sink.events.onConfigComplete.dispatch(
decodedMessage.payloadVariant.value,
);
if (decodedMessage.payloadVariant.value === sink.configId) {
sink.log.info(
Emitter[Emitter.HandleFromRadio],
@@ -162,7 +182,9 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
"⚠️ Received Module Config packet of variant: UNK",
);
}
sink.events.onModuleConfigPacket.dispatch(decodedMessage.payloadVariant.value);
sink.events.onModuleConfigPacket.dispatch(
decodedMessage.payloadVariant.value,
);
break;
}
case "channel": {
@@ -170,7 +192,9 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
Emitter[Emitter.HandleFromRadio],
`🔐 Received Channel: ${decodedMessage.payloadVariant.value.index}`,
);
sink.events.onChannelPacket.dispatch(decodedMessage.payloadVariant.value);
sink.events.onChannelPacket.dispatch(
decodedMessage.payloadVariant.value,
);
break;
}
case "queueStatus": {
@@ -178,7 +202,9 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
Emitter[Emitter.HandleFromRadio],
`🚧 Received Queue Status: ${decodedMessage.payloadVariant.value}`,
);
sink.events.onQueueStatus.dispatch(decodedMessage.payloadVariant.value);
sink.events.onQueueStatus.dispatch(
decodedMessage.payloadVariant.value,
);
break;
}
case "xmodemPacket": {
@@ -187,15 +213,19 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
}
case "metadata": {
if (
Number.parseFloat(decodedMessage.payloadVariant.value.firmwareVersion) <
Constants.minFwVer
Number.parseFloat(
decodedMessage.payloadVariant.value.firmwareVersion,
) < Constants.minFwVer
) {
sink.log.fatal(
Emitter[Emitter.HandleFromRadio],
`Device firmware outdated. Min supported: ${Constants.minFwVer} got: ${decodedMessage.payloadVariant.value.firmwareVersion}`,
);
}
sink.log.debug(Emitter[Emitter.GetMetadata], "🏷️ Received metadata packet");
sink.log.debug(
Emitter[Emitter.GetMetadata],
"🏷️ Received metadata packet",
);
sink.events.onDeviceMetadataPacket.dispatch({
id: decodedMessage.id,
rxTime: new Date(),
@@ -215,7 +245,9 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
Emitter[Emitter.HandleFromRadio],
`📣 Received ClientNotification: ${decodedMessage.payloadVariant.value.message}`,
);
sink.events.onClientNotificationPacket.dispatch(decodedMessage.payloadVariant.value);
sink.events.onClientNotificationPacket.dispatch(
decodedMessage.payloadVariant.value,
);
break;
}
default: {
@@ -230,7 +262,10 @@ export const decodePacket = (sink: PacketSink): WritableStream<DeviceOutput> =>
},
});
function handleMeshPacket(sink: PacketSink, meshPacket: Protobuf.Mesh.MeshPacket): void {
function handleMeshPacket(
sink: PacketSink,
meshPacket: Protobuf.Mesh.MeshPacket,
): void {
sink.events.onMeshPacket.dispatch(meshPacket);
if (meshPacket.from !== sink.myNodeNum) {
sink.events.onMeshHeartbeat.dispatch(new Date());
@@ -283,7 +318,10 @@ function handleDecodedPacket(
case Protobuf.Portnums.PortNum.REMOTE_HARDWARE_APP: {
sink.events.onRemoteHardwarePacket.dispatch({
...packetMetadata,
data: fromBinary(Protobuf.RemoteHardware.HardwareMessageSchema, dataPacket.payload),
data: fromBinary(
Protobuf.RemoteHardware.HardwareMessageSchema,
dataPacket.payload,
),
});
break;
}
@@ -302,11 +340,19 @@ function handleDecodedPacket(
break;
}
case Protobuf.Portnums.PortNum.ROUTING_APP: {
const routingPacket = fromBinary(Protobuf.Mesh.RoutingSchema, dataPacket.payload);
sink.events.onRoutingPacket.dispatch({ ...packetMetadata, data: routingPacket });
const routingPacket = fromBinary(
Protobuf.Mesh.RoutingSchema,
dataPacket.payload,
);
sink.events.onRoutingPacket.dispatch({
...packetMetadata,
data: routingPacket,
});
switch (routingPacket.variant.case) {
case "errorReason": {
if (routingPacket.variant.value === Protobuf.Mesh.Routing_Error.NONE) {
if (
routingPacket.variant.value === Protobuf.Mesh.Routing_Error.NONE
) {
sink.queue.processAck(dataPacket.requestId);
} else {
sink.queue.processError({
@@ -325,10 +371,15 @@ function handleDecodedPacket(
break;
}
case Protobuf.Portnums.PortNum.ADMIN_APP: {
const adminMessage = fromBinary(Protobuf.Admin.AdminMessageSchema, dataPacket.payload);
const adminMessage = fromBinary(
Protobuf.Admin.AdminMessageSchema,
dataPacket.payload,
);
switch (adminMessage.payloadVariant.case) {
case "getChannelResponse": {
sink.events.onChannelPacket.dispatch(adminMessage.payloadVariant.value);
sink.events.onChannelPacket.dispatch(
adminMessage.payloadVariant.value,
);
break;
}
case "getOwnerResponse": {
@@ -339,11 +390,15 @@ function handleDecodedPacket(
break;
}
case "getConfigResponse": {
sink.events.onConfigPacket.dispatch(adminMessage.payloadVariant.value);
sink.events.onConfigPacket.dispatch(
adminMessage.payloadVariant.value,
);
break;
}
case "getModuleConfigResponse": {
sink.events.onModuleConfigPacket.dispatch(adminMessage.payloadVariant.value);
sink.events.onModuleConfigPacket.dispatch(
adminMessage.payloadVariant.value,
);
break;
}
case "getDeviceMetadataResponse": {
@@ -388,7 +443,10 @@ function handleDecodedPacket(
break;
}
case Protobuf.Portnums.PortNum.AUDIO_APP: {
sink.events.onAudioPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onAudioPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.DETECTION_SENSOR_APP: {
@@ -399,11 +457,17 @@ function handleDecodedPacket(
break;
}
case Protobuf.Portnums.PortNum.REPLY_APP: {
sink.events.onPingPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onPingPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.IP_TUNNEL_APP: {
sink.events.onIpTunnelPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onIpTunnelPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.PAXCOUNTER_APP: {
@@ -414,36 +478,57 @@ function handleDecodedPacket(
break;
}
case Protobuf.Portnums.PortNum.SERIAL_APP: {
sink.events.onSerialPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onSerialPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.STORE_FORWARD_APP: {
sink.events.onStoreForwardPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onStoreForwardPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.RANGE_TEST_APP: {
sink.events.onRangeTestPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onRangeTestPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.TELEMETRY_APP: {
sink.events.onTelemetryPacket.dispatch({
...packetMetadata,
data: fromBinary(Protobuf.Telemetry.TelemetrySchema, dataPacket.payload),
data: fromBinary(
Protobuf.Telemetry.TelemetrySchema,
dataPacket.payload,
),
});
break;
}
case Protobuf.Portnums.PortNum.ZPS_APP: {
sink.events.onZpsPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onZpsPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.SIMULATOR_APP: {
sink.events.onSimulatorPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onSimulatorPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.TRACEROUTE_APP: {
sink.events.onTraceRoutePacket.dispatch({
...packetMetadata,
data: fromBinary(Protobuf.Mesh.RouteDiscoverySchema, dataPacket.payload),
data: fromBinary(
Protobuf.Mesh.RouteDiscoverySchema,
dataPacket.payload,
),
});
break;
}
@@ -455,19 +540,31 @@ function handleDecodedPacket(
break;
}
case Protobuf.Portnums.PortNum.ATAK_PLUGIN: {
sink.events.onAtakPluginPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onAtakPluginPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.MAP_REPORT_APP: {
sink.events.onMapReportPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onMapReportPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.PRIVATE_APP: {
sink.events.onPrivatePacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onPrivatePacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
case Protobuf.Portnums.PortNum.ATAK_FORWARDER: {
sink.events.onAtakForwarderPacket.dispatch({ ...packetMetadata, data: dataPacket.payload });
sink.events.onAtakForwarderPacket.dispatch({
...packetMetadata,
data: dataPacket.payload,
});
break;
}
default:
@@ -4,7 +4,10 @@ import type { DeviceOutput } from "../transport/Transport.ts";
* Transforms a raw byte stream from the device into typed DeviceOutput chunks
* by parsing the 0x94 0xC3 framing header and length prefix.
*/
export const fromDeviceStream: () => TransformStream<Uint8Array, DeviceOutput> = () => {
export const fromDeviceStream: () => TransformStream<
Uint8Array,
DeviceOutput
> = () => {
let byteBuffer = new Uint8Array([]);
const textDecoder = new TextDecoder();
return new TransformStream<Uint8Array, DeviceOutput>({
@@ -26,11 +29,18 @@ export const fromDeviceStream: () => TransformStream<Uint8Array, DeviceOutput> =
const msb = byteBuffer[2];
const lsb = byteBuffer[3];
if (msb !== undefined && lsb !== undefined && byteBuffer.length >= 4 + (msb << 8) + lsb) {
if (
msb !== undefined &&
lsb !== undefined &&
byteBuffer.length >= 4 + (msb << 8) + lsb
) {
const packet = byteBuffer.subarray(4, 4 + (msb << 8) + lsb);
const malformedDetectorIndex = packet.indexOf(0x94);
if (malformedDetectorIndex !== -1 && packet[malformedDetectorIndex + 1] === 0xc3) {
if (
malformedDetectorIndex !== -1 &&
packet[malformedDetectorIndex + 1] === 0xc3
) {
console.warn(
`⚠️ Malformed packet found, discarding: ${byteBuffer
.subarray(0, malformedDetectorIndex - 1)
+10 -2
View File
@@ -1,11 +1,19 @@
/**
* Pads outbound packets with the 0x94 0xC3 framing header and length prefix.
*/
export const toDeviceStream: () => TransformStream<Uint8Array, Uint8Array> = () => {
export const toDeviceStream: () => TransformStream<
Uint8Array,
Uint8Array
> = () => {
return new TransformStream<Uint8Array, Uint8Array>({
transform(chunk: Uint8Array, controller): void {
const bufLen = chunk.length;
const header = new Uint8Array([0x94, 0xc3, (bufLen >> 8) & 0xff, bufLen & 0xff]);
const header = new Uint8Array([
0x94,
0xc3,
(bufLen >> 8) & 0xff,
bufLen & 0xff,
]);
controller.enqueue(new Uint8Array([...header, ...chunk]));
},
});
+13 -4
View File
@@ -53,7 +53,9 @@ export class Queue {
return;
}
console.warn(`Packet ${item.id} of type ${decoded.payloadVariant.case} timed out`);
console.warn(
`Packet ${item.id} of type ${decoded.payloadVariant.case} timed out`,
);
reject({
id: item.id,
@@ -79,7 +81,9 @@ export class Queue {
}
public processError(e: PacketError): void {
console.error(`Error received for packet ${e.id}: ${Protobuf.Mesh.Routing_Error[e.error]}`);
console.error(
`Error received for packet ${e.id}: ${Protobuf.Mesh.Routing_Error[e.error]}`,
);
this.errorNotifier.dispatch(e);
}
@@ -91,7 +95,9 @@ export class Queue {
return queueItem.promise;
}
public async processQueue(outputStream: WritableStream<Uint8Array>): Promise<void> {
public async processQueue(
outputStream: WritableStream<Uint8Array>,
): Promise<void> {
if (this.lock) {
return;
}
@@ -109,7 +115,10 @@ export class Queue {
item.sent = true;
} catch (error) {
const err = error as { code?: string };
if (err?.code === "ECONNRESET" || err?.code === "ERR_INVALID_STATE") {
if (
err?.code === "ECONNRESET" ||
err?.code === "ERR_INVALID_STATE"
) {
writer.releaseLock();
this.lock = false;
throw error;
@@ -41,7 +41,9 @@ describe("MeshRegistry", () => {
const reg = new MeshRegistry();
const { transport } = createFakeTransport();
reg.create(1, { transport });
expect(() => reg.create(1, { transport: createFakeTransport().transport })).toThrow();
expect(() =>
reg.create(1, { transport: createFakeTransport().transport }),
).toThrow();
});
it("remove disconnects the client and falls back to another active id", async () => {
@@ -113,9 +113,11 @@ export class MeshRegistry {
}
private snapshot(): void {
this.backing.value = Array.from(this.clients.entries()).map(([id, client]) => ({
id,
client,
}));
this.backing.value = Array.from(this.clients.entries()).map(
([id, client]) => ({
id,
client,
}),
);
}
}
+7 -2
View File
@@ -21,12 +21,17 @@ export interface ReadonlySignal<T> {
* is invoked asynchronously on every change. Returns both the writable signal
* (for slice-internal use) and the readable facade (for external consumption).
*/
export function createStore<T>(initial: T): { write: Signal<T>; read: ReadonlySignal<T> } {
export function createStore<T>(initial: T): {
write: Signal<T>;
read: ReadonlySignal<T>;
} {
const write = signal(initial);
return { write, read: toReadonly(write) };
}
export function toReadonly<T>(s: Signal<T> | PreactReadonlySignal<T>): ReadonlySignal<T> {
export function toReadonly<T>(
s: Signal<T> | PreactReadonlySignal<T>,
): ReadonlySignal<T> {
return {
get value() {
return s.value;
@@ -13,8 +13,12 @@ export interface FakeTransportHandle {
close(): Promise<void>;
}
type MyNodeInfoInit = Parameters<typeof create<typeof Protobuf.Mesh.MyNodeInfoSchema>>[1];
type NodeInfoInit = Parameters<typeof create<typeof Protobuf.Mesh.NodeInfoSchema>>[1];
type MyNodeInfoInit = Parameters<
typeof create<typeof Protobuf.Mesh.MyNodeInfoSchema>
>[1];
type NodeInfoInit = Parameters<
typeof create<typeof Protobuf.Mesh.NodeInfoSchema>
>[1];
export interface FakeResponder {
withMyNodeInfo(info: MyNodeInfoInit & { myNodeNum: number }): void;
@@ -31,7 +35,9 @@ export interface FakeResponder {
*/
export function createFakeTransport(): FakeTransportHandle {
const sent: Uint8Array[] = [];
let fromDeviceController: ReadableStreamDefaultController<Uint8Array> | undefined;
let fromDeviceController:
| ReadableStreamDefaultController<Uint8Array>
| undefined;
const fromDeviceRaw = new ReadableStream<Uint8Array>({
start(controller) {
+4 -1
View File
@@ -1,2 +1,5 @@
export { createFakeTransport } from "./createFakeTransport.ts";
export type { FakeResponder, FakeTransportHandle } from "./createFakeTransport.ts";
export type {
FakeResponder,
FakeTransportHandle,
} from "./createFakeTransport.ts";
+5 -1
View File
@@ -1,6 +1,10 @@
import type * as Protobuf from "@meshtastic/protobufs";
export type { DeviceOutput, HttpRetryConfig, Transport } from "./transport/Transport.ts";
export type {
DeviceOutput,
HttpRetryConfig,
Transport,
} from "./transport/Transport.ts";
export { DeviceStatusEnum } from "./transport/Transport.ts";
export interface QueueItem {
+5 -1
View File
@@ -72,7 +72,11 @@ export class Xmodem {
this.rxBuffer[this.counter] = packet.buffer;
return this.sendCommand(Protobuf.Xmodem.XModem_Control.ACK);
}
return await this.sendCommand(Protobuf.Xmodem.XModem_Control.NAK, undefined, packet.seq);
return await this.sendCommand(
Protobuf.Xmodem.XModem_Control.NAK,
undefined,
packet.seq,
);
}
case Protobuf.Xmodem.XModem_Control.STX: {
break;
@@ -23,7 +23,11 @@ describe("ChannelsClient", () => {
);
expect(client.channels.list.value.length).toBe(2);
expect(client.channels.get(0)?.role).toBe(Protobuf.Channel.Channel_Role.PRIMARY);
expect(client.channels.get(1)?.role).toBe(Protobuf.Channel.Channel_Role.SECONDARY);
expect(client.channels.get(0)?.role).toBe(
Protobuf.Channel.Channel_Role.PRIMARY,
);
expect(client.channels.get(1)?.role).toBe(
Protobuf.Channel.Channel_Role.SECONDARY,
);
});
});
@@ -2,7 +2,11 @@ import type * as Protobuf from "@meshtastic/protobufs";
import type { ResultType } from "better-result";
import type { MeshClient } from "../../core/client/MeshClient.ts";
import type { ReadonlySignal } from "../../core/signals/createStore.ts";
import { clearChannel, getChannel, setChannel } from "./application/ChannelUseCases.ts";
import {
clearChannel,
getChannel,
setChannel,
} from "./application/ChannelUseCases.ts";
import type { Channel } from "./domain/Channel.ts";
import { ChannelMapper } from "./infrastructure/ChannelMapper.ts";
import { ChannelsStore } from "./state/channelsStore.ts";
@@ -26,7 +30,9 @@ export class ChannelsClient {
return this.store.get(index);
}
public set(channel: Protobuf.Channel.Channel): Promise<ResultType<number, Error>> {
public set(
channel: Protobuf.Channel.Channel,
): Promise<ResultType<number, Error>> {
return setChannel(this.client, channel);
}
@@ -10,7 +10,10 @@ export async function setChannel(
channel: Protobuf.Channel.Channel,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "setChannel", value: channel });
const id = await sendAdminMessage(client, {
case: "setChannel",
value: channel,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -22,7 +25,10 @@ export async function getChannel(
index: number,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "getChannelRequest", value: index + 1 });
const id = await sendAdminMessage(client, {
case: "getChannelRequest",
value: index + 1,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -36,10 +36,16 @@ describe("ChatClient drafts", () => {
const { transport } = createFakeTransport();
const client = new MeshClient({ transport, chat: { draftRepository } });
client.chat.drafts.set({ kind: "channel", channel: ChannelNumber.Primary }, "hello");
await new Promise((r) => setTimeout(r, 5));
expect(await draftRepository.load({ kind: "channel", channel: ChannelNumber.Primary })).toBe(
client.chat.drafts.set(
{ kind: "channel", channel: ChannelNumber.Primary },
"hello",
);
await new Promise((r) => setTimeout(r, 5));
expect(
await draftRepository.load({
kind: "channel",
channel: ChannelNumber.Primary,
}),
).toBe("hello");
});
});
@@ -22,7 +22,10 @@ function seedMessage(id: number, ms: number, text: string): Message {
describe("ChatClient persistence", () => {
it("hydrates messages from the repository on first subscription", async () => {
const repository = new InMemoryMessageRepository();
await repository.appendBatch([seedMessage(1, 1000, "first"), seedMessage(2, 2000, "second")]);
await repository.appendBatch([
seedMessage(1, 1000, "first"),
seedMessage(2, 2000, "second"),
]);
const { transport } = createFakeTransport();
const client = new MeshClient({
@@ -60,7 +63,11 @@ describe("ChatClient persistence", () => {
new Date(3000),
50,
);
expect(sig.value.map((m) => m.text)).toEqual(["oldest", "middle", "newest"]);
expect(sig.value.map((m) => m.text)).toEqual([
"oldest",
"middle",
"newest",
]);
});
it("persists inbound messages through the repository", async () => {
@@ -14,9 +14,13 @@ type SendResult = ResultType<number, SendTextError>;
* for us, so this helper polls the queue and acks any pending items
* until the send promise settles.
*/
async function withAckFlush<T>(client: MeshClient, run: () => Promise<T>): Promise<T> {
async function withAckFlush<T>(
client: MeshClient,
run: () => Promise<T>,
): Promise<T> {
const flush = setInterval(() => {
for (const item of client.queue.getState()) client.queue.processAck(item.id);
for (const item of client.queue.getState())
client.queue.processAck(item.id);
}, 5);
try {
return await run();
@@ -112,7 +116,8 @@ describe("ChatClient.send optimistic append", () => {
expect(messages.value[0]!.state).toBe(MessageState.Pending);
// Now drain the queue so the test doesn't dangle.
for (const item of client.queue.getState()) client.queue.processAck(item.id);
for (const item of client.queue.getState())
client.queue.processAck(item.id);
await sendPromise;
});
+40 -11
View File
@@ -3,7 +3,10 @@ import type { ResultType } from "better-result";
import type { MeshClient } from "../../core/client/MeshClient.ts";
import { Constants } from "../../core/constants/index.ts";
import { generatePacketId } from "../../core/identifiers/PacketId.ts";
import { type ReadonlySignal, toReadonly } from "../../core/signals/createStore.ts";
import {
type ReadonlySignal,
toReadonly,
} from "../../core/signals/createStore.ts";
import type { ChannelNumber } from "../../core/types.ts";
import type { DraftRepository } from "./domain/DraftRepository.ts";
import type { Message } from "./domain/Message.ts";
@@ -16,7 +19,11 @@ import { MessageState } from "./domain/MessageState.ts";
import { MessageMapper } from "./infrastructure/MessageMapper.ts";
import { InMemoryDraftRepository } from "./infrastructure/repositories/InMemoryDraftRepository.ts";
import { InMemoryMessageRepository } from "./infrastructure/repositories/InMemoryMessageRepository.ts";
import { type SendTextError, type SendTextInput, sendText } from "./application/SendTextUseCase.ts";
import {
type SendTextError,
type SendTextInput,
sendText,
} from "./application/SendTextUseCase.ts";
import { ChatStore } from "./state/chatStore.ts";
import { DraftStore } from "./state/draftStore.ts";
@@ -80,7 +87,8 @@ export class ChatClient {
this.store = new ChatStore();
this.draftStore = new DraftStore();
this.repository = options.repository ?? new InMemoryMessageRepository();
this.draftRepository = options.draftRepository ?? new InMemoryDraftRepository();
this.draftRepository =
options.draftRepository ?? new InMemoryDraftRepository();
this.retention = options.retention;
this.initialLoadLimit = options.initialLoadLimit ?? 50;
@@ -101,7 +109,10 @@ export class ChatClient {
const message = MessageMapper.fromPacket(packet);
const conv: ConversationKey =
packet.type === "direct" && packet.to !== Constants.broadcastNum
? { kind: "direct", peer: packet.from === client.myNodeNum ? packet.to : packet.from }
? {
kind: "direct",
peer: packet.from === client.myNodeNum ? packet.to : packet.from,
}
: { kind: "channel", channel: packet.channel };
const key = this.keyFor(conv);
@@ -121,7 +132,10 @@ export class ChatClient {
client.events.onRoutingPacket.subscribe((packet) => {
if (packet.data.variant.case === "errorReason") {
const state = packet.data.variant.value === 0 ? MessageState.Ack : MessageState.Failed;
const state =
packet.data.variant.value === 0
? MessageState.Ack
: MessageState.Failed;
this.store.updateState(packet.id, state);
void this.repository.updateState(packet.id, state).catch(() => {});
}
@@ -138,10 +152,15 @@ export class ChatClient {
return this.store.messagesForDirect(peer);
}
public async loadOlder(conv: ConversationKey, before: Date, limit = 50): Promise<Message[]> {
public async loadOlder(
conv: ConversationKey,
before: Date,
limit = 50,
): Promise<Message[]> {
const older = await this.repository.loadBefore(conv, before, limit);
const key = this.keyFor(conv);
for (let i = older.length - 1; i >= 0; i--) this.store.prepend(key, older[i]!);
for (let i = older.length - 1; i >= 0; i--)
this.store.prepend(key, older[i]!);
return older;
}
@@ -174,7 +193,9 @@ export class ChatClient {
}
}
public async send(input: SendTextInput): Promise<ResultType<number, SendTextError>> {
public async send(
input: SendTextInput,
): Promise<ResultType<number, SendTextError>> {
const conv: ConversationKey =
typeof input.destination === "number"
? { kind: "direct", peer: input.destination }
@@ -196,7 +217,10 @@ export class ChatClient {
const message: Message = {
id: packetId,
from: this.client.myNodeNum,
to: typeof input.destination === "number" ? input.destination : Constants.broadcastNum,
to:
typeof input.destination === "number"
? input.destination
: Constants.broadcastNum,
channel: input.channel ?? 0,
rxTime: new Date(),
type: typeof input.destination === "number" ? "direct" : "broadcast",
@@ -216,7 +240,9 @@ export class ChatClient {
// the optimistic message visible but mark it Failed so the user
// sees the error state next to their bubble.
this.store.updateState(packetId, MessageState.Failed);
void this.repository.updateState(packetId, MessageState.Failed).catch(() => {});
void this.repository
.updateState(packetId, MessageState.Failed)
.catch(() => {});
}
return result;
}
@@ -227,7 +253,10 @@ export class ChatClient {
this.hydrated.add(key);
void (async () => {
try {
const recent = await this.repository.loadRecent(conv, this.initialLoadLimit);
const recent = await this.repository.loadRecent(
conv,
this.initialLoadLimit,
);
for (const m of recent) this.store.append(key, m);
} catch {
// adapter may not have history yet; safe to ignore
@@ -41,7 +41,10 @@ describe("ChatClient unread counts", () => {
const client = new MeshClient({ transport });
expect(client.chat.unread.total.value).toBe(0);
expect(
client.chat.unread.count({ kind: "channel", channel: ChannelNumber.Primary }).value,
client.chat.unread.count({
kind: "channel",
channel: ChannelNumber.Primary,
}).value,
).toBe(0);
});
@@ -54,11 +57,17 @@ describe("ChatClient unread counts", () => {
dispatchInbound(client, { from: 200, ms: 2000 });
expect(
client.chat.unread.count({ kind: "channel", channel: ChannelNumber.Primary }).value,
client.chat.unread.count({
kind: "channel",
channel: ChannelNumber.Primary,
}).value,
).toBe(2);
expect(client.chat.unread.total.value).toBe(2);
client.chat.unread.markRead({ kind: "channel", channel: ChannelNumber.Primary });
client.chat.unread.markRead({
kind: "channel",
channel: ChannelNumber.Primary,
});
expect(client.chat.unread.total.value).toBe(0);
});
@@ -76,10 +85,20 @@ describe("ChatClient unread counts", () => {
const client = new MeshClient({ transport });
setMyNode(client);
dispatchInbound(client, { from: 50, to: MY_NODE, type: "direct", ms: 1000 });
expect(client.chat.unread.count({ kind: "direct", peer: 50 }).value).toBe(1);
dispatchInbound(client, {
from: 50,
to: MY_NODE,
type: "direct",
ms: 1000,
});
expect(client.chat.unread.count({ kind: "direct", peer: 50 }).value).toBe(
1,
);
expect(
client.chat.unread.count({ kind: "channel", channel: ChannelNumber.Primary }).value,
client.chat.unread.count({
kind: "channel",
channel: ChannelNumber.Primary,
}).value,
).toBe(0);
});
});
@@ -1,6 +1,10 @@
import { Result } from "better-result";
import { describe, expect, it, vi } from "vitest";
import { EmptyMessageError, MessageTooLongError, sendText } from "./SendTextUseCase.ts";
import {
EmptyMessageError,
MessageTooLongError,
sendText,
} from "./SendTextUseCase.ts";
import type { MeshClient } from "../../../core/client/MeshClient.ts";
function makeClient(sendPacket = vi.fn().mockResolvedValue(123)) {
@@ -2,7 +2,11 @@ import * as Protobuf from "@meshtastic/protobufs";
import { Result } from "better-result";
import type { ResultType } from "better-result";
import type { MeshClient } from "../../../core/client/MeshClient.ts";
import { ChannelNumber, type Destination, Emitter } from "../../../core/types.ts";
import {
ChannelNumber,
type Destination,
Emitter,
} from "../../../core/types.ts";
export interface SendTextInput {
text: string;
@@ -30,7 +34,9 @@ export class EmptyMessageError extends Error {
export class MessageTooLongError extends Error {
readonly _tag = "MessageTooLongError";
constructor(readonly byteLength: number) {
super(`Message text encodes to ${byteLength} bytes; max payload is ~228 bytes`);
super(
`Message text encodes to ${byteLength} bytes; max payload is ~228 bytes`,
);
this.name = "MessageTooLongError";
}
}
@@ -4,7 +4,11 @@ import { Result } from "better-result";
import type { ResultType } from "better-result";
import type { MeshClient } from "../../../core/client/MeshClient.ts";
import { generatePacketId } from "../../../core/identifiers/PacketId.ts";
import { ChannelNumber, type Destination, Emitter } from "../../../core/types.ts";
import {
ChannelNumber,
type Destination,
Emitter,
} from "../../../core/types.ts";
export async function sendWaypoint(
client: MeshClient,
@@ -26,7 +26,11 @@ export interface RetentionPolicy {
*/
export interface MessageRepository {
loadRecent(key: ConversationKey, limit: number): Promise<Message[]>;
loadBefore(key: ConversationKey, cursor: Date, limit: number): Promise<Message[]>;
loadBefore(
key: ConversationKey,
cursor: Date,
limit: number,
): Promise<Message[]>;
append(message: Message): Promise<void>;
appendBatch(messages: ReadonlyArray<Message>): Promise<void>;
updateState(id: number, state: MessageState): Promise<void>;
@@ -38,5 +42,7 @@ export interface MessageRepository {
}
export function conversationKeyString(key: ConversationKey): string {
return key.kind === "channel" ? `channel:${key.channel}` : `direct:${key.peer}`;
return key.kind === "channel"
? `channel:${key.channel}`
: `direct:${key.peer}`;
}
+5 -1
View File
@@ -1,5 +1,9 @@
export { ChatClient } from "./ChatClient.ts";
export type { ChatClientOptions, ChatDrafts, ChatUnread } from "./ChatClient.ts";
export type {
ChatClientOptions,
ChatDrafts,
ChatUnread,
} from "./ChatClient.ts";
export type { DraftRepository } from "./domain/DraftRepository.ts";
export { InMemoryDraftRepository } from "./infrastructure/repositories/InMemoryDraftRepository.ts";
export type { Message } from "./domain/Message.ts";
@@ -3,7 +3,10 @@ import type { Message } from "../domain/Message.ts";
import { MessageState } from "../domain/MessageState.ts";
export const MessageMapper = {
fromPacket(packet: PacketMetadata<string>, state: MessageState = MessageState.Ack): Message {
fromPacket(
packet: PacketMetadata<string>,
state: MessageState = MessageState.Ack,
): Message {
return {
id: packet.id,
from: packet.from,
@@ -1,8 +1,14 @@
import { type ConversationKey, conversationKeyString } from "../../domain/MessageRepository.ts";
import {
type ConversationKey,
conversationKeyString,
} from "../../domain/MessageRepository.ts";
import type { DraftRepository } from "../../domain/DraftRepository.ts";
export class InMemoryDraftRepository implements DraftRepository {
private readonly map = new Map<string, { key: ConversationKey; text: string }>();
private readonly map = new Map<
string,
{ key: ConversationKey; text: string }
>();
async load(key: ConversationKey): Promise<string> {
return this.map.get(conversationKeyString(key))?.text ?? "";
@@ -20,7 +26,9 @@ export class InMemoryDraftRepository implements DraftRepository {
this.map.delete(conversationKeyString(key));
}
async loadAll(): Promise<ReadonlyArray<{ key: ConversationKey; text: string }>> {
async loadAll(): Promise<
ReadonlyArray<{ key: ConversationKey; text: string }>
> {
return Array.from(this.map.values());
}
}
@@ -21,13 +21,21 @@ describe("InMemoryMessageRepository", () => {
it("loadRecent returns the tail of a bucket", async () => {
const repo = new InMemoryMessageRepository();
await repo.appendBatch([msg(1, 1000), msg(2, 2000), msg(3, 3000)]);
const out = await repo.loadRecent({ kind: "channel", channel: ChannelNumber.Primary }, 2);
const out = await repo.loadRecent(
{ kind: "channel", channel: ChannelNumber.Primary },
2,
);
expect(out.map((m) => m.id)).toEqual([2, 3]);
});
it("loadBefore paginates older messages", async () => {
const repo = new InMemoryMessageRepository();
await repo.appendBatch([msg(1, 1000), msg(2, 2000), msg(3, 3000), msg(4, 4000)]);
await repo.appendBatch([
msg(1, 1000),
msg(2, 2000),
msg(3, 3000),
msg(4, 4000),
]);
const out = await repo.loadBefore(
{ kind: "channel", channel: ChannelNumber.Primary },
new Date(3000),
@@ -40,24 +48,41 @@ describe("InMemoryMessageRepository", () => {
const repo = new InMemoryMessageRepository();
await repo.append(msg(42, 1000));
await repo.updateState(42, MessageState.Failed);
const [found] = await repo.loadRecent({ kind: "channel", channel: ChannelNumber.Primary }, 1);
const [found] = await repo.loadRecent(
{ kind: "channel", channel: ChannelNumber.Primary },
1,
);
expect(found?.state).toBe(MessageState.Failed);
});
it("prune enforces maxPerBucket", async () => {
const repo = new InMemoryMessageRepository();
await repo.appendBatch([msg(1, 1000), msg(2, 2000), msg(3, 3000), msg(4, 4000)]);
await repo.appendBatch([
msg(1, 1000),
msg(2, 2000),
msg(3, 3000),
msg(4, 4000),
]);
await repo.prune({ maxPerBucket: 2 });
const out = await repo.loadRecent({ kind: "channel", channel: ChannelNumber.Primary }, 10);
const out = await repo.loadRecent(
{ kind: "channel", channel: ChannelNumber.Primary },
10,
);
expect(out.map((m) => m.id)).toEqual([3, 4]);
});
it("prune enforces olderThanMs", async () => {
const repo = new InMemoryMessageRepository();
const now = Date.now();
await repo.appendBatch([msg(1, now - 1000 * 60 * 60 * 24 * 10), msg(2, now)]);
await repo.appendBatch([
msg(1, now - 1000 * 60 * 60 * 24 * 10),
msg(2, now),
]);
await repo.prune({ olderThanMs: 1000 * 60 * 60 * 24 });
const out = await repo.loadRecent({ kind: "channel", channel: ChannelNumber.Primary }, 10);
const out = await repo.loadRecent(
{ kind: "channel", channel: ChannelNumber.Primary },
10,
);
expect(out.map((m) => m.id)).toEqual([2]);
});
});
@@ -19,7 +19,11 @@ export class InMemoryMessageRepository implements MessageRepository {
return bucket.slice(-limit);
}
async loadBefore(key: ConversationKey, cursor: Date, limit: number): Promise<Message[]> {
async loadBefore(
key: ConversationKey,
cursor: Date,
limit: number,
): Promise<Message[]> {
const bucket = this.buckets.get(conversationKeyString(key)) ?? [];
const idx = bucket.findIndex((m) => m.rxTime >= cursor);
const end = idx === -1 ? bucket.length : idx;
@@ -59,8 +63,13 @@ export class InMemoryMessageRepository implements MessageRepository {
let filtered =
policy.olderThanMs === undefined
? bucket
: bucket.filter((m) => nowMs - m.rxTime.getTime() <= policy.olderThanMs!);
if (policy.maxPerBucket !== undefined && filtered.length > policy.maxPerBucket) {
: bucket.filter(
(m) => nowMs - m.rxTime.getTime() <= policy.olderThanMs!,
);
if (
policy.maxPerBucket !== undefined &&
filtered.length > policy.maxPerBucket
) {
filtered = filtered.slice(-policy.maxPerBucket);
}
this.buckets.set(key, filtered);
@@ -1,6 +1,12 @@
import { type Signal, signal } from "@preact/signals-core";
import { type ReadonlySignal, toReadonly } from "../../../core/signals/createStore.ts";
import { type ConversationKey, conversationKeyString } from "../domain/MessageRepository.ts";
import {
type ReadonlySignal,
toReadonly,
} from "../../../core/signals/createStore.ts";
import {
type ConversationKey,
conversationKeyString,
} from "../domain/MessageRepository.ts";
/**
* Per-conversation draft text exposed as readonly signals. Lazy creation:
@@ -2,7 +2,10 @@ import * as Protobuf from "@meshtastic/protobufs";
import { computed } from "@preact/signals-core";
import type { ResultType } from "better-result";
import type { MeshClient } from "../../core/client/MeshClient.ts";
import { type ReadonlySignal, toReadonly } from "../../core/signals/createStore.ts";
import {
type ReadonlySignal,
toReadonly,
} from "../../core/signals/createStore.ts";
import {
beginEditSettings,
commitEditSettings,
@@ -45,12 +48,17 @@ export class ConfigClient {
computed(() => {
const lora = this.store.radio.write.value.lora;
if (!lora) return false;
return lora.region === Protobuf.Config.Config_LoRaConfig_RegionCode.UNSET;
return (
lora.region === Protobuf.Config.Config_LoRaConfig_RegionCode.UNSET
);
}),
);
client.events.onConfigPacket.subscribe((config) => {
this.store.radio.write.value = ConfigMapper.mergeRadio(this.store.radio.write.value, config);
this.store.radio.write.value = ConfigMapper.mergeRadio(
this.store.radio.write.value,
config,
);
});
client.events.onModuleConfigPacket.subscribe((moduleConfig) => {
this.store.modules.write.value = ConfigMapper.mergeModule(
@@ -68,7 +76,9 @@ export class ConfigClient {
return commitEditSettings(this.client);
}
public setRadio(config: Protobuf.Config.Config): Promise<ResultType<number, Error>> {
public setRadio(
config: Protobuf.Config.Config,
): Promise<ResultType<number, Error>> {
return setConfig(this.client, config);
}
@@ -18,7 +18,9 @@ function mqttPacket(enabled: boolean): Protobuf.ModuleConfig.ModuleConfig {
return create(Protobuf.ModuleConfig.ModuleConfigSchema, {
payloadVariant: {
case: "mqtt",
value: create(Protobuf.ModuleConfig.ModuleConfig_MQTTConfigSchema, { enabled }),
value: create(Protobuf.ModuleConfig.ModuleConfig_MQTTConfigSchema, {
enabled,
}),
},
});
}
@@ -43,7 +45,10 @@ describe("ConfigEditor", () => {
expect(editor.isDirty.value).toBe(false);
expect(editor.radio.value.lora?.region).toBe(4);
editor.setRadioSection("lora", create(Protobuf.Config.Config_LoRaConfigSchema, { region: 7 }));
editor.setRadioSection(
"lora",
create(Protobuf.Config.Config_LoRaConfigSchema, { region: 7 }),
);
expect(editor.isDirty.value).toBe(true);
expect(editor.dirtyRadioSections.value).toEqual(["lora"]);
});
@@ -58,7 +63,9 @@ describe("ConfigEditor", () => {
editor.setModuleSection(
"mqtt",
create(Protobuf.ModuleConfig.ModuleConfig_MQTTConfigSchema, { enabled: true }),
create(Protobuf.ModuleConfig.ModuleConfig_MQTTConfigSchema, {
enabled: true,
}),
);
expect(editor.dirtyModuleSections.value).toEqual(["mqtt"]);
@@ -73,7 +80,10 @@ describe("ConfigEditor", () => {
const editor = client.config.editor;
client.events.onConfigPacket.dispatch(loraPacket(4));
editor.setRadioSection("lora", create(Protobuf.Config.Config_LoRaConfigSchema, { region: 7 }));
editor.setRadioSection(
"lora",
create(Protobuf.Config.Config_LoRaConfigSchema, { region: 7 }),
);
expect(editor.isDirty.value).toBe(true);
editor.reset();
@@ -88,7 +98,10 @@ describe("ConfigEditor", () => {
// User edits lora region
client.events.onConfigPacket.dispatch(loraPacket(4));
editor.setRadioSection("lora", create(Protobuf.Config.Config_LoRaConfigSchema, { region: 7 }));
editor.setRadioSection(
"lora",
create(Protobuf.Config.Config_LoRaConfigSchema, { region: 7 }),
);
// Device pushes a baseline change for the same section while user is editing
client.events.onConfigPacket.dispatch(loraPacket(8));
@@ -104,7 +117,10 @@ describe("ConfigEditor", () => {
const editor = client.config.editor;
client.events.onConfigPacket.dispatch(loraPacket(4));
editor.setRadioSection("lora", create(Protobuf.Config.Config_LoRaConfigSchema, { region: 7 }));
editor.setRadioSection(
"lora",
create(Protobuf.Config.Config_LoRaConfigSchema, { region: 7 }),
);
expect(editor.isDirty.value).toBe(true);
client.events.onDeviceStatus.dispatch(DeviceStatusEnum.DeviceDisconnected);
@@ -4,20 +4,30 @@ import type { ResultType } from "better-result";
import type { MeshClient } from "../../../core/client/MeshClient.ts";
import { sendAdminMessage } from "../../device/infrastructure/AdminMessageSender.ts";
export async function beginEditSettings(client: MeshClient): Promise<ResultType<number, Error>> {
export async function beginEditSettings(
client: MeshClient,
): Promise<ResultType<number, Error>> {
client.events.onPendingSettingsChange.dispatch(true);
try {
const id = await sendAdminMessage(client, { case: "beginEditSettings", value: true });
const id = await sendAdminMessage(client, {
case: "beginEditSettings",
value: true,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
}
}
export async function commitEditSettings(client: MeshClient): Promise<ResultType<number, Error>> {
export async function commitEditSettings(
client: MeshClient,
): Promise<ResultType<number, Error>> {
client.events.onPendingSettingsChange.dispatch(false);
try {
const id = await sendAdminMessage(client, { case: "commitEditSettings", value: true });
const id = await sendAdminMessage(client, {
case: "commitEditSettings",
value: true,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -29,7 +39,10 @@ export async function setConfig(
config: Protobuf.Config.Config,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "setConfig", value: config });
const id = await sendAdminMessage(client, {
case: "setConfig",
value: config,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -41,7 +54,10 @@ export async function getConfig(
type: Protobuf.Admin.AdminMessage_ConfigType,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "getConfigRequest", value: type });
const id = await sendAdminMessage(client, {
case: "getConfigRequest",
value: type,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -53,7 +69,10 @@ export async function setModuleConfig(
moduleConfig: Protobuf.ModuleConfig.ModuleConfig,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "setModuleConfig", value: moduleConfig });
const id = await sendAdminMessage(client, {
case: "setModuleConfig",
value: moduleConfig,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -65,7 +84,10 @@ export async function getModuleConfig(
type: Protobuf.Admin.AdminMessage_ModuleConfigType,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "getModuleConfigRequest", value: type });
const id = await sendAdminMessage(client, {
case: "getModuleConfigRequest",
value: type,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -4,7 +4,10 @@ import { signal } from "@preact/signals-core";
import { Result } from "better-result";
import type { ResultType } from "better-result";
import type { MeshClient } from "../../../core/client/MeshClient.ts";
import { type ReadonlySignal, toReadonly } from "../../../core/signals/createStore.ts";
import {
type ReadonlySignal,
toReadonly,
} from "../../../core/signals/createStore.ts";
import { setChannel } from "../../channels/application/ChannelUseCases.ts";
import { DeviceStatusEnum } from "../../device/domain/DeviceStatus.ts";
import { sendAdminMessage } from "../../device/infrastructure/AdminMessageSender.ts";
@@ -39,40 +42,56 @@ export class ConfigEditor {
private readonly client: MeshClient;
private readonly baselineRadio = signal<RadioConfig>({});
private readonly baselineModules = signal<ModuleConfig>({});
private readonly baselineChannels = signal<ReadonlyMap<number, Protobuf.Channel.Channel>>(
new Map(),
);
private readonly baselineChannels = signal<
ReadonlyMap<number, Protobuf.Channel.Channel>
>(new Map());
private readonly workingRadio = signal<RadioConfig>({});
private readonly workingModules = signal<ModuleConfig>({});
private readonly workingChannels = signal<ReadonlyMap<number, Protobuf.Channel.Channel>>(
new Map(),
private readonly workingChannels = signal<
ReadonlyMap<number, Protobuf.Channel.Channel>
>(new Map());
private readonly baselineOwner = signal<Protobuf.Mesh.User | undefined>(
undefined,
);
private readonly baselineOwner = signal<Protobuf.Mesh.User | undefined>(undefined);
private readonly workingOwner = signal<Protobuf.Mesh.User | undefined>(undefined);
private readonly queuedAdminMessages = signal<readonly Protobuf.Admin.AdminMessage[]>([]);
private readonly _dirtyRadioSections = signal<readonly RadioConfigSection[]>([]);
private readonly _dirtyModuleSections = signal<readonly ModuleConfigSection[]>([]);
private readonly workingOwner = signal<Protobuf.Mesh.User | undefined>(
undefined,
);
private readonly queuedAdminMessages = signal<
readonly Protobuf.Admin.AdminMessage[]
>([]);
private readonly _dirtyRadioSections = signal<readonly RadioConfigSection[]>(
[],
);
private readonly _dirtyModuleSections = signal<
readonly ModuleConfigSection[]
>([]);
private readonly _dirtyChannels = signal<readonly number[]>([]);
private readonly _isOwnerDirty = signal<boolean>(false);
private readonly _isDirty = signal<boolean>(false);
public readonly radio: ReadonlySignal<RadioConfig> = toReadonly(this.workingRadio);
public readonly modules: ReadonlySignal<ModuleConfig> = toReadonly(this.workingModules);
public readonly channels: ReadonlySignal<ReadonlyMap<number, Protobuf.Channel.Channel>> =
toReadonly(this.workingChannels);
public readonly dirtyRadioSections: ReadonlySignal<readonly RadioConfigSection[]> = toReadonly(
this._dirtyRadioSections,
public readonly radio: ReadonlySignal<RadioConfig> = toReadonly(
this.workingRadio,
);
public readonly dirtyModuleSections: ReadonlySignal<readonly ModuleConfigSection[]> = toReadonly(
this._dirtyModuleSections,
public readonly modules: ReadonlySignal<ModuleConfig> = toReadonly(
this.workingModules,
);
public readonly channels: ReadonlySignal<
ReadonlyMap<number, Protobuf.Channel.Channel>
> = toReadonly(this.workingChannels);
public readonly dirtyRadioSections: ReadonlySignal<
readonly RadioConfigSection[]
> = toReadonly(this._dirtyRadioSections);
public readonly dirtyModuleSections: ReadonlySignal<
readonly ModuleConfigSection[]
> = toReadonly(this._dirtyModuleSections);
public readonly dirtyChannels: ReadonlySignal<readonly number[]> = toReadonly(
this._dirtyChannels,
);
public readonly owner: ReadonlySignal<Protobuf.Mesh.User | undefined> = toReadonly(
this.workingOwner,
public readonly owner: ReadonlySignal<Protobuf.Mesh.User | undefined> =
toReadonly(this.workingOwner);
public readonly isOwnerDirty: ReadonlySignal<boolean> = toReadonly(
this._isOwnerDirty,
);
public readonly isOwnerDirty: ReadonlySignal<boolean> = toReadonly(this._isOwnerDirty);
public readonly isDirty: ReadonlySignal<boolean> = toReadonly(this._isDirty);
constructor(client: MeshClient) {
@@ -88,14 +107,20 @@ export class ConfigEditor {
// their edit in place; the dirty bookkeeping will refresh below.
const wasDirty = this._dirtyRadioSections.peek().includes(variant);
if (!wasDirty) {
this.workingRadio.value = { ...this.workingRadio.peek(), [variant]: next[variant] };
this.workingRadio.value = {
...this.workingRadio.peek(),
[variant]: next[variant],
};
}
}
this.recomputeDirty();
});
client.events.onModuleConfigPacket.subscribe((moduleConfig) => {
const next = ConfigMapper.mergeModule(this.baselineModules.peek(), moduleConfig);
const next = ConfigMapper.mergeModule(
this.baselineModules.peek(),
moduleConfig,
);
this.baselineModules.value = next;
const variant = moduleConfig.payloadVariant.case;
if (variant) {
@@ -168,7 +193,10 @@ export class ConfigEditor {
* alongside their other edits.
*/
public queueAdminMessage(message: Protobuf.Admin.AdminMessage): void {
this.queuedAdminMessages.value = [...this.queuedAdminMessages.peek(), message];
this.queuedAdminMessages.value = [
...this.queuedAdminMessages.peek(),
message,
];
this.recomputeDirty();
}
@@ -279,10 +307,16 @@ export class ConfigEditor {
const radioBase = this.baselineRadio.peek();
const radioWorking = this.workingRadio.peek();
const radioDirty: RadioConfigSection[] = [];
const radioKeys = new Set<string>([...Object.keys(radioBase), ...Object.keys(radioWorking)]);
const radioKeys = new Set<string>([
...Object.keys(radioBase),
...Object.keys(radioWorking),
]);
for (const key of radioKeys) {
if (
!shallowEqual(radioBase[key as keyof RadioConfig], radioWorking[key as keyof RadioConfig])
!shallowEqual(
radioBase[key as keyof RadioConfig],
radioWorking[key as keyof RadioConfig],
)
) {
radioDirty.push(key as RadioConfigSection);
}
@@ -291,7 +325,10 @@ export class ConfigEditor {
const moduleBase = this.baselineModules.peek();
const moduleWorking = this.workingModules.peek();
const moduleDirty: ModuleConfigSection[] = [];
const moduleKeys = new Set<string>([...Object.keys(moduleBase), ...Object.keys(moduleWorking)]);
const moduleKeys = new Set<string>([
...Object.keys(moduleBase),
...Object.keys(moduleWorking),
]);
for (const key of moduleKeys) {
if (
!shallowEqual(
@@ -306,14 +343,20 @@ export class ConfigEditor {
const channelDirty: number[] = [];
const channelBase = this.baselineChannels.peek();
const channelWorking = this.workingChannels.peek();
const channelKeys = new Set<number>([...channelBase.keys(), ...channelWorking.keys()]);
const channelKeys = new Set<number>([
...channelBase.keys(),
...channelWorking.keys(),
]);
for (const idx of channelKeys) {
if (!shallowEqual(channelBase.get(idx), channelWorking.get(idx))) {
channelDirty.push(idx);
}
}
const ownerDirty = !shallowEqual(this.baselineOwner.peek(), this.workingOwner.peek());
const ownerDirty = !shallowEqual(
this.baselineOwner.peek(),
this.workingOwner.peek(),
);
const hasQueuedAdmin = this.queuedAdminMessages.peek().length > 0;
this._dirtyRadioSections.value = radioDirty;
@@ -329,9 +372,15 @@ export class ConfigEditor {
}
}
function buildRadioConfig(variant: RadioConfigSection, value: unknown): Protobuf.Config.Config {
function buildRadioConfig(
variant: RadioConfigSection,
value: unknown,
): Protobuf.Config.Config {
return create(Protobuf.Config.ConfigSchema, {
payloadVariant: { case: variant, value } as Protobuf.Config.Config["payloadVariant"],
payloadVariant: {
case: variant,
value,
} as Protobuf.Config.Config["payloadVariant"],
});
}
+4 -1
View File
@@ -1,5 +1,8 @@
export { ConfigClient } from "./ConfigClient.ts";
export { ConfigEditor } from "./domain/ConfigEditor.ts";
export type { RadioConfig, RadioConfigSection } from "./domain/RadioConfig.ts";
export type { ModuleConfig, ModuleConfigSection } from "./domain/ModuleConfig.ts";
export type {
ModuleConfig,
ModuleConfigSection,
} from "./domain/ModuleConfig.ts";
export { ConfigMapper } from "./infrastructure/ConfigMapper.ts";
@@ -3,12 +3,18 @@ import type { ModuleConfig } from "../domain/ModuleConfig.ts";
import type { RadioConfig } from "../domain/RadioConfig.ts";
export const ConfigMapper = {
mergeRadio(existing: RadioConfig, incoming: Protobuf.Config.Config): RadioConfig {
mergeRadio(
existing: RadioConfig,
incoming: Protobuf.Config.Config,
): RadioConfig {
const variant = incoming.payloadVariant;
if (!variant.case) return existing;
return { ...existing, [variant.case]: variant.value };
},
mergeModule(existing: ModuleConfig, incoming: Protobuf.ModuleConfig.ModuleConfig): ModuleConfig {
mergeModule(
existing: ModuleConfig,
incoming: Protobuf.ModuleConfig.ModuleConfig,
): ModuleConfig {
const variant = incoming.payloadVariant;
if (!variant.case) return existing;
return { ...existing, [variant.case]: variant.value };
@@ -18,8 +18,12 @@ export class DeviceClient {
public readonly isConfigured: ReadonlySignal<boolean>;
public readonly pendingSettingsChanges: ReadonlySignal<boolean>;
public readonly myNodeNum: ReadonlySignal<number | undefined>;
public readonly metadata: ReadonlySignal<Protobuf.Mesh.DeviceMetadata | undefined>;
public readonly myNodeInfo: ReadonlySignal<Protobuf.Mesh.MyNodeInfo | undefined>;
public readonly metadata: ReadonlySignal<
Protobuf.Mesh.DeviceMetadata | undefined
>;
public readonly myNodeInfo: ReadonlySignal<
Protobuf.Mesh.MyNodeInfo | undefined
>;
constructor(client: MeshClient) {
this.client = client;
@@ -2,7 +2,10 @@ import type { MeshClient } from "../../../core/client/MeshClient.ts";
import { ChannelNumber } from "../../../core/types.ts";
import { sendAdminMessage } from "../infrastructure/AdminMessageSender.ts";
export function getMetadata(client: MeshClient, nodeNum: number): Promise<number> {
export function getMetadata(
client: MeshClient,
nodeNum: number,
): Promise<number> {
return sendAdminMessage(
client,
{ case: "getDeviceMetadataRequest", value: true },
@@ -9,7 +9,10 @@ export function reboot(client: MeshClient, seconds: number): Promise<number> {
return sendAdminMessage(client, { case: "rebootSeconds", value: seconds });
}
export function rebootOta(client: MeshClient, seconds: number): Promise<number> {
export function rebootOta(
client: MeshClient,
seconds: number,
): Promise<number> {
return sendAdminMessage(client, { case: "rebootOtaSeconds", value: seconds });
}
+5 -1
View File
@@ -1,3 +1,7 @@
export { DeviceClient } from "./DeviceClient.ts";
export type { Device } from "./domain/Device.ts";
export { DeviceStatusEnum, isConfigured, isConnected } from "./domain/DeviceStatus.ts";
export {
DeviceStatusEnum,
isConfigured,
isConnected,
} from "./domain/DeviceStatus.ts";
@@ -8,14 +8,27 @@ import { DeviceStatusEnum } from "../../../core/transport/Transport.ts";
* DeviceClient.
*/
export function createDeviceStore() {
const status = createStore<DeviceStatusEnum>(DeviceStatusEnum.DeviceDisconnected);
const status = createStore<DeviceStatusEnum>(
DeviceStatusEnum.DeviceDisconnected,
);
const isConfigured = createStore(false);
const pendingSettingsChanges = createStore(false);
const myNodeNum = createStore<number | undefined>(undefined);
const metadata = createStore<Protobuf.Mesh.DeviceMetadata | undefined>(undefined);
const myNodeInfo = createStore<Protobuf.Mesh.MyNodeInfo | undefined>(undefined);
const metadata = createStore<Protobuf.Mesh.DeviceMetadata | undefined>(
undefined,
);
const myNodeInfo = createStore<Protobuf.Mesh.MyNodeInfo | undefined>(
undefined,
);
return { status, isConfigured, pendingSettingsChanges, myNodeNum, metadata, myNodeInfo };
return {
status,
isConfigured,
pendingSettingsChanges,
myNodeNum,
metadata,
myNodeInfo,
};
}
export type DeviceStore = ReturnType<typeof createDeviceStore>;
@@ -42,7 +42,9 @@ export class FilesClient {
}
}
public async download(filename: string): Promise<ResultType<FileTransfer, Error>> {
public async download(
filename: string,
): Promise<ResultType<FileTransfer, Error>> {
const id = generatePacketId();
const transfer: FileTransfer = {
id,
@@ -81,7 +81,9 @@ describe("NodesClient PKI error tracking", () => {
}),
});
expect(client.nodes.errorFor(99)?.error).toBe(Protobuf.Mesh.Routing_Error.PKI_UNKNOWN_PUBKEY);
expect(client.nodes.errorFor(99)?.error).toBe(
Protobuf.Mesh.Routing_Error.PKI_UNKNOWN_PUBKEY,
);
});
it("clearError / clearAllErrors", () => {
@@ -11,7 +11,9 @@ describe("NodesClient", () => {
expect(client.nodes.list.value).toEqual([]);
client.events.onNodeInfoPacket.dispatch(create(Protobuf.Mesh.NodeInfoSchema, { num: 1 }));
client.events.onNodeInfoPacket.dispatch(
create(Protobuf.Mesh.NodeInfoSchema, { num: 1 }),
);
client.events.onNodeInfoPacket.dispatch(
create(Protobuf.Mesh.NodeInfoSchema, { num: 2, isFavorite: true }),
);
+18 -5
View File
@@ -11,9 +11,18 @@ import { InMemoryNodesRepository } from "./infrastructure/repositories/InMemoryN
import { validateIncomingNode } from "./infrastructure/nodeValidation.ts";
import { NodeErrorsStore } from "./state/nodeErrorsStore.ts";
import { NodesStore } from "./state/nodesStore.ts";
import { favoriteNode, removeFavoriteNode } from "./application/FavoriteNodeUseCase.ts";
import { ignoreNode, removeIgnoredNode } from "./application/IgnoreNodeUseCase.ts";
import { removeNodeByNum, resetNodes } from "./application/RemoveNodeUseCase.ts";
import {
favoriteNode,
removeFavoriteNode,
} from "./application/FavoriteNodeUseCase.ts";
import {
ignoreNode,
removeIgnoredNode,
} from "./application/IgnoreNodeUseCase.ts";
import {
removeNodeByNum,
resetNodes,
} from "./application/RemoveNodeUseCase.ts";
export interface NodesClientOptions {
repository?: NodesRepository;
@@ -42,7 +51,9 @@ export class NodesClient {
this.list = this.store.read;
this.errors = this.errorsStore.read;
client.events.onNodeInfoPacket.subscribe((info) => this.handleIncoming(info));
client.events.onNodeInfoPacket.subscribe((info) =>
this.handleIncoming(info),
);
client.events.onUserPacket.subscribe((packet) => {
this.patch(packet.from, { user: packet.data });
@@ -191,7 +202,9 @@ export class NodesClient {
* removeAllNodes(true) + resetNodes flow that the ResetNodeDb dialog
* relied on.
*/
public async reset(options: { keepMyNode?: boolean } = {}): Promise<ResultType<number, Error>> {
public async reset(
options: { keepMyNode?: boolean } = {},
): Promise<ResultType<number, Error>> {
const myNodeNum = this.client.device.myNodeNum.value;
if (options.keepMyNode && myNodeNum !== undefined) {
const me = this.store.get(myNodeNum);
@@ -8,7 +8,10 @@ export async function favoriteNode(
nodeNum: number,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "setFavoriteNode", value: nodeNum });
const id = await sendAdminMessage(client, {
case: "setFavoriteNode",
value: nodeNum,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -20,7 +23,10 @@ export async function removeFavoriteNode(
nodeNum: number,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "removeFavoriteNode", value: nodeNum });
const id = await sendAdminMessage(client, {
case: "removeFavoriteNode",
value: nodeNum,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -8,7 +8,10 @@ export async function ignoreNode(
nodeNum: number,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "setIgnoredNode", value: nodeNum });
const id = await sendAdminMessage(client, {
case: "setIgnoredNode",
value: nodeNum,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -20,7 +23,10 @@ export async function removeIgnoredNode(
nodeNum: number,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "removeIgnoredNode", value: nodeNum });
const id = await sendAdminMessage(client, {
case: "removeIgnoredNode",
value: nodeNum,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -8,16 +8,24 @@ export async function removeNodeByNum(
nodeNum: number,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "removeByNodenum", value: nodeNum });
const id = await sendAdminMessage(client, {
case: "removeByNodenum",
value: nodeNum,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
}
}
export async function resetNodes(client: MeshClient): Promise<ResultType<number, Error>> {
export async function resetNodes(
client: MeshClient,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "nodedbReset", value: true });
const id = await sendAdminMessage(client, {
case: "nodedbReset",
value: true,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -9,7 +9,10 @@ export async function setOwner(
owner: Protobuf.Mesh.User,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(client, { case: "setOwner", value: owner });
const id = await sendAdminMessage(client, {
case: "setOwner",
value: owner,
});
return Result.ok(id);
} catch (e) {
return Result.err(e instanceof Error ? e : new Error(String(e)));
@@ -7,7 +7,10 @@ import type * as Protobuf from "@meshtastic/protobufs";
* stored public-key history. Routing errors come straight off the wire
* via `Routing_Error`.
*/
export type NodeErrorType = Protobuf.Mesh.Routing_Error | "MISMATCH_PKI" | "DUPLICATE_PKI";
export type NodeErrorType =
| Protobuf.Mesh.Routing_Error
| "MISMATCH_PKI"
| "DUPLICATE_PKI";
export interface NodeError {
readonly node: number;
+4 -1
View File
@@ -4,5 +4,8 @@ export type { Node } from "./domain/Node.ts";
export type { NodeError, NodeErrorType } from "./domain/NodeError.ts";
export type { NodesRepository } from "./domain/NodesRepository.ts";
export { NodeMapper } from "./infrastructure/NodeMapper.ts";
export { equalKey, validateIncomingNode } from "./infrastructure/nodeValidation.ts";
export {
equalKey,
validateIncomingNode,
} from "./infrastructure/nodeValidation.ts";
export { InMemoryNodesRepository } from "./infrastructure/repositories/InMemoryNodesRepository.ts";
@@ -6,7 +6,10 @@ import type { NodeErrorType } from "../domain/NodeError.ts";
* Byte-equal compare for two public keys. Empty / undefined keys never
* match each other.
*/
export function equalKey(a?: Uint8Array | null, b?: Uint8Array | null): boolean {
export function equalKey(
a?: Uint8Array | null,
b?: Uint8Array | null,
): boolean {
if (!a || !b) return false;
if (a === b) return true;
if (a.byteLength !== b.byteLength) return false;
@@ -31,7 +31,10 @@ export class PositionClient {
return this.store.get(nodeNum);
}
public setFixed(latitude: number, longitude: number): Promise<ResultType<number, Error>> {
public setFixed(
latitude: number,
longitude: number,
): Promise<ResultType<number, Error>> {
return setFixedPosition(this.client, latitude, longitude);
}
@@ -39,7 +42,9 @@ export class PositionClient {
return removeFixedPosition(this.client);
}
public set(position: Protobuf.Mesh.Position): Promise<ResultType<number, Error>> {
public set(
position: Protobuf.Mesh.Position,
): Promise<ResultType<number, Error>> {
return setPosition(this.client, position);
}
@@ -29,7 +29,9 @@ export async function setFixedPosition(
}
}
export async function removeFixedPosition(client: MeshClient): Promise<ResultType<number, Error>> {
export async function removeFixedPosition(
client: MeshClient,
): Promise<ResultType<number, Error>> {
try {
const id = await sendAdminMessage(
client,
@@ -6,7 +6,12 @@ import { createFakeTransport } from "../../core/testing/createFakeTransport.ts";
import { ChannelNumber } from "../../core/types.ts";
import { InMemoryTelemetryRepository } from "./infrastructure/repositories/InMemoryTelemetryRepository.ts";
const dispatch = (client: MeshClient, from: number, ms: number, battery = 80): void => {
const dispatch = (
client: MeshClient,
from: number,
ms: number,
battery = 80,
): void => {
client.events.onTelemetryPacket.dispatch({
id: ms,
from,
@@ -18,7 +23,9 @@ const dispatch = (client: MeshClient, from: number, ms: number, battery = 80): v
time: ms,
variant: {
case: "deviceMetrics",
value: create(Protobuf.Telemetry.DeviceMetricsSchema, { batteryLevel: battery }),
value: create(Protobuf.Telemetry.DeviceMetricsSchema, {
batteryLevel: battery,
}),
},
}),
});
@@ -46,7 +53,9 @@ describe("TelemetryClient persistence", () => {
nodeNum: 200,
time: new Date(500),
kind: "deviceMetrics",
value: create(Protobuf.Telemetry.DeviceMetricsSchema, { batteryLevel: 60 }),
value: create(Protobuf.Telemetry.DeviceMetricsSchema, {
batteryLevel: 60,
}),
});
const { transport } = createFakeTransport();
@@ -21,7 +21,9 @@ describe("TelemetryClient", () => {
time: 1000,
variant: {
case: "deviceMetrics",
value: create(Protobuf.Telemetry.DeviceMetricsSchema, { batteryLevel: 80 }),
value: create(Protobuf.Telemetry.DeviceMetricsSchema, {
batteryLevel: 80,
}),
},
}),
});
@@ -36,7 +38,9 @@ describe("TelemetryClient", () => {
time: 2000,
variant: {
case: "deviceMetrics",
value: create(Protobuf.Telemetry.DeviceMetricsSchema, { batteryLevel: 70 }),
value: create(Protobuf.Telemetry.DeviceMetricsSchema, {
batteryLevel: 70,
}),
},
}),
});
@@ -56,7 +56,11 @@ export class TelemetryClient {
* repository. Caller is responsible for merging with the in-memory store
* if it wants the result reflected in the `history(nodeNum)` signal.
*/
public loadBefore(nodeNum: number, cursor: Date, limit: number): Promise<TelemetryReading[]> {
public loadBefore(
nodeNum: number,
cursor: Date,
limit: number,
): Promise<TelemetryReading[]> {
return this.repository.loadBefore(nodeNum, cursor, limit);
}
@@ -17,7 +17,11 @@ export interface TelemetryRetentionPolicy {
*/
export interface TelemetryRepository {
loadRecent(nodeNum: number, limit: number): Promise<TelemetryReading[]>;
loadBefore(nodeNum: number, cursor: Date, limit: number): Promise<TelemetryReading[]>;
loadBefore(
nodeNum: number,
cursor: Date,
limit: number,
): Promise<TelemetryReading[]>;
append(reading: TelemetryReading): Promise<void>;
appendBatch(readings: ReadonlyArray<TelemetryReading>): Promise<void>;
prune(policy: TelemetryRetentionPolicy): Promise<void>;
+4 -1
View File
@@ -1,6 +1,9 @@
export { TelemetryClient } from "./TelemetryClient.ts";
export type { TelemetryClientOptions } from "./TelemetryClient.ts";
export type { TelemetryKind, TelemetryReading } from "./domain/TelemetryReading.ts";
export type {
TelemetryKind,
TelemetryReading,
} from "./domain/TelemetryReading.ts";
export type {
TelemetryRepository,
TelemetryRetentionPolicy,
@@ -3,7 +3,9 @@ import type * as Protobuf from "@meshtastic/protobufs";
import type { TelemetryReading } from "../domain/TelemetryReading.ts";
export const TelemetryMapper = {
fromPacket(packet: PacketMetadata<Protobuf.Telemetry.Telemetry>): TelemetryReading {
fromPacket(
packet: PacketMetadata<Protobuf.Telemetry.Telemetry>,
): TelemetryReading {
return {
nodeNum: packet.from,
time: packet.rxTime,
@@ -36,7 +36,10 @@ describe("InMemoryTelemetryRepository", () => {
it("prune drops readings older than olderThanMs", async () => {
const repo = new InMemoryTelemetryRepository();
await repo.appendBatch([reading(1, Date.now() - 10_000_000), reading(1, Date.now() - 1_000)]);
await repo.appendBatch([
reading(1, Date.now() - 10_000_000),
reading(1, Date.now() - 1_000),
]);
await repo.prune({ olderThanMs: 5_000_000 });
const remaining = await repo.loadRecent(1, 10);
expect(remaining).toHaveLength(1);
@@ -11,11 +11,18 @@ import type {
export class InMemoryTelemetryRepository implements TelemetryRepository {
private readonly buckets = new Map<number, TelemetryReading[]>();
async loadRecent(nodeNum: number, limit: number): Promise<TelemetryReading[]> {
async loadRecent(
nodeNum: number,
limit: number,
): Promise<TelemetryReading[]> {
return (this.buckets.get(nodeNum) ?? []).slice(-limit);
}
async loadBefore(nodeNum: number, cursor: Date, limit: number): Promise<TelemetryReading[]> {
async loadBefore(
nodeNum: number,
cursor: Date,
limit: number,
): Promise<TelemetryReading[]> {
const bucket = this.buckets.get(nodeNum) ?? [];
const idx = bucket.findIndex((r) => r.time >= cursor);
const end = idx === -1 ? bucket.length : idx;
@@ -37,7 +44,9 @@ export class InMemoryTelemetryRepository implements TelemetryRepository {
}
async prune(policy: TelemetryRetentionPolicy): Promise<void> {
const cutoff = policy.olderThanMs ? Date.now() - policy.olderThanMs : undefined;
const cutoff = policy.olderThanMs
? Date.now() - policy.olderThanMs
: undefined;
for (const [nodeNum, bucket] of this.buckets.entries()) {
let next = bucket;
if (cutoff !== undefined) {
@@ -1,14 +1,26 @@
import { type Signal, signal } from "@preact/signals-core";
import { type ReadonlySignal, toReadonly } from "../../../core/signals/createStore.ts";
import {
type ReadonlySignal,
toReadonly,
} from "../../../core/signals/createStore.ts";
import type { TelemetryReading } from "../domain/TelemetryReading.ts";
const MAX_HISTORY = 256;
export class TelemetryStore {
private readonly latest = new Map<number, Signal<TelemetryReading | undefined>>();
private readonly latestRead = new Map<number, ReadonlySignal<TelemetryReading | undefined>>();
private readonly latest = new Map<
number,
Signal<TelemetryReading | undefined>
>();
private readonly latestRead = new Map<
number,
ReadonlySignal<TelemetryReading | undefined>
>();
private readonly history = new Map<number, Signal<TelemetryReading[]>>();
private readonly historyRead = new Map<number, ReadonlySignal<TelemetryReading[]>>();
private readonly historyRead = new Map<
number,
ReadonlySignal<TelemetryReading[]>
>();
append(reading: TelemetryReading): void {
const latestSig = this.ensureLatest(reading.nodeNum);
@@ -9,7 +9,9 @@ export async function runTraceRoute(
destination: number,
): Promise<ResultType<number, Error>> {
try {
const routeDiscovery = create(Protobuf.Mesh.RouteDiscoverySchema, { route: [] });
const routeDiscovery = create(Protobuf.Mesh.RouteDiscoverySchema, {
route: [],
});
const id = await client.sendPacket(
toBinary(Protobuf.Mesh.RouteDiscoverySchema, routeDiscovery),
Protobuf.Portnums.PortNum.TRACEROUTE_APP,
+95 -25
View File
@@ -9,12 +9,20 @@
import { create, toBinary } from "@bufbuild/protobuf";
import * as Protobuf from "@meshtastic/protobufs";
import type { Logger } from "tslog";
import { MeshClient, type MeshClientOptions } from "../core/client/MeshClient.ts";
import {
MeshClient,
type MeshClientOptions,
} from "../core/client/MeshClient.ts";
import type { EventBus } from "../core/event-bus/EventBus.ts";
import type { Queue } from "../core/queue/Queue.ts";
import type { Transport } from "../core/transport/Transport.ts";
import { DeviceStatusEnum } from "../core/transport/Transport.ts";
import { ChannelNumber, type Destination, Emitter, type PacketMetadata } from "../core/types.ts";
import {
ChannelNumber,
type Destination,
Emitter,
type PacketMetadata,
} from "../core/types.ts";
import type { Xmodem } from "../core/xmodem/Xmodem.ts";
import { sendAdminMessage } from "../features/device/infrastructure/AdminMessageSender.ts";
@@ -36,7 +44,10 @@ export class MeshDevice {
protected pendingSettingsChanges: boolean;
private myNodeInfo: Protobuf.Mesh.MyNodeInfo;
constructor(transport: Transport, configIdOrOptions?: number | MeshDeviceOptions) {
constructor(
transport: Transport,
configIdOrOptions?: number | MeshDeviceOptions,
) {
const options: MeshClientOptions =
typeof configIdOrOptions === "number"
? { transport, configId: configIdOrOptions }
@@ -56,8 +67,10 @@ export class MeshDevice {
this.client.events.onDeviceStatus.subscribe((status) => {
this.deviceStatus = status;
if (status === DeviceStatusEnum.DeviceConfigured) this.isConfigured = true;
else if (status === DeviceStatusEnum.DeviceConfiguring) this.isConfigured = false;
if (status === DeviceStatusEnum.DeviceConfigured)
this.isConfigured = true;
else if (status === DeviceStatusEnum.DeviceConfiguring)
this.isConfigured = false;
});
this.client.events.onMyNodeInfo.subscribe((info) => {
this.myNodeInfo = info;
@@ -164,8 +177,13 @@ export class MeshDevice {
return sendAdminMessage(this.client, { case: "setConfig", value: config });
}
public setModuleConfig(config: Protobuf.ModuleConfig.ModuleConfig): Promise<number> {
return sendAdminMessage(this.client, { case: "setModuleConfig", value: config });
public setModuleConfig(
config: Protobuf.ModuleConfig.ModuleConfig,
): Promise<number> {
return sendAdminMessage(this.client, {
case: "setModuleConfig",
value: config,
});
}
public setCannedMessages(
@@ -182,11 +200,17 @@ export class MeshDevice {
}
public setChannel(channel: Protobuf.Channel.Channel): Promise<number> {
return sendAdminMessage(this.client, { case: "setChannel", value: channel });
return sendAdminMessage(this.client, {
case: "setChannel",
value: channel,
});
}
public enterDfuMode(): Promise<number> {
return sendAdminMessage(this.client, { case: "enterDfuModeRequest", value: true });
return sendAdminMessage(this.client, {
case: "enterDfuModeRequest",
value: true,
});
}
public setPosition(position: Protobuf.Mesh.Position): Promise<number> {
@@ -197,7 +221,10 @@ export class MeshDevice {
);
}
public setFixedPosition(latitude: number, longitude: number): Promise<number> {
public setFixedPosition(
latitude: number,
longitude: number,
): Promise<number> {
const position = create(Protobuf.Mesh.PositionSchema, {
latitudeI: Math.floor(latitude / 1e-7),
longitudeI: Math.floor(longitude / 1e-7),
@@ -224,19 +251,35 @@ export class MeshDevice {
}
public getChannel(index: number): Promise<number> {
return sendAdminMessage(this.client, { case: "getChannelRequest", value: index + 1 });
return sendAdminMessage(this.client, {
case: "getChannelRequest",
value: index + 1,
});
}
public getConfig(type: Protobuf.Admin.AdminMessage_ConfigType): Promise<number> {
return sendAdminMessage(this.client, { case: "getConfigRequest", value: type });
public getConfig(
type: Protobuf.Admin.AdminMessage_ConfigType,
): Promise<number> {
return sendAdminMessage(this.client, {
case: "getConfigRequest",
value: type,
});
}
public getModuleConfig(type: Protobuf.Admin.AdminMessage_ModuleConfigType): Promise<number> {
return sendAdminMessage(this.client, { case: "getModuleConfigRequest", value: type });
public getModuleConfig(
type: Protobuf.Admin.AdminMessage_ModuleConfigType,
): Promise<number> {
return sendAdminMessage(this.client, {
case: "getModuleConfigRequest",
value: type,
});
}
public getOwner(): Promise<number> {
return sendAdminMessage(this.client, { case: "getOwnerRequest", value: true });
return sendAdminMessage(this.client, {
case: "getOwnerRequest",
value: true,
});
}
public getMetadata(nodeNum: number): Promise<number> {
@@ -253,12 +296,18 @@ export class MeshDevice {
index,
role: Protobuf.Channel.Channel_Role.DISABLED,
});
return sendAdminMessage(this.client, { case: "setChannel", value: channel });
return sendAdminMessage(this.client, {
case: "setChannel",
value: channel,
});
}
public commitEditSettings(): Promise<number> {
this.events.onPendingSettingsChange.dispatch(false);
return sendAdminMessage(this.client, { case: "commitEditSettings", value: true });
return sendAdminMessage(this.client, {
case: "commitEditSettings",
value: true,
});
}
public resetNodes(): Promise<number> {
@@ -266,27 +315,45 @@ export class MeshDevice {
}
public removeNodeByNum(nodeNum: number): Promise<number> {
return sendAdminMessage(this.client, { case: "removeByNodenum", value: nodeNum });
return sendAdminMessage(this.client, {
case: "removeByNodenum",
value: nodeNum,
});
}
public shutdown(time: number): Promise<number> {
return sendAdminMessage(this.client, { case: "shutdownSeconds", value: time });
return sendAdminMessage(this.client, {
case: "shutdownSeconds",
value: time,
});
}
public reboot(time: number): Promise<number> {
return sendAdminMessage(this.client, { case: "rebootSeconds", value: time });
return sendAdminMessage(this.client, {
case: "rebootSeconds",
value: time,
});
}
public rebootOta(time: number): Promise<number> {
return sendAdminMessage(this.client, { case: "rebootOtaSeconds", value: time });
return sendAdminMessage(this.client, {
case: "rebootOtaSeconds",
value: time,
});
}
public factoryResetDevice(): Promise<number> {
return sendAdminMessage(this.client, { case: "factoryResetDevice", value: 1 });
return sendAdminMessage(this.client, {
case: "factoryResetDevice",
value: 1,
});
}
public factoryResetConfig(): Promise<number> {
return sendAdminMessage(this.client, { case: "factoryResetConfig", value: 1 });
return sendAdminMessage(this.client, {
case: "factoryResetConfig",
value: 1,
});
}
public traceRoute(destination: number): Promise<number> {
@@ -318,7 +385,10 @@ export class MeshDevice {
}
/** Exposes an optimistic packet echo (was called by sendPacket echoResponse). */
public echoLocalPacket<T>(metadata: Omit<PacketMetadata<T>, "data">, data: T): void {
public echoLocalPacket<T>(
metadata: Omit<PacketMetadata<T>, "data">,
data: T,
): void {
void metadata;
void data;
}
+6 -1
View File
@@ -17,4 +17,9 @@ export type {
QueueItem,
Transport,
} from "../core/types.ts";
export { ChannelNumber, DeviceStatusEnum, Emitter, EmitterScope } from "../core/types.ts";
export {
ChannelNumber,
DeviceStatusEnum,
Emitter,
EmitterScope,
} from "../core/types.ts";