refactor: switch to using Bun (#718)
This commit is contained in:
@@ -5,6 +5,6 @@ const broadcastNum = 0xffffffff;
|
||||
const minFwVer = 2.2;
|
||||
|
||||
export const Constants = {
|
||||
broadcastNum,
|
||||
minFwVer,
|
||||
broadcastNum,
|
||||
minFwVer,
|
||||
};
|
||||
|
||||
+1150
-1173
File diff suppressed because it is too large
Load Diff
+84
-84
@@ -1,45 +1,45 @@
|
||||
import type * as Protobuf from "@meshtastic/protobufs";
|
||||
|
||||
interface Packet {
|
||||
type: "packet";
|
||||
data: Uint8Array;
|
||||
type: "packet";
|
||||
data: Uint8Array;
|
||||
}
|
||||
|
||||
interface DebugLog {
|
||||
type: "debug";
|
||||
data: string;
|
||||
type: "debug";
|
||||
data: string;
|
||||
}
|
||||
|
||||
export type DeviceOutput = Packet | DebugLog;
|
||||
|
||||
export interface Transport {
|
||||
toDevice: WritableStream<Uint8Array>;
|
||||
fromDevice: ReadableStream<DeviceOutput>;
|
||||
toDevice: WritableStream<Uint8Array>;
|
||||
fromDevice: ReadableStream<DeviceOutput>;
|
||||
}
|
||||
|
||||
export interface QueueItem {
|
||||
id: number;
|
||||
data: Uint8Array;
|
||||
sent: boolean;
|
||||
added: Date;
|
||||
promise: Promise<number>;
|
||||
id: number;
|
||||
data: Uint8Array;
|
||||
sent: boolean;
|
||||
added: Date;
|
||||
promise: Promise<number>;
|
||||
}
|
||||
|
||||
export interface HttpRetryConfig {
|
||||
maxRetries: number;
|
||||
initialDelayMs: number;
|
||||
maxDelayMs: number;
|
||||
backoffFactor: number;
|
||||
maxRetries: number;
|
||||
initialDelayMs: number;
|
||||
maxDelayMs: number;
|
||||
backoffFactor: number;
|
||||
}
|
||||
|
||||
export enum DeviceStatusEnum {
|
||||
DeviceRestarting = 1,
|
||||
DeviceDisconnected = 2,
|
||||
DeviceConnecting = 3,
|
||||
DeviceReconnecting = 4,
|
||||
DeviceConnected = 5,
|
||||
DeviceConfiguring = 6,
|
||||
DeviceConfigured = 7,
|
||||
DeviceRestarting = 1,
|
||||
DeviceDisconnected = 2,
|
||||
DeviceConnecting = 3,
|
||||
DeviceReconnecting = 4,
|
||||
DeviceConnected = 5,
|
||||
DeviceConfiguring = 6,
|
||||
DeviceConfigured = 7,
|
||||
}
|
||||
|
||||
export type LogEventPacket = LogEvent & { date: Date };
|
||||
@@ -47,83 +47,83 @@ export type LogEventPacket = LogEvent & { date: Date };
|
||||
export type PacketDestination = "broadcast" | "direct";
|
||||
|
||||
export interface PacketMetadata<T> {
|
||||
id: number;
|
||||
rxTime: Date;
|
||||
type: PacketDestination;
|
||||
from: number;
|
||||
to: number;
|
||||
channel: ChannelNumber;
|
||||
data: T;
|
||||
id: number;
|
||||
rxTime: Date;
|
||||
type: PacketDestination;
|
||||
from: number;
|
||||
to: number;
|
||||
channel: ChannelNumber;
|
||||
data: T;
|
||||
}
|
||||
|
||||
export enum EmitterScope {
|
||||
MeshDevice = 1,
|
||||
SerialConnection = 2,
|
||||
NodeSerialConnection = 3,
|
||||
BleConnection = 4,
|
||||
HttpConnection = 5,
|
||||
MeshDevice = 1,
|
||||
SerialConnection = 2,
|
||||
NodeSerialConnection = 3,
|
||||
BleConnection = 4,
|
||||
HttpConnection = 5,
|
||||
}
|
||||
|
||||
export enum Emitter {
|
||||
Constructor = 0,
|
||||
SendText = 1,
|
||||
SendWaypoint = 2,
|
||||
SendPacket = 3,
|
||||
SendRaw = 4,
|
||||
SetConfig = 5,
|
||||
SetModuleConfig = 6,
|
||||
ConfirmSetConfig = 7,
|
||||
SetOwner = 8,
|
||||
SetChannel = 9,
|
||||
ConfirmSetChannel = 10,
|
||||
ClearChannel = 11,
|
||||
GetChannel = 12,
|
||||
GetAllChannels = 13,
|
||||
GetConfig = 14,
|
||||
GetModuleConfig = 15,
|
||||
GetOwner = 16,
|
||||
Configure = 17,
|
||||
HandleFromRadio = 18,
|
||||
HandleMeshPacket = 19,
|
||||
Connect = 20,
|
||||
Ping = 21,
|
||||
ReadFromRadio = 22,
|
||||
WriteToRadio = 23,
|
||||
SetDebugMode = 24,
|
||||
GetMetadata = 25,
|
||||
ResetNodes = 26,
|
||||
Shutdown = 27,
|
||||
Reboot = 28,
|
||||
RebootOta = 29,
|
||||
FactoryReset = 30,
|
||||
EnterDfuMode = 31,
|
||||
RemoveNodeByNum = 32,
|
||||
SetCannedMessages = 33,
|
||||
Disconnect = 34,
|
||||
Constructor = 0,
|
||||
SendText = 1,
|
||||
SendWaypoint = 2,
|
||||
SendPacket = 3,
|
||||
SendRaw = 4,
|
||||
SetConfig = 5,
|
||||
SetModuleConfig = 6,
|
||||
ConfirmSetConfig = 7,
|
||||
SetOwner = 8,
|
||||
SetChannel = 9,
|
||||
ConfirmSetChannel = 10,
|
||||
ClearChannel = 11,
|
||||
GetChannel = 12,
|
||||
GetAllChannels = 13,
|
||||
GetConfig = 14,
|
||||
GetModuleConfig = 15,
|
||||
GetOwner = 16,
|
||||
Configure = 17,
|
||||
HandleFromRadio = 18,
|
||||
HandleMeshPacket = 19,
|
||||
Connect = 20,
|
||||
Ping = 21,
|
||||
ReadFromRadio = 22,
|
||||
WriteToRadio = 23,
|
||||
SetDebugMode = 24,
|
||||
GetMetadata = 25,
|
||||
ResetNodes = 26,
|
||||
Shutdown = 27,
|
||||
Reboot = 28,
|
||||
RebootOta = 29,
|
||||
FactoryReset = 30,
|
||||
EnterDfuMode = 31,
|
||||
RemoveNodeByNum = 32,
|
||||
SetCannedMessages = 33,
|
||||
Disconnect = 34,
|
||||
}
|
||||
|
||||
export interface LogEvent {
|
||||
scope: EmitterScope;
|
||||
emitter: Emitter;
|
||||
message: string;
|
||||
level: Protobuf.Mesh.LogRecord_Level;
|
||||
packet?: Uint8Array;
|
||||
scope: EmitterScope;
|
||||
emitter: Emitter;
|
||||
message: string;
|
||||
level: Protobuf.Mesh.LogRecord_Level;
|
||||
packet?: Uint8Array;
|
||||
}
|
||||
|
||||
export enum ChannelNumber {
|
||||
Primary = 0,
|
||||
Channel1 = 1,
|
||||
Channel2 = 2,
|
||||
Channel3 = 3,
|
||||
Channel4 = 4,
|
||||
Channel5 = 5,
|
||||
Channel6 = 6,
|
||||
Admin = 7,
|
||||
Primary = 0,
|
||||
Channel1 = 1,
|
||||
Channel2 = 2,
|
||||
Channel3 = 3,
|
||||
Channel4 = 4,
|
||||
Channel5 = 5,
|
||||
Channel6 = 6,
|
||||
Admin = 7,
|
||||
}
|
||||
|
||||
export type Destination = number | "self" | "broadcast";
|
||||
|
||||
export interface PacketError {
|
||||
id: number;
|
||||
error: Protobuf.Mesh.Routing_Error;
|
||||
id: number;
|
||||
error: Protobuf.Mesh.Routing_Error;
|
||||
}
|
||||
|
||||
@@ -1,467 +1,381 @@
|
||||
import { SimpleEventDispatcher } from "ste-simple-events";
|
||||
import type * as Protobuf from "@meshtastic/protobufs";
|
||||
import { SimpleEventDispatcher } from "ste-simple-events";
|
||||
import type { PacketMetadata } from "../types.ts";
|
||||
import type * as Types from "../types.ts";
|
||||
|
||||
export class EventSystem {
|
||||
/**
|
||||
* Fires when a new FromRadio message has been received from the device
|
||||
*
|
||||
* @event onLogEvent
|
||||
*/
|
||||
public readonly onLogEvent: SimpleEventDispatcher<
|
||||
Types.LogEventPacket
|
||||
> = new SimpleEventDispatcher<
|
||||
Types.LogEventPacket
|
||||
>();
|
||||
/**
|
||||
* Fires when a new FromRadio message has been received from the device
|
||||
*
|
||||
* @event onLogEvent
|
||||
*/
|
||||
public readonly onLogEvent: SimpleEventDispatcher<Types.LogEventPacket> =
|
||||
new SimpleEventDispatcher<Types.LogEventPacket>();
|
||||
|
||||
/**
|
||||
* Fires when a new FromRadio message has been received from the device
|
||||
*
|
||||
* @event onFromRadio
|
||||
*/
|
||||
public readonly onFromRadio: SimpleEventDispatcher<
|
||||
Protobuf.Mesh.FromRadio
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.Mesh.FromRadio
|
||||
>();
|
||||
/**
|
||||
* Fires when a new FromRadio message has been received from the device
|
||||
*
|
||||
* @event onFromRadio
|
||||
*/
|
||||
public readonly onFromRadio: SimpleEventDispatcher<Protobuf.Mesh.FromRadio> =
|
||||
new SimpleEventDispatcher<Protobuf.Mesh.FromRadio>();
|
||||
|
||||
/**
|
||||
* Fires when a new FromRadio message containing a Data packet has been
|
||||
* received from the device
|
||||
*
|
||||
* @event onMeshPacket
|
||||
*/
|
||||
public readonly onMeshPacket: SimpleEventDispatcher<
|
||||
Protobuf.Mesh.MeshPacket
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.Mesh.MeshPacket
|
||||
>();
|
||||
/**
|
||||
* Fires when a new FromRadio message containing a Data packet has been
|
||||
* received from the device
|
||||
*
|
||||
* @event onMeshPacket
|
||||
*/
|
||||
public readonly onMeshPacket: SimpleEventDispatcher<Protobuf.Mesh.MeshPacket> =
|
||||
new SimpleEventDispatcher<Protobuf.Mesh.MeshPacket>();
|
||||
|
||||
/**
|
||||
* Fires when a new MyNodeInfo message has been received from the device
|
||||
*
|
||||
* @event onMyNodeInfo
|
||||
*/
|
||||
public readonly onMyNodeInfo: SimpleEventDispatcher<
|
||||
Protobuf.Mesh.MyNodeInfo
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.Mesh.MyNodeInfo
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MyNodeInfo message has been received from the device
|
||||
*
|
||||
* @event onMyNodeInfo
|
||||
*/
|
||||
public readonly onMyNodeInfo: SimpleEventDispatcher<Protobuf.Mesh.MyNodeInfo> =
|
||||
new SimpleEventDispatcher<Protobuf.Mesh.MyNodeInfo>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a NodeInfo packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onNodeInfoPacket
|
||||
*/
|
||||
public readonly onNodeInfoPacket: SimpleEventDispatcher<
|
||||
Protobuf.Mesh.NodeInfo
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.Mesh.NodeInfo
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a NodeInfo packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onNodeInfoPacket
|
||||
*/
|
||||
public readonly onNodeInfoPacket: SimpleEventDispatcher<Protobuf.Mesh.NodeInfo> =
|
||||
new SimpleEventDispatcher<Protobuf.Mesh.NodeInfo>();
|
||||
|
||||
/**
|
||||
* Fires when a new Channel message is received
|
||||
*
|
||||
* @event onChannelPacket
|
||||
*/
|
||||
public readonly onChannelPacket: SimpleEventDispatcher<
|
||||
Protobuf.Channel.Channel
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.Channel.Channel
|
||||
>();
|
||||
/**
|
||||
* Fires when a new Channel message is received
|
||||
*
|
||||
* @event onChannelPacket
|
||||
*/
|
||||
public readonly onChannelPacket: SimpleEventDispatcher<Protobuf.Channel.Channel> =
|
||||
new SimpleEventDispatcher<Protobuf.Channel.Channel>();
|
||||
|
||||
/**
|
||||
* Fires when a new Config message is received
|
||||
*
|
||||
* @event onConfigPacket
|
||||
*/
|
||||
public readonly onConfigPacket: SimpleEventDispatcher<
|
||||
Protobuf.Config.Config
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.Config.Config
|
||||
>();
|
||||
/**
|
||||
* Fires when a new Config message is received
|
||||
*
|
||||
* @event onConfigPacket
|
||||
*/
|
||||
public readonly onConfigPacket: SimpleEventDispatcher<Protobuf.Config.Config> =
|
||||
new SimpleEventDispatcher<Protobuf.Config.Config>();
|
||||
|
||||
/**
|
||||
* Fires when a new ModuleConfig message is received
|
||||
*
|
||||
* @event onModuleConfigPacket
|
||||
*/
|
||||
public readonly onModuleConfigPacket: SimpleEventDispatcher<
|
||||
Protobuf.ModuleConfig.ModuleConfig
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.ModuleConfig.ModuleConfig
|
||||
>();
|
||||
/**
|
||||
* Fires when a new ModuleConfig message is received
|
||||
*
|
||||
* @event onModuleConfigPacket
|
||||
*/
|
||||
public readonly onModuleConfigPacket: SimpleEventDispatcher<Protobuf.ModuleConfig.ModuleConfig> =
|
||||
new SimpleEventDispatcher<Protobuf.ModuleConfig.ModuleConfig>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a ATAK packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onAtakPacket
|
||||
*/
|
||||
public readonly onAtakPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a ATAK packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onAtakPacket
|
||||
*/
|
||||
public readonly onAtakPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Text packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onMessagePacket
|
||||
*/
|
||||
public readonly onMessagePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<string>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<string>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Text packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onMessagePacket
|
||||
*/
|
||||
public readonly onMessagePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<string>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<string>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Remote Hardware packet has
|
||||
* been received from device
|
||||
*
|
||||
* @event onRemoteHardwarePacket
|
||||
*/
|
||||
public readonly onRemoteHardwarePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.RemoteHardware.HardwareMessage>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.RemoteHardware.HardwareMessage>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Remote Hardware packet has
|
||||
* been received from device
|
||||
*
|
||||
* @event onRemoteHardwarePacket
|
||||
*/
|
||||
public readonly onRemoteHardwarePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.RemoteHardware.HardwareMessage>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.RemoteHardware.HardwareMessage>
|
||||
>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Position packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onPositionPacket
|
||||
*/
|
||||
public readonly onPositionPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Position>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Position>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Position packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onPositionPacket
|
||||
*/
|
||||
public readonly onPositionPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Position>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.Mesh.Position>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a User packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onUserPacket
|
||||
*/
|
||||
public readonly onUserPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.User>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.User>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a User packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onUserPacket
|
||||
*/
|
||||
public readonly onUserPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.User>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.Mesh.User>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Routing packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onRoutingPacket
|
||||
*/
|
||||
public readonly onRoutingPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Routing>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Routing>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Routing packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onRoutingPacket
|
||||
*/
|
||||
public readonly onRoutingPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Routing>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.Mesh.Routing>>();
|
||||
|
||||
/**
|
||||
* Fires when the device receives a Metadata packet
|
||||
*
|
||||
* @event onDeviceMetadataPacket
|
||||
*/
|
||||
public readonly onDeviceMetadataPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.DeviceMetadata>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.DeviceMetadata>
|
||||
>();
|
||||
/**
|
||||
* Fires when the device receives a Metadata packet
|
||||
*
|
||||
* @event onDeviceMetadataPacket
|
||||
*/
|
||||
public readonly onDeviceMetadataPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.DeviceMetadata>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.Mesh.DeviceMetadata>>();
|
||||
|
||||
/**
|
||||
* Fires when the device receives a Canned Message Module message packet
|
||||
*
|
||||
* @event onCannedMessageModulePacket
|
||||
*/
|
||||
public readonly onCannedMessageModulePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<string>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<string>
|
||||
>();
|
||||
/**
|
||||
* Fires when the device receives a Canned Message Module message packet
|
||||
*
|
||||
* @event onCannedMessageModulePacket
|
||||
*/
|
||||
public readonly onCannedMessageModulePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<string>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<string>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Waypoint packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onWaypointPacket
|
||||
*/
|
||||
public readonly onWaypointPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Waypoint>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Waypoint>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Waypoint packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onWaypointPacket
|
||||
*/
|
||||
public readonly onWaypointPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.Waypoint>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.Mesh.Waypoint>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing an Audio packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onAudioPacket
|
||||
*/
|
||||
public readonly onAudioPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing an Audio packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onAudioPacket
|
||||
*/
|
||||
public readonly onAudioPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Detection Sensor packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onDetectionSensorPacket
|
||||
*/
|
||||
public readonly onDetectionSensorPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Detection Sensor packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onDetectionSensorPacket
|
||||
*/
|
||||
public readonly onDetectionSensorPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Ping packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onPingPacket
|
||||
*/
|
||||
public readonly onPingPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Ping packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onPingPacket
|
||||
*/
|
||||
public readonly onPingPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a IP Tunnel packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onIpTunnelPacket
|
||||
*/
|
||||
public readonly onIpTunnelPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a IP Tunnel packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onIpTunnelPacket
|
||||
*/
|
||||
public readonly onIpTunnelPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Paxcounter packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onPaxcounterPacket
|
||||
*/
|
||||
public readonly onPaxcounterPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.PaxCount.Paxcount>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.PaxCount.Paxcount>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Paxcounter packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onPaxcounterPacket
|
||||
*/
|
||||
public readonly onPaxcounterPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.PaxCount.Paxcount>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.PaxCount.Paxcount>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Serial packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onSerialPacket
|
||||
*/
|
||||
public readonly onSerialPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Serial packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onSerialPacket
|
||||
*/
|
||||
public readonly onSerialPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Store and Forward packet
|
||||
* has been received from device
|
||||
*
|
||||
* @event onStoreForwardPacket
|
||||
*/
|
||||
public readonly onStoreForwardPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Store and Forward packet
|
||||
* has been received from device
|
||||
*
|
||||
* @event onStoreForwardPacket
|
||||
*/
|
||||
public readonly onStoreForwardPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Store and Forward packet
|
||||
* has been received from device
|
||||
*
|
||||
* @event onRangeTestPacket
|
||||
*/
|
||||
public readonly onRangeTestPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Store and Forward packet
|
||||
* has been received from device
|
||||
*
|
||||
* @event onRangeTestPacket
|
||||
*/
|
||||
public readonly onRangeTestPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Telemetry packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onTelemetryPacket
|
||||
*/
|
||||
public readonly onTelemetryPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Telemetry.Telemetry>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Telemetry.Telemetry>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Telemetry packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onTelemetryPacket
|
||||
*/
|
||||
public readonly onTelemetryPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Telemetry.Telemetry>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.Telemetry.Telemetry>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a ZPS packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onZPSPacket
|
||||
*/
|
||||
public readonly onZpsPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a ZPS packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onZPSPacket
|
||||
*/
|
||||
public readonly onZpsPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Simulator packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onSimulatorPacket
|
||||
*/
|
||||
public readonly onSimulatorPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Simulator packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onSimulatorPacket
|
||||
*/
|
||||
public readonly onSimulatorPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Trace Route packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onTraceRoutePacket
|
||||
*/
|
||||
public readonly onTraceRoutePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.RouteDiscovery>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.RouteDiscovery>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Trace Route packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onTraceRoutePacket
|
||||
*/
|
||||
public readonly onTraceRoutePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.RouteDiscovery>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.Mesh.RouteDiscovery>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Neighbor Info packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onNeighborInfoPacket
|
||||
*/
|
||||
public readonly onNeighborInfoPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.NeighborInfo>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.NeighborInfo>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Neighbor Info packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onNeighborInfoPacket
|
||||
*/
|
||||
public readonly onNeighborInfoPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Protobuf.Mesh.NeighborInfo>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Protobuf.Mesh.NeighborInfo>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing an ATAK packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onAtakPluginPacket
|
||||
*/
|
||||
public readonly onAtakPluginPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing an ATAK packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onAtakPluginPacket
|
||||
*/
|
||||
public readonly onAtakPluginPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Map Report packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onMapReportPacket
|
||||
*/
|
||||
public readonly onMapReportPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Map Report packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onMapReportPacket
|
||||
*/
|
||||
public readonly onMapReportPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Private packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onPrivatePacket
|
||||
*/
|
||||
public readonly onPrivatePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing a Private packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onPrivatePacket
|
||||
*/
|
||||
public readonly onPrivatePacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing an ATAK Forwarder packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onAtakForwarderPacket
|
||||
*/
|
||||
public readonly onAtakForwarderPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
>();
|
||||
/**
|
||||
* Fires when a new MeshPacket message containing an ATAK Forwarder packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onAtakForwarderPacket
|
||||
*/
|
||||
public readonly onAtakForwarderPacket: SimpleEventDispatcher<
|
||||
PacketMetadata<Uint8Array>
|
||||
> = new SimpleEventDispatcher<PacketMetadata<Uint8Array>>();
|
||||
|
||||
/**
|
||||
* Fires when the devices connection or configuration status changes
|
||||
*
|
||||
* @event onDeviceStatus
|
||||
*/
|
||||
public readonly onDeviceStatus: SimpleEventDispatcher<
|
||||
Types.DeviceStatusEnum
|
||||
> = new SimpleEventDispatcher<
|
||||
Types.DeviceStatusEnum
|
||||
>();
|
||||
/**
|
||||
* Fires when the devices connection or configuration status changes
|
||||
*
|
||||
* @event onDeviceStatus
|
||||
*/
|
||||
public readonly onDeviceStatus: SimpleEventDispatcher<Types.DeviceStatusEnum> =
|
||||
new SimpleEventDispatcher<Types.DeviceStatusEnum>();
|
||||
|
||||
/**
|
||||
* Fires when a new FromRadio message containing a LogRecord packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onLogRecord
|
||||
*/
|
||||
public readonly onLogRecord: SimpleEventDispatcher<
|
||||
Protobuf.Mesh.LogRecord
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.Mesh.LogRecord
|
||||
>();
|
||||
/**
|
||||
* Fires when a new FromRadio message containing a LogRecord packet has been
|
||||
* received from device
|
||||
*
|
||||
* @event onLogRecord
|
||||
*/
|
||||
public readonly onLogRecord: SimpleEventDispatcher<Protobuf.Mesh.LogRecord> =
|
||||
new SimpleEventDispatcher<Protobuf.Mesh.LogRecord>();
|
||||
|
||||
/**
|
||||
* Fires when the device receives a meshPacket, returns a timestamp
|
||||
*
|
||||
* @event onMeshHeartbeat
|
||||
*/
|
||||
public readonly onMeshHeartbeat: SimpleEventDispatcher<Date> =
|
||||
new SimpleEventDispatcher<Date>();
|
||||
/**
|
||||
* Fires when the device receives a meshPacket, returns a timestamp
|
||||
*
|
||||
* @event onMeshHeartbeat
|
||||
*/
|
||||
public readonly onMeshHeartbeat: SimpleEventDispatcher<Date> =
|
||||
new SimpleEventDispatcher<Date>();
|
||||
|
||||
/**
|
||||
* Outputs any debug log data (currently serial connections only)
|
||||
*
|
||||
* @event onDeviceDebugLog
|
||||
*/
|
||||
public readonly onDeviceDebugLog: SimpleEventDispatcher<Uint8Array> =
|
||||
new SimpleEventDispatcher<Uint8Array>();
|
||||
/**
|
||||
* Outputs any debug log data (currently serial connections only)
|
||||
*
|
||||
* @event onDeviceDebugLog
|
||||
*/
|
||||
public readonly onDeviceDebugLog: SimpleEventDispatcher<Uint8Array> =
|
||||
new SimpleEventDispatcher<Uint8Array>();
|
||||
|
||||
/**
|
||||
* Outputs status of pending settings changes
|
||||
*
|
||||
* @event onpendingSettingsChange
|
||||
*/
|
||||
public readonly onPendingSettingsChange: SimpleEventDispatcher<
|
||||
boolean
|
||||
> = new SimpleEventDispatcher<
|
||||
boolean
|
||||
>();
|
||||
/**
|
||||
* Outputs status of pending settings changes
|
||||
*
|
||||
* @event onpendingSettingsChange
|
||||
*/
|
||||
public readonly onPendingSettingsChange: SimpleEventDispatcher<boolean> =
|
||||
new SimpleEventDispatcher<boolean>();
|
||||
|
||||
/**
|
||||
* Fires when a QueueStatus message is generated
|
||||
*
|
||||
* @event onQueueStatus
|
||||
*/
|
||||
public readonly onQueueStatus: SimpleEventDispatcher<
|
||||
Protobuf.Mesh.QueueStatus
|
||||
> = new SimpleEventDispatcher<
|
||||
Protobuf.Mesh.QueueStatus
|
||||
>();
|
||||
/**
|
||||
* Fires when a QueueStatus message is generated
|
||||
*
|
||||
* @event onQueueStatus
|
||||
*/
|
||||
public readonly onQueueStatus: SimpleEventDispatcher<Protobuf.Mesh.QueueStatus> =
|
||||
new SimpleEventDispatcher<Protobuf.Mesh.QueueStatus>();
|
||||
}
|
||||
|
||||
+101
-101
@@ -1,119 +1,119 @@
|
||||
import { SimpleEventDispatcher } from "ste-simple-events";
|
||||
import { fromBinary } from "@bufbuild/protobuf";
|
||||
import * as Protobuf from "@meshtastic/protobufs";
|
||||
import { SimpleEventDispatcher } from "ste-simple-events";
|
||||
import type { PacketError, QueueItem } from "../types.ts";
|
||||
|
||||
export class Queue {
|
||||
private queue: QueueItem[] = [];
|
||||
private lock = false;
|
||||
private ackNotifier = new SimpleEventDispatcher<number>();
|
||||
private errorNotifier = new SimpleEventDispatcher<PacketError>();
|
||||
private timeout: number;
|
||||
private queue: QueueItem[] = [];
|
||||
private lock = false;
|
||||
private ackNotifier = new SimpleEventDispatcher<number>();
|
||||
private errorNotifier = new SimpleEventDispatcher<PacketError>();
|
||||
private timeout: number;
|
||||
|
||||
constructor() {
|
||||
this.timeout = 60000;
|
||||
}
|
||||
constructor() {
|
||||
this.timeout = 60000;
|
||||
}
|
||||
|
||||
public getState(): QueueItem[] {
|
||||
return this.queue;
|
||||
}
|
||||
public getState(): QueueItem[] {
|
||||
return this.queue;
|
||||
}
|
||||
|
||||
public clear(): void {
|
||||
this.queue = [];
|
||||
}
|
||||
public clear(): void {
|
||||
this.queue = [];
|
||||
}
|
||||
|
||||
public push(item: Omit<QueueItem, "promise" | "sent" | "added">): void {
|
||||
const queueItem: QueueItem = {
|
||||
...item,
|
||||
sent: false,
|
||||
added: new Date(),
|
||||
promise: new Promise<number>((resolve, reject) => {
|
||||
this.ackNotifier.subscribe((id) => {
|
||||
if (item.id === id) {
|
||||
this.remove(item.id);
|
||||
resolve(id);
|
||||
}
|
||||
});
|
||||
this.errorNotifier.subscribe((e) => {
|
||||
if (item.id === e.id) {
|
||||
this.remove(item.id);
|
||||
reject(e);
|
||||
}
|
||||
});
|
||||
setTimeout(() => {
|
||||
if (this.queue.findIndex((qi) => qi.id === item.id) !== -1) {
|
||||
this.remove(item.id);
|
||||
const decoded = fromBinary(Protobuf.Mesh.ToRadioSchema, item.data);
|
||||
console.warn(
|
||||
`Packet ${item.id} of type ${decoded.payloadVariant.case} timed out`,
|
||||
);
|
||||
public push(item: Omit<QueueItem, "promise" | "sent" | "added">): void {
|
||||
const queueItem: QueueItem = {
|
||||
...item,
|
||||
sent: false,
|
||||
added: new Date(),
|
||||
promise: new Promise<number>((resolve, reject) => {
|
||||
this.ackNotifier.subscribe((id) => {
|
||||
if (item.id === id) {
|
||||
this.remove(item.id);
|
||||
resolve(id);
|
||||
}
|
||||
});
|
||||
this.errorNotifier.subscribe((e) => {
|
||||
if (item.id === e.id) {
|
||||
this.remove(item.id);
|
||||
reject(e);
|
||||
}
|
||||
});
|
||||
setTimeout(() => {
|
||||
if (this.queue.findIndex((qi) => qi.id === item.id) !== -1) {
|
||||
this.remove(item.id);
|
||||
const decoded = fromBinary(Protobuf.Mesh.ToRadioSchema, item.data);
|
||||
console.warn(
|
||||
`Packet ${item.id} of type ${decoded.payloadVariant.case} timed out`,
|
||||
);
|
||||
|
||||
reject({
|
||||
id: item.id,
|
||||
error: Protobuf.Mesh.Routing_Error.TIMEOUT,
|
||||
});
|
||||
}
|
||||
}, this.timeout);
|
||||
}),
|
||||
};
|
||||
this.queue.push(queueItem);
|
||||
}
|
||||
reject({
|
||||
id: item.id,
|
||||
error: Protobuf.Mesh.Routing_Error.TIMEOUT,
|
||||
});
|
||||
}
|
||||
}, this.timeout);
|
||||
}),
|
||||
};
|
||||
this.queue.push(queueItem);
|
||||
}
|
||||
|
||||
public remove(id: number): void {
|
||||
if (this.lock) {
|
||||
setTimeout(() => this.remove(id), 100);
|
||||
return;
|
||||
}
|
||||
this.queue = this.queue.filter((item) => item.id !== id);
|
||||
}
|
||||
public remove(id: number): void {
|
||||
if (this.lock) {
|
||||
setTimeout(() => this.remove(id), 100);
|
||||
return;
|
||||
}
|
||||
this.queue = this.queue.filter((item) => item.id !== id);
|
||||
}
|
||||
|
||||
public processAck(id: number): void {
|
||||
this.ackNotifier.dispatch(id);
|
||||
}
|
||||
public processAck(id: number): void {
|
||||
this.ackNotifier.dispatch(id);
|
||||
}
|
||||
|
||||
public processError(e: PacketError): void {
|
||||
console.error(
|
||||
`Error received for packet ${e.id}: ${
|
||||
Protobuf.Mesh.Routing_Error[e.error]
|
||||
}`,
|
||||
);
|
||||
this.errorNotifier.dispatch(e);
|
||||
}
|
||||
public processError(e: PacketError): void {
|
||||
console.error(
|
||||
`Error received for packet ${e.id}: ${
|
||||
Protobuf.Mesh.Routing_Error[e.error]
|
||||
}`,
|
||||
);
|
||||
this.errorNotifier.dispatch(e);
|
||||
}
|
||||
|
||||
public wait(id: number): Promise<number> {
|
||||
const queueItem = this.queue.find((qi) => qi.id === id);
|
||||
if (!queueItem) {
|
||||
throw new Error("Packet does not exist");
|
||||
}
|
||||
return queueItem.promise;
|
||||
}
|
||||
public wait(id: number): Promise<number> {
|
||||
const queueItem = this.queue.find((qi) => qi.id === id);
|
||||
if (!queueItem) {
|
||||
throw new Error("Packet does not exist");
|
||||
}
|
||||
return queueItem.promise;
|
||||
}
|
||||
|
||||
public async processQueue(
|
||||
outputStream: WritableStream<Uint8Array>,
|
||||
): Promise<void> {
|
||||
if (this.lock) {
|
||||
return;
|
||||
}
|
||||
public async processQueue(
|
||||
outputStream: WritableStream<Uint8Array>,
|
||||
): Promise<void> {
|
||||
if (this.lock) {
|
||||
return;
|
||||
}
|
||||
|
||||
this.lock = true;
|
||||
const writer = outputStream.getWriter();
|
||||
this.lock = true;
|
||||
const writer = outputStream.getWriter();
|
||||
|
||||
try {
|
||||
while (this.queue.filter((p) => !p.sent).length > 0) {
|
||||
const item = this.queue.filter((p) => !p.sent)[0];
|
||||
if (item) {
|
||||
await new Promise((resolve) => setTimeout(resolve, 200));
|
||||
try {
|
||||
await writer.write(item.data);
|
||||
item.sent = true;
|
||||
} catch (error) {
|
||||
console.error(`Error sending packet ${item.id}`, error);
|
||||
}
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
writer.releaseLock();
|
||||
this.lock = false;
|
||||
}
|
||||
}
|
||||
try {
|
||||
while (this.queue.filter((p) => !p.sent).length > 0) {
|
||||
const item = this.queue.filter((p) => !p.sent)[0];
|
||||
if (item) {
|
||||
await new Promise((resolve) => setTimeout(resolve, 200));
|
||||
try {
|
||||
await writer.write(item.data);
|
||||
item.sent = true;
|
||||
} catch (error) {
|
||||
console.error(`Error sending packet ${item.id}`, error);
|
||||
}
|
||||
}
|
||||
}
|
||||
} finally {
|
||||
writer.releaseLock();
|
||||
this.lock = false;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,222 +1,222 @@
|
||||
import { fromBinary } from "@bufbuild/protobuf";
|
||||
import type { DeviceOutput } from "../../types.ts";
|
||||
import { Constants, Protobuf, Types } from "../../../mod.ts";
|
||||
import type { MeshDevice } from "../../../mod.ts";
|
||||
import type { DeviceOutput } from "../../types.ts";
|
||||
|
||||
export const decodePacket = (device: MeshDevice) =>
|
||||
new WritableStream<DeviceOutput>({
|
||||
write(chunk) {
|
||||
switch (chunk.type) {
|
||||
case "debug": {
|
||||
break;
|
||||
}
|
||||
case "packet": {
|
||||
const decodedMessage = fromBinary(
|
||||
Protobuf.Mesh.FromRadioSchema,
|
||||
chunk.data,
|
||||
);
|
||||
device.events.onFromRadio.dispatch(decodedMessage);
|
||||
new WritableStream<DeviceOutput>({
|
||||
write(chunk) {
|
||||
switch (chunk.type) {
|
||||
case "debug": {
|
||||
break;
|
||||
}
|
||||
case "packet": {
|
||||
const decodedMessage = fromBinary(
|
||||
Protobuf.Mesh.FromRadioSchema,
|
||||
chunk.data,
|
||||
);
|
||||
device.events.onFromRadio.dispatch(decodedMessage);
|
||||
|
||||
/** @todo Add map here when `all=true` gets fixed. */
|
||||
switch (decodedMessage.payloadVariant.case) {
|
||||
case "packet": {
|
||||
device.handleMeshPacket(decodedMessage.payloadVariant.value);
|
||||
break;
|
||||
}
|
||||
/** @todo Add map here when `all=true` gets fixed. */
|
||||
switch (decodedMessage.payloadVariant.case) {
|
||||
case "packet": {
|
||||
device.handleMeshPacket(decodedMessage.payloadVariant.value);
|
||||
break;
|
||||
}
|
||||
|
||||
case "myInfo": {
|
||||
device.events.onMyNodeInfo.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
device.log.info(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
"📱 Received Node info for this device",
|
||||
);
|
||||
break;
|
||||
}
|
||||
case "myInfo": {
|
||||
device.events.onMyNodeInfo.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
device.log.info(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
"📱 Received Node info for this device",
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
case "nodeInfo": {
|
||||
device.log.info(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`📱 Received Node Info packet for node: ${decodedMessage.payloadVariant.value.num}`,
|
||||
);
|
||||
case "nodeInfo": {
|
||||
device.log.info(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`📱 Received Node Info packet for node: ${decodedMessage.payloadVariant.value.num}`,
|
||||
);
|
||||
|
||||
device.events.onNodeInfoPacket.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
device.events.onNodeInfoPacket.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
|
||||
//TODO: HERE
|
||||
if (decodedMessage.payloadVariant.value.position) {
|
||||
device.events.onPositionPacket.dispatch({
|
||||
id: decodedMessage.id,
|
||||
rxTime: new Date(),
|
||||
from: decodedMessage.payloadVariant.value.num,
|
||||
to: decodedMessage.payloadVariant.value.num,
|
||||
type: "direct",
|
||||
channel: Types.ChannelNumber.Primary,
|
||||
data: decodedMessage.payloadVariant.value.position,
|
||||
});
|
||||
}
|
||||
//TODO: HERE
|
||||
if (decodedMessage.payloadVariant.value.position) {
|
||||
device.events.onPositionPacket.dispatch({
|
||||
id: decodedMessage.id,
|
||||
rxTime: new Date(),
|
||||
from: decodedMessage.payloadVariant.value.num,
|
||||
to: decodedMessage.payloadVariant.value.num,
|
||||
type: "direct",
|
||||
channel: Types.ChannelNumber.Primary,
|
||||
data: decodedMessage.payloadVariant.value.position,
|
||||
});
|
||||
}
|
||||
|
||||
//TODO: HERE
|
||||
if (decodedMessage.payloadVariant.value.user) {
|
||||
device.events.onUserPacket.dispatch({
|
||||
id: decodedMessage.id,
|
||||
rxTime: new Date(),
|
||||
from: decodedMessage.payloadVariant.value.num,
|
||||
to: decodedMessage.payloadVariant.value.num,
|
||||
type: "direct",
|
||||
channel: Types.ChannelNumber.Primary,
|
||||
data: decodedMessage.payloadVariant.value.user,
|
||||
});
|
||||
}
|
||||
break;
|
||||
}
|
||||
//TODO: HERE
|
||||
if (decodedMessage.payloadVariant.value.user) {
|
||||
device.events.onUserPacket.dispatch({
|
||||
id: decodedMessage.id,
|
||||
rxTime: new Date(),
|
||||
from: decodedMessage.payloadVariant.value.num,
|
||||
to: decodedMessage.payloadVariant.value.num,
|
||||
type: "direct",
|
||||
channel: Types.ChannelNumber.Primary,
|
||||
data: decodedMessage.payloadVariant.value.user,
|
||||
});
|
||||
}
|
||||
break;
|
||||
}
|
||||
|
||||
case "config": {
|
||||
if (decodedMessage.payloadVariant.value.payloadVariant.case) {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`💾 Received Config packet of variant: ${decodedMessage.payloadVariant.value.payloadVariant.case}`,
|
||||
);
|
||||
} else {
|
||||
device.log.warn(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`⚠️ Received Config packet of variant: ${"UNK"}`,
|
||||
);
|
||||
}
|
||||
case "config": {
|
||||
if (decodedMessage.payloadVariant.value.payloadVariant.case) {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`💾 Received Config packet of variant: ${decodedMessage.payloadVariant.value.payloadVariant.case}`,
|
||||
);
|
||||
} else {
|
||||
device.log.warn(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`⚠️ Received Config packet of variant: ${"UNK"}`,
|
||||
);
|
||||
}
|
||||
|
||||
device.events.onConfigPacket.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
device.events.onConfigPacket.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
case "logRecord": {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
"Received onLogRecord",
|
||||
);
|
||||
device.events.onLogRecord.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
case "logRecord": {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
"Received onLogRecord",
|
||||
);
|
||||
device.events.onLogRecord.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
case "configCompleteId": {
|
||||
if (decodedMessage.payloadVariant.value !== device.configId) {
|
||||
device.log.error(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`❌ Invalid config id received from device, expected ${device.configId} but received ${decodedMessage.payloadVariant.value}`,
|
||||
);
|
||||
}
|
||||
case "configCompleteId": {
|
||||
if (decodedMessage.payloadVariant.value !== device.configId) {
|
||||
device.log.error(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`❌ Invalid config id received from device, expected ${device.configId} but received ${decodedMessage.payloadVariant.value}`,
|
||||
);
|
||||
}
|
||||
|
||||
device.log.info(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`⚙️ Valid config id received from device: ${device.configId}`,
|
||||
);
|
||||
device.log.info(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`⚙️ Valid config id received from device: ${device.configId}`,
|
||||
);
|
||||
|
||||
device.updateDeviceStatus(
|
||||
Types.DeviceStatusEnum.DeviceConfigured,
|
||||
);
|
||||
break;
|
||||
}
|
||||
device.updateDeviceStatus(
|
||||
Types.DeviceStatusEnum.DeviceConfigured,
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
case "rebooted": {
|
||||
device.configure().catch(() => {
|
||||
// TODO: FIX, workaround for `wantConfigId` not getting acks.
|
||||
});
|
||||
break;
|
||||
}
|
||||
case "rebooted": {
|
||||
device.configure().catch(() => {
|
||||
// TODO: FIX, workaround for `wantConfigId` not getting acks.
|
||||
});
|
||||
break;
|
||||
}
|
||||
|
||||
case "moduleConfig": {
|
||||
if (decodedMessage.payloadVariant.value.payloadVariant.case) {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`💾 Received Module Config packet of variant: ${decodedMessage.payloadVariant.value.payloadVariant.case}`,
|
||||
);
|
||||
} else {
|
||||
device.log.warn(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
"⚠️ Received Module Config packet of variant: UNK",
|
||||
);
|
||||
}
|
||||
case "moduleConfig": {
|
||||
if (decodedMessage.payloadVariant.value.payloadVariant.case) {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`💾 Received Module Config packet of variant: ${decodedMessage.payloadVariant.value.payloadVariant.case}`,
|
||||
);
|
||||
} else {
|
||||
device.log.warn(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
"⚠️ Received Module Config packet of variant: UNK",
|
||||
);
|
||||
}
|
||||
|
||||
device.events.onModuleConfigPacket.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
device.events.onModuleConfigPacket.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
case "channel": {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`🔐 Received Channel: ${decodedMessage.payloadVariant.value.index}`,
|
||||
);
|
||||
case "channel": {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`🔐 Received Channel: ${decodedMessage.payloadVariant.value.index}`,
|
||||
);
|
||||
|
||||
device.events.onChannelPacket.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
device.events.onChannelPacket.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
case "queueStatus": {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`🚧 Received Queue Status: ${decodedMessage.payloadVariant.value}`,
|
||||
);
|
||||
case "queueStatus": {
|
||||
device.log.trace(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`🚧 Received Queue Status: ${decodedMessage.payloadVariant.value}`,
|
||||
);
|
||||
|
||||
device.events.onQueueStatus.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
device.events.onQueueStatus.dispatch(
|
||||
decodedMessage.payloadVariant.value,
|
||||
);
|
||||
break;
|
||||
}
|
||||
|
||||
case "xmodemPacket": {
|
||||
device.xModem.handlePacket(decodedMessage.payloadVariant.value);
|
||||
break;
|
||||
}
|
||||
case "xmodemPacket": {
|
||||
device.xModem.handlePacket(decodedMessage.payloadVariant.value);
|
||||
break;
|
||||
}
|
||||
|
||||
case "metadata": {
|
||||
if (
|
||||
Number.parseFloat(
|
||||
decodedMessage.payloadVariant.value.firmwareVersion,
|
||||
) < Constants.minFwVer
|
||||
) {
|
||||
device.log.fatal(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`Device firmware outdated. Min supported: ${Constants.minFwVer} got : ${decodedMessage.payloadVariant.value.firmwareVersion}`,
|
||||
);
|
||||
}
|
||||
device.log.debug(
|
||||
Types.Emitter[Types.Emitter.GetMetadata],
|
||||
"🏷️ Received metadata packet",
|
||||
);
|
||||
case "metadata": {
|
||||
if (
|
||||
Number.parseFloat(
|
||||
decodedMessage.payloadVariant.value.firmwareVersion,
|
||||
) < Constants.minFwVer
|
||||
) {
|
||||
device.log.fatal(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`Device firmware outdated. Min supported: ${Constants.minFwVer} got : ${decodedMessage.payloadVariant.value.firmwareVersion}`,
|
||||
);
|
||||
}
|
||||
device.log.debug(
|
||||
Types.Emitter[Types.Emitter.GetMetadata],
|
||||
"🏷️ Received metadata packet",
|
||||
);
|
||||
|
||||
device.events.onDeviceMetadataPacket.dispatch({
|
||||
id: decodedMessage.id,
|
||||
rxTime: new Date(),
|
||||
from: 0,
|
||||
to: 0,
|
||||
type: "direct",
|
||||
channel: Types.ChannelNumber.Primary,
|
||||
data: decodedMessage.payloadVariant.value,
|
||||
});
|
||||
break;
|
||||
}
|
||||
device.events.onDeviceMetadataPacket.dispatch({
|
||||
id: decodedMessage.id,
|
||||
rxTime: new Date(),
|
||||
from: 0,
|
||||
to: 0,
|
||||
type: "direct",
|
||||
channel: Types.ChannelNumber.Primary,
|
||||
data: decodedMessage.payloadVariant.value,
|
||||
});
|
||||
break;
|
||||
}
|
||||
|
||||
case "mqttClientProxyMessage": {
|
||||
break;
|
||||
}
|
||||
case "mqttClientProxyMessage": {
|
||||
break;
|
||||
}
|
||||
|
||||
default: {
|
||||
device.log.warn(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`⚠️ Unhandled payload variant: ${decodedMessage.payloadVariant.case}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
});
|
||||
default: {
|
||||
device.log.warn(
|
||||
Types.Emitter[Types.Emitter.HandleFromRadio],
|
||||
`⚠️ Unhandled payload variant: ${decodedMessage.payloadVariant.case}`,
|
||||
);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
});
|
||||
|
||||
@@ -1,73 +1,71 @@
|
||||
import type { DeviceOutput } from "../../types.ts";
|
||||
|
||||
export const fromDeviceStream: () => TransformStream<Uint8Array, DeviceOutput> =
|
||||
(
|
||||
// onReleaseEvent: SimpleEventDispatcher<boolean>,
|
||||
) => {
|
||||
let byteBuffer = new Uint8Array([]);
|
||||
const textDecoder = new TextDecoder();
|
||||
return new TransformStream<Uint8Array, DeviceOutput>({
|
||||
transform(chunk: Uint8Array, controller): void {
|
||||
// onReleaseEvent.subscribe(() => {
|
||||
// controller.terminate();
|
||||
// });
|
||||
byteBuffer = new Uint8Array([...byteBuffer, ...chunk]);
|
||||
let processingExhausted = false;
|
||||
while (byteBuffer.length !== 0 && !processingExhausted) {
|
||||
const framingIndex = byteBuffer.findIndex((byte) => byte === 0x94);
|
||||
const framingByte2 = byteBuffer[framingIndex + 1];
|
||||
if (framingByte2 === 0xc3) {
|
||||
if (byteBuffer.subarray(0, framingIndex).length) {
|
||||
controller.enqueue({
|
||||
type: "debug",
|
||||
data: textDecoder.decode(byteBuffer.subarray(0, framingIndex)),
|
||||
});
|
||||
byteBuffer = byteBuffer.subarray(framingIndex);
|
||||
}
|
||||
(
|
||||
// onReleaseEvent: SimpleEventDispatcher<boolean>,
|
||||
) => {
|
||||
let byteBuffer = new Uint8Array([]);
|
||||
const textDecoder = new TextDecoder();
|
||||
return new TransformStream<Uint8Array, DeviceOutput>({
|
||||
transform(chunk: Uint8Array, controller): void {
|
||||
// onReleaseEvent.subscribe(() => {
|
||||
// controller.terminate();
|
||||
// });
|
||||
byteBuffer = new Uint8Array([...byteBuffer, ...chunk]);
|
||||
let processingExhausted = false;
|
||||
while (byteBuffer.length !== 0 && !processingExhausted) {
|
||||
const framingIndex = byteBuffer.findIndex((byte) => byte === 0x94);
|
||||
const framingByte2 = byteBuffer[framingIndex + 1];
|
||||
if (framingByte2 === 0xc3) {
|
||||
if (byteBuffer.subarray(0, framingIndex).length) {
|
||||
controller.enqueue({
|
||||
type: "debug",
|
||||
data: textDecoder.decode(byteBuffer.subarray(0, framingIndex)),
|
||||
});
|
||||
byteBuffer = byteBuffer.subarray(framingIndex);
|
||||
}
|
||||
|
||||
const msb = byteBuffer[2];
|
||||
const lsb = byteBuffer[3];
|
||||
const msb = byteBuffer[2];
|
||||
const lsb = byteBuffer[3];
|
||||
|
||||
if (
|
||||
msb !== undefined &&
|
||||
lsb !== undefined &&
|
||||
byteBuffer.length >= 4 + (msb << 8) + lsb
|
||||
) {
|
||||
const packet = byteBuffer.subarray(4, 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.findIndex(
|
||||
(byte) => byte === 0x94,
|
||||
);
|
||||
if (
|
||||
malformedDetectorIndex !== -1 &&
|
||||
packet[malformedDetectorIndex + 1] === 0xc3
|
||||
) {
|
||||
console.warn(
|
||||
`⚠️ Malformed packet found, discarding: ${
|
||||
byteBuffer
|
||||
.subarray(0, malformedDetectorIndex - 1)
|
||||
.toString()
|
||||
}`,
|
||||
);
|
||||
const malformedDetectorIndex = packet.findIndex(
|
||||
(byte) => byte === 0x94,
|
||||
);
|
||||
if (
|
||||
malformedDetectorIndex !== -1 &&
|
||||
packet[malformedDetectorIndex + 1] === 0xc3
|
||||
) {
|
||||
console.warn(
|
||||
`⚠️ Malformed packet found, discarding: ${byteBuffer
|
||||
.subarray(0, malformedDetectorIndex - 1)
|
||||
.toString()}`,
|
||||
);
|
||||
|
||||
byteBuffer = byteBuffer.subarray(malformedDetectorIndex);
|
||||
} else {
|
||||
byteBuffer = byteBuffer.subarray(3 + (msb << 8) + lsb + 1);
|
||||
byteBuffer = byteBuffer.subarray(malformedDetectorIndex);
|
||||
} else {
|
||||
byteBuffer = byteBuffer.subarray(3 + (msb << 8) + lsb + 1);
|
||||
|
||||
controller.enqueue({
|
||||
type: "packet",
|
||||
data: packet,
|
||||
});
|
||||
}
|
||||
} else {
|
||||
/** Only partioal message in buffer, wait for the rest */
|
||||
processingExhausted = true;
|
||||
}
|
||||
} else {
|
||||
/** Message not complete, only 1 byte in buffer */
|
||||
processingExhausted = true;
|
||||
}
|
||||
}
|
||||
},
|
||||
});
|
||||
};
|
||||
controller.enqueue({
|
||||
type: "packet",
|
||||
data: packet,
|
||||
});
|
||||
}
|
||||
} else {
|
||||
/** Only partioal message in buffer, wait for the rest */
|
||||
processingExhausted = true;
|
||||
}
|
||||
} else {
|
||||
/** Message not complete, only 1 byte in buffer */
|
||||
processingExhausted = true;
|
||||
}
|
||||
}
|
||||
},
|
||||
});
|
||||
};
|
||||
|
||||
@@ -2,15 +2,15 @@
|
||||
* Pads packets with appropriate framing information before writing to the output stream.
|
||||
*/
|
||||
export const toDeviceStream: TransformStream<Uint8Array, Uint8Array> =
|
||||
new TransformStream<Uint8Array, Uint8Array>({
|
||||
transform(chunk: Uint8Array, controller): void {
|
||||
const bufLen = chunk.length;
|
||||
const header = new Uint8Array([
|
||||
0x94,
|
||||
0xC3,
|
||||
(bufLen >> 8) & 0xFF,
|
||||
bufLen & 0xFF,
|
||||
]);
|
||||
controller.enqueue(new Uint8Array([...header, ...chunk]));
|
||||
},
|
||||
});
|
||||
new TransformStream<Uint8Array, Uint8Array>({
|
||||
transform(chunk: Uint8Array, controller): void {
|
||||
const bufLen = chunk.length;
|
||||
const header = new Uint8Array([
|
||||
0x94,
|
||||
0xc3,
|
||||
(bufLen >> 8) & 0xff,
|
||||
bufLen & 0xff,
|
||||
]);
|
||||
controller.enqueue(new Uint8Array([...header, ...chunk]));
|
||||
},
|
||||
});
|
||||
|
||||
+117
-117
@@ -1,135 +1,135 @@
|
||||
import crc16ccitt from "crc/calculators/crc16ccitt";
|
||||
import { create, toBinary } from "@bufbuild/protobuf";
|
||||
import * as Protobuf from "@meshtastic/protobufs";
|
||||
import crc16ccitt from "crc/calculators/crc16ccitt";
|
||||
|
||||
//if counter > 35 then reset counter/clear/error/reject promise
|
||||
type XmodemProps = (toRadio: Uint8Array, id?: number) => Promise<number>;
|
||||
|
||||
export class Xmodem {
|
||||
private sendRaw: XmodemProps;
|
||||
private rxBuffer: Uint8Array[];
|
||||
private txBuffer: Uint8Array[];
|
||||
private textEncoder: TextEncoder;
|
||||
private counter: number;
|
||||
private sendRaw: XmodemProps;
|
||||
private rxBuffer: Uint8Array[];
|
||||
private txBuffer: Uint8Array[];
|
||||
private textEncoder: TextEncoder;
|
||||
private counter: number;
|
||||
|
||||
constructor(sendRaw: XmodemProps) {
|
||||
this.sendRaw = sendRaw;
|
||||
this.rxBuffer = [];
|
||||
this.txBuffer = [];
|
||||
this.textEncoder = new TextEncoder();
|
||||
this.counter = 0;
|
||||
}
|
||||
constructor(sendRaw: XmodemProps) {
|
||||
this.sendRaw = sendRaw;
|
||||
this.rxBuffer = [];
|
||||
this.txBuffer = [];
|
||||
this.textEncoder = new TextEncoder();
|
||||
this.counter = 0;
|
||||
}
|
||||
|
||||
async downloadFile(filename: string): Promise<number> {
|
||||
return await this.sendCommand(
|
||||
Protobuf.Xmodem.XModem_Control.STX,
|
||||
this.textEncoder.encode(filename),
|
||||
0,
|
||||
);
|
||||
}
|
||||
async downloadFile(filename: string): Promise<number> {
|
||||
return await this.sendCommand(
|
||||
Protobuf.Xmodem.XModem_Control.STX,
|
||||
this.textEncoder.encode(filename),
|
||||
0,
|
||||
);
|
||||
}
|
||||
|
||||
async uploadFile(filename: string, data: Uint8Array): Promise<number> {
|
||||
for (let i = 0; i < data.length; i += 128) {
|
||||
this.txBuffer.push(data.slice(i, i + 128));
|
||||
}
|
||||
async uploadFile(filename: string, data: Uint8Array): Promise<number> {
|
||||
for (let i = 0; i < data.length; i += 128) {
|
||||
this.txBuffer.push(data.slice(i, i + 128));
|
||||
}
|
||||
|
||||
return await this.sendCommand(
|
||||
Protobuf.Xmodem.XModem_Control.SOH,
|
||||
this.textEncoder.encode(filename),
|
||||
0,
|
||||
);
|
||||
}
|
||||
return await this.sendCommand(
|
||||
Protobuf.Xmodem.XModem_Control.SOH,
|
||||
this.textEncoder.encode(filename),
|
||||
0,
|
||||
);
|
||||
}
|
||||
|
||||
async sendCommand(
|
||||
command: Protobuf.Xmodem.XModem_Control,
|
||||
buffer?: Uint8Array,
|
||||
sequence?: number,
|
||||
crc16?: number,
|
||||
): Promise<number> {
|
||||
const toRadio = create(Protobuf.Mesh.ToRadioSchema, {
|
||||
payloadVariant: {
|
||||
case: "xmodemPacket",
|
||||
value: {
|
||||
buffer,
|
||||
control: command,
|
||||
seq: sequence,
|
||||
crc16: crc16,
|
||||
},
|
||||
},
|
||||
});
|
||||
return await this.sendRaw(toBinary(Protobuf.Mesh.ToRadioSchema, toRadio));
|
||||
}
|
||||
async sendCommand(
|
||||
command: Protobuf.Xmodem.XModem_Control,
|
||||
buffer?: Uint8Array,
|
||||
sequence?: number,
|
||||
crc16?: number,
|
||||
): Promise<number> {
|
||||
const toRadio = create(Protobuf.Mesh.ToRadioSchema, {
|
||||
payloadVariant: {
|
||||
case: "xmodemPacket",
|
||||
value: {
|
||||
buffer,
|
||||
control: command,
|
||||
seq: sequence,
|
||||
crc16: crc16,
|
||||
},
|
||||
},
|
||||
});
|
||||
return await this.sendRaw(toBinary(Protobuf.Mesh.ToRadioSchema, toRadio));
|
||||
}
|
||||
|
||||
async handlePacket(packet: Protobuf.Xmodem.XModem): Promise<number> {
|
||||
await new Promise((resolve) => setTimeout(resolve, 100));
|
||||
async handlePacket(packet: Protobuf.Xmodem.XModem): Promise<number> {
|
||||
await new Promise((resolve) => setTimeout(resolve, 100));
|
||||
|
||||
switch (packet.control) {
|
||||
case Protobuf.Xmodem.XModem_Control.NUL: {
|
||||
// nothing
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.SOH: {
|
||||
this.counter = packet.seq;
|
||||
if (this.validateCrc16(packet)) {
|
||||
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,
|
||||
);
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.STX: {
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.EOT: {
|
||||
// end of transmission
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.ACK: {
|
||||
this.counter++;
|
||||
if (this.txBuffer[this.counter - 1]) {
|
||||
return this.sendCommand(
|
||||
Protobuf.Xmodem.XModem_Control.SOH,
|
||||
this.txBuffer[this.counter - 1],
|
||||
this.counter,
|
||||
crc16ccitt(this.txBuffer[this.counter - 1] ?? new Uint8Array()),
|
||||
);
|
||||
}
|
||||
if (this.counter === this.txBuffer.length + 1) {
|
||||
return this.sendCommand(Protobuf.Xmodem.XModem_Control.EOT);
|
||||
}
|
||||
this.clear();
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.NAK: {
|
||||
return this.sendCommand(
|
||||
Protobuf.Xmodem.XModem_Control.SOH,
|
||||
this.txBuffer[this.counter],
|
||||
this.counter,
|
||||
crc16ccitt(this.txBuffer[this.counter - 1] ?? new Uint8Array()),
|
||||
);
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.CAN: {
|
||||
this.clear();
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.CTRLZ: {
|
||||
break;
|
||||
}
|
||||
}
|
||||
switch (packet.control) {
|
||||
case Protobuf.Xmodem.XModem_Control.NUL: {
|
||||
// nothing
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.SOH: {
|
||||
this.counter = packet.seq;
|
||||
if (this.validateCrc16(packet)) {
|
||||
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,
|
||||
);
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.STX: {
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.EOT: {
|
||||
// end of transmission
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.ACK: {
|
||||
this.counter++;
|
||||
if (this.txBuffer[this.counter - 1]) {
|
||||
return this.sendCommand(
|
||||
Protobuf.Xmodem.XModem_Control.SOH,
|
||||
this.txBuffer[this.counter - 1],
|
||||
this.counter,
|
||||
crc16ccitt(this.txBuffer[this.counter - 1] ?? new Uint8Array()),
|
||||
);
|
||||
}
|
||||
if (this.counter === this.txBuffer.length + 1) {
|
||||
return this.sendCommand(Protobuf.Xmodem.XModem_Control.EOT);
|
||||
}
|
||||
this.clear();
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.NAK: {
|
||||
return this.sendCommand(
|
||||
Protobuf.Xmodem.XModem_Control.SOH,
|
||||
this.txBuffer[this.counter],
|
||||
this.counter,
|
||||
crc16ccitt(this.txBuffer[this.counter - 1] ?? new Uint8Array()),
|
||||
);
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.CAN: {
|
||||
this.clear();
|
||||
break;
|
||||
}
|
||||
case Protobuf.Xmodem.XModem_Control.CTRLZ: {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
return Promise.resolve(0);
|
||||
}
|
||||
return Promise.resolve(0);
|
||||
}
|
||||
|
||||
validateCrc16(packet: Protobuf.Xmodem.XModem): boolean {
|
||||
return crc16ccitt(packet.buffer) === packet.crc16;
|
||||
}
|
||||
validateCrc16(packet: Protobuf.Xmodem.XModem): boolean {
|
||||
return crc16ccitt(packet.buffer) === packet.crc16;
|
||||
}
|
||||
|
||||
clear() {
|
||||
this.counter = 0;
|
||||
this.rxBuffer = [];
|
||||
this.txBuffer = [];
|
||||
}
|
||||
clear() {
|
||||
this.counter = 0;
|
||||
this.rxBuffer = [];
|
||||
this.txBuffer = [];
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user