1use 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
14pub trait Plugin: Sized {
18 type Config: DeserializeOwned + JsonSchema;
22 type SystemConfig: DeserializeOwned + JsonSchema;
25
26 fn manifest() -> Manifest;
29
30 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 fn shutdown(&mut self, _grace: std::time::Duration) {}
44}
45
46pub use crate::schema::NoConfig;
47
48#[derive(Clone, Debug)]
50pub struct Setup<C, S> {
51 pub config: C,
52 pub systems: Vec<System<S>>,
54 pub host: HostInfo,
55 pub capture_dir: PathBuf,
56 pub data_dir: PathBuf,
58 pub audio_formats: Vec<String>,
60}
61
62impl<C, S> Setup<C, S> {
63 pub fn system(&self, index: u16) -> Option<&System<S>> {
66 self.systems.iter().find(|s| s.index == index)
67 }
68 pub fn system_named(&self, short_name: &str) -> Option<&System<S>> {
71 self.systems.iter().find(|s| s.short_name == short_name)
72 }
73 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 pub index: u16,
83 pub short_name: String,
85 pub kind: String,
87 pub config: Option<S>,
88}
89
90#[derive(Clone)]
94pub struct Host {
95 out: Arc<Mutex<Box<dyn Write + Send>>>,
96}
97
98impl Host {
99 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 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 pub fn status(&self, state: State, message: impl Into<String>) {
130 self.send(&PluginMessage::Status { state, message: message.into() });
131 }
132
133 pub fn metrics(&self, m: &Metrics) {
136 self.send(&PluginMessage::Metrics(m.clone()));
137 }
138
139 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
145pub 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
156pub 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
189pub 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
261fn 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#[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}