Skip to main content

trunk_recorder_plugin/
protocol.rs

1//! The wire protocol: JSON lines.
2//!
3//! The recorder starts a plugin's executable with no arguments and writes one
4//! [`HostMessage`] per line to its stdin; the plugin writes one
5//! [`PluginMessage`] per line to its stdout. stderr is the plugin's log.
6//! Either side ignores message types and fields it doesn't know, so both can
7//! grow within an [`API_VERSION`].
8//!
9//! `plugin --describe` prints the plugin's [`Manifest`] (one JSON object) and exits.
10
11use std::path::PathBuf;
12
13use serde::{Deserialize, Deserializer, Serialize};
14use serde_json::{Map, Value};
15
16/// The protocol version. A plugin states the version it was built for in its
17/// manifest; the recorder runs plugins whose `api` it supports.
18pub const API_VERSION: u32 = 1;
19
20/// Exit status of a plugin that can't run with its configuration: the recorder
21/// shows the error and doesn't restart it until recording next starts.
22/// (Any other exit is a crash, and the plugin is restarted with backoff.)
23pub const EXIT_CONFIG: i32 = 78;
24
25/// What a plugin can subscribe to ([`Manifest::subscribe`]).
26pub mod topic {
27    /// A call started ([`super::CallInfo`]).
28    pub const CALL_START: &str = "call.start";
29    /// A call ended ([`super::CallInfo`]); its files follow in `call.concluded`.
30    pub const CALL_END: &str = "call.end";
31    /// A recorded call's files are on disk ([`super::ConcludedCall`]).
32    pub const CALL_CONCLUDED: &str = "call.concluded";
33    /// A radio registered, affiliated, … ([`super::UnitEvent`]).
34    pub const UNIT: &str = "unit";
35    /// Live audio of recording calls ([`super::AudioChunk`]).
36    pub const AUDIO: &str = "audio";
37    /// Every system's state, every few seconds ([`super::Status`]).
38    pub const STATUS: &str = "status";
39
40    pub const ALL: &[&str] = &[CALL_START, CALL_END, CALL_CONCLUDED, UNIT, AUDIO, STATUS];
41}
42
43/// Audio formats a plugin can ask for ([`Manifest::audio_formats`]). WAV is always there.
44pub mod format {
45    /// AAC in an MP4 container (OpenMHz, Broadcastify Calls).
46    pub const M4A: &str = "m4a";
47}
48
49/// Who a plugin is and what it wants — printed by `plugin --describe`.
50#[derive(Clone, Debug, Default, Serialize, Deserialize, PartialEq)]
51#[serde(default)]
52pub struct Manifest {
53    /// Unique, lowercase, `a-z0-9-` (the registry key and the install folder).
54    pub id: String,
55    /// Shown in the interface.
56    pub name: String,
57    /// Semver of the plugin.
58    pub version: String,
59    pub description: String,
60    /// The protocol version the plugin speaks ([`API_VERSION`]).
61    pub api: u32,
62    /// Topics ([`topic`]) to receive. Nothing else is sent — or even built.
63    pub subscribe: Vec<String>,
64    /// Extra audio formats ([`format`](mod@format)) for `call.concluded`. The recorder
65    /// encodes each call once for every plugin that asks, when it can.
66    pub audio_formats: Vec<String>,
67    /// JSON Schema of the plugin's settings (an object), for the settings form.
68    #[serde(skip_serializing_if = "Option::is_none")]
69    pub config: Option<Value>,
70    /// JSON Schema of the settings the plugin takes for each system (an
71    /// object); the form repeats it under each system.
72    #[serde(skip_serializing_if = "Option::is_none")]
73    pub system_config: Option<Value>,
74    #[serde(skip_serializing_if = "String::is_empty")]
75    pub homepage: String,
76    #[serde(skip_serializing_if = "String::is_empty")]
77    pub repository: String,
78    #[serde(skip_serializing_if = "Vec::is_empty")]
79    pub authors: Vec<String>,
80    #[serde(skip_serializing_if = "String::is_empty")]
81    pub license: String,
82}
83
84impl Manifest {
85    pub fn subscribes(&self, topic: &str) -> bool {
86        self.subscribe.iter().any(|t| t == topic)
87    }
88    pub fn wants_format(&self, format: &str) -> bool {
89        self.audio_formats.iter().any(|f| f == format)
90    }
91}
92
93// ─── Recorder → plugin ────────────────────────────────────────────────────
94
95// (Messages are read, handled and dropped one at a time: their size doesn't matter.)
96#[allow(clippy::large_enum_variant)]
97#[derive(Clone, Debug, Serialize, Deserialize)]
98#[serde(tag = "type")]
99pub enum HostMessage {
100    /// Always first, once.
101    #[serde(rename = "hello")]
102    Hello(Hello),
103    #[serde(rename = "call.start")]
104    CallStart(CallInfo),
105    #[serde(rename = "call.end")]
106    CallEnd(CallInfo),
107    #[serde(rename = "call.concluded")]
108    CallConcluded(ConcludedCall),
109    #[serde(rename = "unit")]
110    Unit(UnitEvent),
111    #[serde(rename = "audio")]
112    Audio(AudioChunk),
113    #[serde(rename = "status")]
114    Status(Status),
115    /// The recorder is stopping: finish what's queued, then exit. stdin closes
116    /// after this, and the process is killed if it hasn't exited within
117    /// [`Shutdown::grace_s`].
118    #[serde(rename = "shutdown")]
119    Shutdown(Shutdown),
120    /// A message type from a newer recorder.
121    #[serde(other)]
122    Unknown,
123}
124
125#[derive(Clone, Debug, Default, Serialize, Deserialize)]
126#[serde(default)]
127pub struct Hello {
128    /// The recorder's protocol version.
129    pub api: u32,
130    pub host: HostInfo,
131    /// The plugin's settings, as the user entered them (see [`Manifest::config`]).
132    pub config: Value,
133    /// Every system calls can come from, including conventional channels.
134    pub systems: Vec<SystemInfo>,
135    /// Where calls are stored.
136    pub capture_dir: PathBuf,
137    /// A folder of the plugin's own that survives restarts and upgrades (queues, state).
138    pub data_dir: PathBuf,
139    /// The audio formats `call.concluded` will carry this run: always "wav",
140    /// plus those the plugin asked for that the recorder can encode.
141    pub audio_formats: Vec<String>,
142}
143
144#[derive(Clone, Debug, Default, Serialize, Deserialize)]
145#[serde(default)]
146pub struct HostInfo {
147    pub name: String,
148    pub version: String,
149}
150
151/// The index calls and events use for conventional channels.
152pub const CONVENTIONAL: u16 = 65535;
153
154#[derive(Clone, Debug, Default, Serialize, Deserialize)]
155#[serde(default)]
156pub struct SystemInfo {
157    /// What events call it ([`CONVENTIONAL`] for conventional channels).
158    pub index: u16,
159    /// The system's short name: its folder, and what users know it by.
160    pub short_name: String,
161    /// "p25" | "smartnet" | "dmr" | "conventional"
162    pub kind: String,
163    /// The plugin's settings for this system (see [`Manifest::system_config`]);
164    /// null when the user left them empty.
165    pub config: Value,
166}
167
168/// A call as it starts or ends.
169#[derive(Clone, Debug, Default, Serialize, Deserialize)]
170#[serde(default)]
171pub struct CallInfo {
172    /// Unique while the recorder runs (Trunk Recorder's call_num).
173    pub id: u32,
174    pub system: u16,
175    pub short_name: String,
176    pub talkgroup: u32,
177    /// The talkgroup's alpha tag, when known.
178    pub talkgroup_tag: String,
179    pub freq_hz: u64,
180    /// Phase II TDMA slot.
181    pub tdma_slot: Option<u8>,
182    pub analog: bool,
183    pub encrypted: bool,
184    pub emergency: bool,
185    /// Recorded (false: only followed, or not recorded — see `reason`).
186    pub recording: bool,
187    /// Why it isn't recorded: "encrypted", "unknown_tg", "no_recorder", "no_source", …
188    pub reason: Option<String>,
189    /// Unix seconds.
190    pub start_time: f64,
191    /// Radios heard on the call so far.
192    pub units: Vec<u32>,
193    /// Talkgroups patched with this one so far, its own included, ascending;
194    /// empty when it isn't patched.
195    pub patched_talkgroups: Vec<u32>,
196}
197
198/// A recorded call whose files are on disk.
199#[derive(Clone, Debug, Default, Serialize, Deserialize)]
200#[serde(default)]
201pub struct ConcludedCall {
202    /// The call's key: its files' path relative to the capture folder, without
203    /// an extension (`<shortName>/<Y>/<M>/<D>/<tg>-<start>_<freq>`).
204    pub path: String,
205    pub system: u16,
206    /// Its call JSON, in Trunk Recorder's format (the `.json` file's contents).
207    pub call: CallRecord,
208    pub files: CallFiles,
209}
210
211#[derive(Clone, Debug, Default, Serialize, Deserialize)]
212#[serde(default)]
213pub struct CallFiles {
214    pub json: PathBuf,
215    /// 16-bit mono WAV, 8 kHz.
216    pub wav: PathBuf,
217    /// When the plugin asked for M4A and the recorder could encode it.
218    #[serde(skip_serializing_if = "Option::is_none")]
219    pub m4a: Option<PathBuf>,
220}
221
222/// Trunk Recorder's call JSON. Fields this crate doesn't name are in `extra`.
223#[derive(Clone, Debug, Default, Serialize, Deserialize)]
224#[serde(default)]
225pub struct CallRecord {
226    pub call_num: u64,
227    pub short_name: String,
228    pub talkgroup: u32,
229    pub talkgroup_tag: String,
230    pub talkgroup_description: String,
231    pub talkgroup_group_tag: String,
232    pub talkgroup_group: String,
233    /// Hz.
234    pub freq: u64,
235    /// Unix seconds.
236    pub start_time: i64,
237    pub stop_time: i64,
238    /// Seconds of audio.
239    pub call_length: f64,
240    #[serde(deserialize_with = "flag")]
241    pub emergency: bool,
242    #[serde(deserialize_with = "flag")]
243    pub encrypted: bool,
244    pub priority: i64,
245    #[serde(deserialize_with = "flag")]
246    pub phase2_tdma: bool,
247    pub tdma_slot: i64,
248    /// "digital" | "digital tdma" | "analog"
249    pub audio_type: String,
250    #[serde(rename = "freqList")]
251    pub freq_list: Vec<FreqEntry>,
252    #[serde(rename = "srcList")]
253    pub src_list: Vec<SrcEntry>,
254    /// Every talkgroup patched with this one during the call, its own
255    /// included; absent (empty) when it wasn't patched with another.
256    #[serde(skip_serializing_if = "Vec::is_empty")]
257    pub patched_talkgroups: Vec<u32>,
258    #[serde(flatten)]
259    pub extra: Map<String, Value>,
260}
261
262impl CallRecord {
263    /// Voice frames with bit errors, over the call.
264    pub fn error_count(&self) -> u64 {
265        self.freq_list.iter().map(|f| f.error_count).sum()
266    }
267    pub fn spike_count(&self) -> u64 {
268        self.freq_list.iter().map(|f| f.spike_count).sum()
269    }
270}
271
272#[derive(Clone, Debug, Default, Serialize, Deserialize)]
273#[serde(default)]
274pub struct FreqEntry {
275    pub freq: u64,
276    /// Unix seconds.
277    pub time: i64,
278    /// Seconds into the audio.
279    pub pos: f64,
280    /// Seconds.
281    pub len: f64,
282    pub error_count: u64,
283    pub spike_count: u64,
284    #[serde(flatten)]
285    pub extra: Map<String, Value>,
286}
287
288/// A radio heard on the call.
289#[derive(Clone, Debug, Default, Serialize, Deserialize)]
290#[serde(default)]
291pub struct SrcEntry {
292    /// Radio ID.
293    pub src: i64,
294    /// Unix seconds.
295    pub time: i64,
296    /// Seconds into the audio.
297    pub pos: f64,
298    #[serde(deserialize_with = "flag")]
299    pub emergency: bool,
300    pub signal_system: String,
301    /// The radio's name from the unit tags file.
302    pub tag: String,
303    /// Its talker alias, heard over the air.
304    pub tag_ota: String,
305    #[serde(flatten)]
306    pub extra: Map<String, Value>,
307}
308
309/// Trunk Recorder writes flags as 0/1 and as booleans; take either.
310fn flag<'de, D: Deserializer<'de>>(d: D) -> Result<bool, D::Error> {
311    Ok(match Value::deserialize(d)? {
312        Value::Bool(b) => b,
313        Value::Number(n) => n.as_f64().is_some_and(|v| v != 0.0),
314        Value::String(s) => s == "true" || s == "1",
315        _ => false,
316    })
317}
318
319/// Something a radio did on the control channel.
320#[derive(Clone, Debug, Default, Serialize, Deserialize)]
321#[serde(default)]
322pub struct UnitEvent {
323    pub system: u16,
324    pub short_name: String,
325    /// "registration" | "deregistration" | "affiliation" | "acknowledge" |
326    /// "location" | "data_grant" | "answer_request" | "call_alert"
327    pub kind: String,
328    /// Radio ID.
329    pub unit: u32,
330    /// The talkgroup involved (affiliation, location, answer request, call alert).
331    pub talkgroup: Option<u32>,
332    /// Unix seconds.
333    pub time: f64,
334}
335
336/// A slice of a recording call's audio: 16-bit mono PCM, little-endian, base64.
337#[derive(Clone, Debug, Default, Serialize, Deserialize)]
338#[serde(default)]
339pub struct AudioChunk {
340    pub call_id: u32,
341    pub system: u16,
342    pub talkgroup: u32,
343    pub sample_rate: u32,
344    pub pcm: String,
345}
346
347impl AudioChunk {
348    pub fn new(call_id: u32, system: u16, talkgroup: u32, sample_rate: u32, samples: &[i16]) -> Self {
349        let bytes: Vec<u8> = samples.iter().flat_map(|s| s.to_le_bytes()).collect();
350        AudioChunk { call_id, system, talkgroup, sample_rate, pcm: base64::encode(&bytes) }
351    }
352    /// The samples.
353    pub fn samples(&self) -> Vec<i16> {
354        let b = base64::decode(&self.pcm);
355        b.chunks_exact(2).map(|c| i16::from_le_bytes([c[0], c[1]])).collect()
356    }
357}
358
359#[derive(Clone, Debug, Default, Serialize, Deserialize)]
360#[serde(default)]
361pub struct Status {
362    /// Unix seconds.
363    pub time: f64,
364    pub systems: Vec<SystemStatus>,
365}
366
367#[derive(Clone, Debug, Default, Serialize, Deserialize)]
368#[serde(default)]
369pub struct SystemStatus {
370    pub index: u16,
371    pub short_name: String,
372    /// The control channel being decoded (none: searching).
373    pub control_channel_hz: Option<u64>,
374    /// Control channel messages decoded per second.
375    pub decode_rate: f64,
376    pub active_calls: u32,
377    pub recording: u32,
378}
379
380#[derive(Clone, Debug, Default, Serialize, Deserialize)]
381#[serde(default)]
382pub struct Shutdown {
383    /// Seconds the plugin has to exit.
384    pub grace_s: f64,
385}
386
387// ─── Plugin → recorder ────────────────────────────────────────────────────
388
389#[derive(Clone, Debug, Serialize, Deserialize)]
390#[serde(tag = "type")]
391pub enum PluginMessage {
392    /// The plugin read its hello and is running.
393    #[serde(rename = "ready")]
394    Ready,
395    #[serde(rename = "log")]
396    Log { level: Level, message: String },
397    /// The plugin's health, shown next to it in the interface.
398    #[serde(rename = "status")]
399    Status {
400        state: State,
401        #[serde(default)]
402        message: String,
403    },
404    /// What became of a concluded call (`path`: [`ConcludedCall::path`]).
405    #[serde(rename = "call.result")]
406    CallResult {
407        path: String,
408        outcome: Outcome,
409        #[serde(default, skip_serializing_if = "String::is_empty")]
410        message: String,
411        /// Where the call can be found now (a link to it on the service).
412        #[serde(default, skip_serializing_if = "String::is_empty")]
413        url: String,
414    },
415    #[serde(other)]
416    Unknown,
417}
418
419#[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq, PartialOrd, Ord)]
420#[serde(rename_all = "lowercase")]
421pub enum Level {
422    Error,
423    Warn,
424    Info,
425    Debug,
426}
427
428#[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq)]
429#[serde(rename_all = "lowercase")]
430pub enum State {
431    /// Working.
432    Ok,
433    /// Working, but something needs attention (a service down, calls queued).
434    Warning,
435    /// Not working (bad settings, nothing it can do).
436    Error,
437}
438
439#[derive(Clone, Copy, Debug, Serialize, Deserialize, PartialEq, Eq)]
440#[serde(rename_all = "lowercase")]
441pub enum Outcome {
442    /// Done (uploaded, sent, stored).
443    Ok,
444    /// Deliberately not handled (a system it isn't set up for, a filtered talkgroup).
445    Skipped,
446    /// Gave up on it.
447    Failed,
448}
449
450/// Standard base64 (for [`AudioChunk::pcm`]).
451pub mod base64 {
452    const A: &[u8; 64] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
453
454    pub fn encode(b: &[u8]) -> String {
455        let mut o = String::with_capacity(b.len().div_ceil(3) * 4);
456        for c in b.chunks(3) {
457            let n = (c[0] as u32) << 16 | (*c.get(1).unwrap_or(&0) as u32) << 8 | *c.get(2).unwrap_or(&0) as u32;
458            for i in 0..4 {
459                if i <= c.len() {
460                    o.push(A[(n >> (18 - 6 * i) & 63) as usize] as char);
461                } else {
462                    o.push('=');
463                }
464            }
465        }
466        o
467    }
468
469    /// Invalid characters are skipped.
470    pub fn decode(s: &str) -> Vec<u8> {
471        let mut o = Vec::with_capacity(s.len() / 4 * 3);
472        let (mut acc, mut bits) = (0u32, 0);
473        for ch in s.bytes() {
474            let v = match ch {
475                b'A'..=b'Z' => ch - b'A',
476                b'a'..=b'z' => ch - b'a' + 26,
477                b'0'..=b'9' => ch - b'0' + 52,
478                b'+' => 62,
479                b'/' => 63,
480                _ => continue,
481            };
482            acc = acc << 6 | v as u32;
483            bits += 6;
484            if bits >= 8 {
485                bits -= 8;
486                o.push((acc >> bits) as u8);
487            }
488        }
489        o
490    }
491}
492
493#[cfg(test)]
494mod tests {
495    use super::*;
496
497    #[test]
498    fn base64_roundtrip() {
499        for n in 0..10 {
500            let b: Vec<u8> = (0..n).map(|i| (i * 37 + 5) as u8).collect();
501            assert_eq!(base64::decode(&base64::encode(&b)), b);
502        }
503        assert_eq!(base64::encode(b"hello"), "aGVsbG8=");
504    }
505
506    #[test]
507    fn unknown_types_and_fields_are_ignored() {
508        let m: HostMessage = serde_json::from_str(r#"{"type":"from.the.future","x":1}"#).unwrap();
509        assert!(matches!(m, HostMessage::Unknown));
510        let m: HostMessage = serde_json::from_str(r#"{"type":"shutdown","grace_s":5,"new_field":true}"#).unwrap();
511        assert!(matches!(m, HostMessage::Shutdown(Shutdown { grace_s }) if grace_s == 5.0));
512    }
513
514    #[test]
515    fn call_record_reads_trunk_recorder_json() {
516        let j = r#"{"call_num":7,"freq":851012500,"start_time":1700000000,"stop_time":1700000005,"emergency":0,"encrypted":1,
517            "call_length":5,"talkgroup":101,"talkgroup_tag":"Fire Disp","audio_type":"digital","short_name":"dcfd","phase2_tdma":0,
518            "freqList":[{"freq":851012500,"time":1700000000,"pos":0,"len":5,"error_count":3,"spike_count":1}],
519            "srcList":[{"src":1234,"time":1700000000,"pos":0.5,"emergency":0,"signal_system":"","tag":"","tag_ota":"E1"}],
520            "color_code":-1,"patched_talkgroups":[101,65001]}"#;
521        let c: CallRecord = serde_json::from_str(j).unwrap();
522        assert!(c.encrypted && !c.emergency);
523        assert_eq!(c.patched_talkgroups, [101, 65001]);
524        assert_eq!((c.error_count(), c.spike_count()), (3, 1));
525        assert_eq!(c.src_list[0].tag_ota, "E1");
526        assert_eq!(c.extra["color_code"], -1);
527        // And back, with the unknown field kept.
528        let v = serde_json::to_value(&c).unwrap();
529        assert_eq!(v["color_code"], -1);
530        assert_eq!(v["srcList"][0]["src"], 1234);
531        assert_eq!(v["patched_talkgroups"], serde_json::json!([101, 65001]));
532        let unpatched = serde_json::to_value(CallRecord::default()).unwrap();
533        assert!(unpatched.get("patched_talkgroups").is_none());
534    }
535}