Skip to main content

trunk_recorder_plugin/
sdk.rs

1//! Writing a plugin: implement [`Plugin`] and call [`run`] from `main`.
2
3use std::io::{self, BufRead, IsTerminal, Write};
4use std::path::PathBuf;
5use std::sync::{Arc, Mutex};
6
7use schemars::JsonSchema;
8use serde::de::DeserializeOwned;
9use serde_json::Value;
10
11use crate::protocol::*;
12use crate::schema;
13
14/// A plugin. Every method runs on the one thread that reads the recorder's
15/// messages, one message at a time: return quickly, and hand slow work
16/// (network, disk, other processes) to a thread of your own.
17pub trait Plugin: Sized {
18    /// The plugin's settings. `#[serde(default)]` on the struct lets a user
19    /// leave any of them empty. The settings form is drawn from its schema
20    /// (see [`crate::schema`]). [`NoConfig`] for none.
21    type Config: DeserializeOwned + JsonSchema;
22    /// Settings for each system, when the plugin needs some (an API key per
23    /// system, say). [`NoConfig`] for none.
24    type SystemConfig: DeserializeOwned + JsonSchema;
25
26    /// Who the plugin is and what it subscribes to — start from [`crate::manifest!`].
27    /// `config` and `system_config` are filled in from the types above.
28    fn manifest() -> Manifest;
29
30    /// Start with the user's settings. An `Err` is shown to the user, and the
31    /// plugin isn't started again until recording next starts.
32    fn start(host: Host, setup: Setup<Self::Config, Self::SystemConfig>) -> Result<Self, String>;
33
34    fn call_start(&mut self, _call: CallInfo) {}
35    fn call_end(&mut self, _call: CallInfo) {}
36    fn call_concluded(&mut self, _call: ConcludedCall) {}
37    fn unit(&mut self, _event: UnitEvent) {}
38    fn audio(&mut self, _chunk: AudioChunk) {}
39    fn status(&mut self, _status: Status) {}
40
41    /// The recorder is stopping (or the plugin's stdin closed): finish or save
42    /// what's in flight within `grace` — the process exits when this returns.
43    fn shutdown(&mut self, _grace: std::time::Duration) {}
44}
45
46pub use crate::schema::NoConfig;
47
48/// What [`Plugin::start`] gets.
49#[derive(Clone, Debug)]
50pub struct Setup<C, S> {
51    pub config: C,
52    /// Every system, with the plugin's settings for it (`None`: left empty).
53    pub systems: Vec<System<S>>,
54    pub host: HostInfo,
55    pub capture_dir: PathBuf,
56    /// The plugin's own folder, kept across restarts and upgrades.
57    pub data_dir: PathBuf,
58    /// Audio formats `call.concluded` will carry (always "wav").
59    pub audio_formats: Vec<String>,
60}
61
62impl<C, S> Setup<C, S> {
63    /// The system an event's `system` names. (Its number is for this run
64    /// only; keep anything you store by [`System::short_name`].)
65    pub fn system(&self, index: u16) -> Option<&System<S>> {
66        self.systems.iter().find(|s| s.index == index)
67    }
68    /// The system with this short name — its identity: every event carries
69    /// it (`short_name`, or a concluded call's `call.short_name`).
70    pub fn system_named(&self, short_name: &str) -> Option<&System<S>> {
71        self.systems.iter().find(|s| s.short_name == short_name)
72    }
73    /// Whether `call.concluded` will carry `format` ([`format`](mod@format)).
74    pub fn has_format(&self, format: &str) -> bool {
75        self.audio_formats.iter().any(|f| f == format)
76    }
77}
78
79#[derive(Clone, Debug)]
80pub struct System<S> {
81    /// The number events carry for it, for this run only.
82    pub index: u16,
83    /// Its identity: unique among the recorder's systems, and what users know it by.
84    pub short_name: String,
85    /// "p25" | "smartnet" | "dmr" | "conventional"
86    pub kind: String,
87    pub config: Option<S>,
88}
89
90/// The way back to the recorder: logs, health, results. Cheap to clone and
91/// usable from any thread. (A plugin's stdout belongs to the protocol — never
92/// `println!` in a plugin; use this, or `eprintln!` for the raw log.)
93#[derive(Clone)]
94pub struct Host {
95    out: Arc<Mutex<Box<dyn Write + Send>>>,
96}
97
98impl Host {
99    /// A host writing to `w` instead of stdout (see [`crate::testing`]).
100    pub fn to_writer(w: impl Write + Send + 'static) -> Host {
101        Host { out: Arc::new(Mutex::new(Box::new(w))) }
102    }
103
104    pub fn send(&self, m: &PluginMessage) {
105        let mut line = serde_json::to_string(m).unwrap_or_default();
106        line.push('\n');
107        let mut o = self.out.lock().unwrap_or_else(|e| e.into_inner());
108        // The recorder has gone when this fails; stdin's end will stop us.
109        let _ = o.write_all(line.as_bytes()).and_then(|_| o.flush());
110    }
111
112    pub fn log(&self, level: Level, message: impl Into<String>) {
113        self.send(&PluginMessage::Log { level, message: message.into() });
114    }
115    pub fn error(&self, message: impl Into<String>) {
116        self.log(Level::Error, message)
117    }
118    pub fn warn(&self, message: impl Into<String>) {
119        self.log(Level::Warn, message)
120    }
121    pub fn info(&self, message: impl Into<String>) {
122        self.log(Level::Info, message)
123    }
124    pub fn debug(&self, message: impl Into<String>) {
125        self.log(Level::Debug, message)
126    }
127
128    /// The plugin's health, shown in the interface until the next.
129    pub fn status(&self, state: State, message: impl Into<String>) {
130        self.send(&PluginMessage::Status { state, message: message.into() });
131    }
132
133    /// How the plugin's work is going (see [`Metrics`]): queue, timing,
134    /// the services it talks to. At most every few seconds.
135    pub fn metrics(&self, m: &Metrics) {
136        self.send(&PluginMessage::Metrics(m.clone()));
137    }
138
139    /// What became of a concluded call.
140    pub fn call_result(&self, path: &str, outcome: Outcome, message: impl Into<String>, url: impl Into<String>) {
141        self.send(&PluginMessage::CallResult { path: path.to_string(), outcome, message: message.into(), url: url.into() });
142    }
143}
144
145/// The manifest with the settings schemas filled in.
146pub fn describe<P: Plugin>() -> Manifest {
147    let mut m = P::manifest();
148    m.api = API_VERSION;
149    let c = schema::schema_for::<P::Config>();
150    m.config = (!schema::is_empty(&c)).then_some(c);
151    let s = schema::schema_for::<P::SystemConfig>();
152    m.system_config = (!schema::is_empty(&s)).then_some(s);
153    m
154}
155
156/// Run the plugin: `--describe` prints its manifest; otherwise it talks to the
157/// recorder on stdin/stdout until told to stop. Call from `main`.
158pub fn run<P: Plugin>() {
159    let arg = std::env::args().nth(1);
160    match arg.as_deref() {
161        Some("--describe") => {
162            println!("{}", serde_json::to_string_pretty(&describe::<P>()).unwrap_or_default());
163            return;
164        }
165        Some("--version" | "-V") => {
166            println!("{}", P::manifest().version);
167            return;
168        }
169        Some(a) => {
170            eprintln!("unknown argument {a} (a plugin takes --describe or --version; the recorder runs it with none)");
171            std::process::exit(2);
172        }
173        None => {}
174    }
175    let stdin = io::stdin();
176    if stdin.is_terminal() {
177        let m = P::manifest();
178        eprintln!(
179            "{} {} is a Trunk Recorder Pro plugin: the recorder runs it and talks to it over stdin/stdout.\n\
180             `--describe` prints its manifest. Try it with `trunk-pro plugin run <this binary>`,\n\
181             or paste protocol lines here (a hello first).",
182            m.name, m.version
183        );
184    }
185    let code = serve::<P>(stdin.lock(), Host { out: Arc::new(Mutex::new(Box::new(io::stdout()))) });
186    std::process::exit(code);
187}
188
189/// The message loop over `input`; the process's exit status.
190pub fn serve<P: Plugin>(input: impl BufRead, host: Host) -> i32 {
191    let mut plugin: Option<P> = None;
192    let mut grace = std::time::Duration::from_secs(10);
193    for line in input.lines() {
194        let Ok(line) = line else { break };
195        if line.trim().is_empty() {
196            continue;
197        }
198        let msg: HostMessage = match serde_json::from_str(&line) {
199            Ok(m) => m,
200            Err(e) => {
201                host.warn(format!("unreadable message from the recorder ({e}): {line}"));
202                continue;
203            }
204        };
205        let Some(p) = plugin.as_mut() else {
206            let HostMessage::Hello(hello) = msg else {
207                host.warn("a message before hello, ignored");
208                continue;
209            };
210            match start::<P>(host.clone(), hello) {
211                Ok(p) => {
212                    plugin = Some(p);
213                    host.send(&PluginMessage::Ready);
214                }
215                Err(e) => {
216                    host.status(State::Error, e.clone());
217                    host.error(e);
218                    return EXIT_CONFIG;
219                }
220            }
221            continue;
222        };
223        match msg {
224            HostMessage::Hello(_) => host.warn("a second hello, ignored"),
225            HostMessage::CallStart(c) => p.call_start(c),
226            HostMessage::CallEnd(c) => p.call_end(c),
227            HostMessage::CallConcluded(c) => p.call_concluded(c),
228            HostMessage::Unit(u) => p.unit(u),
229            HostMessage::Audio(a) => p.audio(a),
230            HostMessage::Status(s) => p.status(s),
231            HostMessage::Shutdown(s) => {
232                grace = std::time::Duration::from_secs_f64(s.grace_s.max(0.0));
233                break;
234            }
235            HostMessage::Unknown => {}
236        }
237    }
238    if let Some(p) = plugin.as_mut() {
239        p.shutdown(grace);
240    }
241    0
242}
243
244fn start<P: Plugin>(host: Host, hello: Hello) -> Result<P, String> {
245    if hello.api != 0 && hello.api < API_VERSION {
246        return Err(format!("this plugin needs a newer recorder (plugin API {API_VERSION}, the recorder has {})", hello.api));
247    }
248    let config: P::Config = parse(hello.config, "settings")?;
249    let mut systems = Vec::new();
250    for s in hello.systems {
251        let config = match s.config {
252            Value::Null => None,
253            v => Some(parse::<P::SystemConfig>(v, &format!("settings for {}", s.short_name))?),
254        };
255        systems.push(System { index: s.index, short_name: s.short_name, kind: s.kind, config });
256    }
257    let setup = Setup { config, systems, host: hello.host, capture_dir: hello.capture_dir, data_dir: hello.data_dir, audio_formats: hello.audio_formats };
258    P::start(host, setup)
259}
260
261/// Settings left out entirely count as `{}`, so `#[serde(default)]` applies.
262fn parse<T: DeserializeOwned>(v: Value, what: &str) -> Result<T, String> {
263    let v = if v.is_null() { Value::Object(Default::default()) } else { v };
264    serde_json::from_value(v).map_err(|e| format!("bad {what}: {e}"))
265}
266
267/// A manifest with the id, version, description, homepage, repository,
268/// authors and license from Cargo.toml; the name is the id until you set one:
269///
270/// ```ignore
271/// fn manifest() -> Manifest {
272///     Manifest { name: "OpenMHz".into(), subscribe: vec![topic::CALL_CONCLUDED.into()], ..trunk_recorder_plugin::manifest!() }
273/// }
274/// ```
275#[macro_export]
276macro_rules! manifest {
277    () => {{
278        let authors: Vec<String> = env!("CARGO_PKG_AUTHORS").split(':').filter(|a| !a.is_empty()).map(String::from).collect();
279        $crate::Manifest {
280            id: env!("CARGO_PKG_NAME").to_string(),
281            name: env!("CARGO_PKG_NAME").to_string(),
282            version: env!("CARGO_PKG_VERSION").to_string(),
283            description: env!("CARGO_PKG_DESCRIPTION").to_string(),
284            api: $crate::API_VERSION,
285            homepage: env!("CARGO_PKG_HOMEPAGE").to_string(),
286            repository: env!("CARGO_PKG_REPOSITORY").to_string(),
287            authors,
288            license: env!("CARGO_PKG_LICENSE").to_string(),
289            ..Default::default()
290        }
291    }};
292}