use std::io::{self, BufWriter, Write};
use std::sync::{Arc, OnceLock};
use parking_lot::Mutex;
use plushie_widget_sdk::protocol::{DiagnosticMessage, OutgoingEvent, ScreenshotResponse};
pub type SinkMutex = parking_lot::Mutex<Box<dyn EventSink>>;
pub trait EventSink: Send {
fn emit_event(&mut self, event: OutgoingEvent) -> io::Result<()>;
fn emit_effect_response(
&mut self,
response: plushie_widget_sdk::protocol::EffectResponse,
) -> io::Result<()>;
fn emit_query_response(
&mut self,
kind: &str,
tag: &str,
data: &serde_json::Value,
) -> io::Result<()>;
fn emit_screenshot_response(
&mut self,
id: &str,
name: &str,
hash: &str,
width: u32,
height: u32,
rgba_bytes: &[u8],
) -> io::Result<()>;
fn emit_hello(
&mut self,
mode: &str,
backend: &str,
native_widgets: &[&str],
widget_set_names: &[&str],
transport: &str,
) -> io::Result<()>;
fn emit_diagnostic(&mut self, message: DiagnosticMessage) -> io::Result<()>;
fn write_raw(&mut self, bytes: &[u8]) -> io::Result<()>;
fn flush_output(&mut self) -> io::Result<()> {
Ok(())
}
}
static EVENT_SINK: OnceLock<Arc<SinkMutex>> = OnceLock::new();
pub fn init_sink(sink: Box<dyn EventSink>) {
let arc = Arc::new(Mutex::new(sink));
if EVENT_SINK.set(arc).is_err() {
panic!("event sink already initialized");
}
plushie_widget_sdk::diagnostics::set_hook(Box::new(|level, diag| {
if let Some(sink_lock) = EVENT_SINK.get() {
let msg = DiagnosticMessage::new(level, diag.clone());
let mut guard = sink_lock.lock();
if let Err(e) = guard.emit_diagnostic(msg) {
log::debug!("emit_diagnostic write error: {e}");
}
}
}));
}
pub fn sink_arc() -> Arc<SinkMutex> {
EVENT_SINK
.get()
.expect("event sink not initialized")
.clone()
}
fn with_sink<R>(f: impl FnOnce(&mut dyn EventSink) -> io::Result<R>) -> io::Result<R> {
let sink = EVENT_SINK
.get()
.ok_or_else(|| io::Error::new(io::ErrorKind::NotConnected, "event sink not initialized"))?;
let mut guard = sink.lock();
f(&mut **guard)
}
pub struct WriterSink {
writer: BufWriter<Box<dyn std::io::Write + Send>>,
codec: plushie_renderer_engine::Codec,
}
const WRITER_BUF_CAP: usize = 8 * 1024;
impl WriterSink {
pub fn new(
writer: Box<dyn std::io::Write + Send>,
codec: plushie_renderer_engine::Codec,
) -> Self {
Self {
writer: BufWriter::with_capacity(WRITER_BUF_CAP, writer),
codec,
}
}
}
impl EventSink for WriterSink {
fn emit_event(&mut self, event: OutgoingEvent) -> io::Result<()> {
let bytes = self.codec.encode(&event).map_err(io::Error::other)?;
self.writer.write_all(&bytes)
}
fn emit_effect_response(
&mut self,
response: plushie_widget_sdk::protocol::EffectResponse,
) -> io::Result<()> {
let bytes = self.codec.encode(&response).map_err(io::Error::other)?;
self.writer.write_all(&bytes)?;
self.writer.flush()
}
fn emit_query_response(
&mut self,
kind: &str,
tag: &str,
data: &serde_json::Value,
) -> io::Result<()> {
let msg = serde_json::json!({
"type": "op_query_response",
"session": "",
"kind": kind,
"tag": tag,
"data": data,
});
let bytes = self.codec.encode(&msg).map_err(io::Error::other)?;
self.writer.write_all(&bytes)?;
self.writer.flush()
}
fn emit_screenshot_response(
&mut self,
id: &str,
name: &str,
hash: &str,
width: u32,
height: u32,
rgba_bytes: &[u8],
) -> io::Result<()> {
let response = ScreenshotResponse::new(
id.to_string(),
name.to_string(),
hash.to_string(),
width,
height,
);
let map = match serde_json::to_value(&response).map_err(io::Error::other)? {
serde_json::Value::Object(map) => map,
_ => unreachable!("ScreenshotResponse must serialize as a JSON object"),
};
let binary = if rgba_bytes.is_empty() {
None
} else {
Some(("rgba", rgba_bytes))
};
let bytes = self
.codec
.encode_binary_message(map, binary)
.map_err(io::Error::other)?;
self.writer.write_all(&bytes)?;
self.writer.flush()
}
fn emit_hello(
&mut self,
mode: &str,
backend: &str,
native_widgets: &[&str],
widget_set_names: &[&str],
transport: &str,
) -> io::Result<()> {
let builtin = plushie_widget_sdk::runtime::IcedWidgetSet::type_names();
let mut all_widgets: Vec<String> = builtin
.iter()
.cloned()
.chain(native_widgets.iter().map(|s| s.to_string()))
.collect();
all_widgets.sort();
all_widgets.dedup();
let mut native_sorted: Vec<&str> = native_widgets.to_vec();
native_sorted.sort_unstable();
let msg = serde_json::json!({
"type": "hello",
"session": "",
"protocol_version": plushie_widget_sdk::protocol::PROTOCOL_VERSION,
"protocol": plushie_widget_sdk::protocol::PROTOCOL_VERSION,
"codec": self.codec.to_string(),
"version": env!("CARGO_PKG_VERSION"),
"name": "plushie-renderer",
"mode": mode,
"backend": backend,
"transport": transport,
"native_widgets": native_sorted,
"widget_sets": widget_set_names,
"widgets": all_widgets,
});
let bytes = self.codec.encode(&msg).map_err(io::Error::other)?;
self.writer.write_all(&bytes)?;
self.writer.flush()
}
fn emit_diagnostic(&mut self, message: DiagnosticMessage) -> io::Result<()> {
let bytes = self.codec.encode(&message).map_err(io::Error::other)?;
self.writer.write_all(&bytes)?;
self.writer.flush()
}
fn write_raw(&mut self, bytes: &[u8]) -> io::Result<()> {
self.writer.write_all(bytes)?;
self.writer.flush()
}
fn flush_output(&mut self) -> io::Result<()> {
self.writer.flush()
}
}
pub fn write_output(bytes: &[u8]) -> io::Result<()> {
with_sink(|sink| sink.write_raw(bytes))
}
pub fn emit_hello(
mode: &str,
backend: &str,
native_widgets: &[&str],
widget_set_names: &[&str],
transport: &str,
) -> io::Result<()> {
with_sink(|sink| sink.emit_hello(mode, backend, native_widgets, widget_set_names, transport))
}
fn panic_payload_message(payload: &(dyn std::any::Any + Send)) -> &str {
payload
.downcast_ref::<&'static str>()
.copied()
.or_else(|| payload.downcast_ref::<String>().map(|s| s.as_str()))
.unwrap_or("(non-string panic)")
}
fn emit_panic_events(msg: &str, location: &str) {
if let Some(sink_lock) = EVENT_SINK.get() {
let Some(mut guard) = sink_lock.try_lock() else {
return;
};
let error_event = plushie_widget_sdk::protocol::OutgoingEvent::generic(
"session_error",
"",
Some(serde_json::json!({
"code": "renderer_panic",
"error": msg,
"location": location,
})),
);
let _ = guard.emit_event(error_event);
let closed_event = plushie_widget_sdk::protocol::OutgoingEvent::generic(
"session_closed",
"",
Some(serde_json::json!({ "reason": "panic" })),
);
let _ = guard.emit_event(closed_event);
}
}
#[cfg(all(panic = "abort", not(target_arch = "wasm32")))]
compile_error!(
"plushie-renderer-lib requires `panic = \"unwind\"` on native targets because \
catch_unwind in the panic hook is load-bearing for widget panic isolation and \
multiplexed-session isolation. Building with `panic = \"abort\"` would \
silently make this a no-op."
);
pub fn install_panic_hook() {
let previous = std::panic::take_hook();
std::panic::set_hook(Box::new(move |info| {
let msg = panic_payload_message(info.payload());
let location = info
.location()
.map(|l| format!("{}:{}:{}", l.file(), l.line(), l.column()))
.unwrap_or_else(|| "<unknown>".to_string());
log::error!("renderer panic at {location}: {msg}");
emit_panic_events(msg, &location);
previous(info);
}));
}
#[cfg(test)]
mod tests {
use super::*;
use base64::Engine as _;
use std::sync::Arc;
use std::sync::Mutex as StdMutex;
use parking_lot::Mutex as PlMutex;
#[derive(Default)]
struct RecordingSink {
events: Arc<StdMutex<Vec<OutgoingEvent>>>,
}
#[derive(Clone, Default)]
struct SharedBuffer(Arc<StdMutex<Vec<u8>>>);
impl std::io::Write for SharedBuffer {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.0.lock().unwrap().extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
impl EventSink for RecordingSink {
fn emit_event(&mut self, event: OutgoingEvent) -> io::Result<()> {
self.events.lock().unwrap().push(event);
Ok(())
}
fn emit_effect_response(
&mut self,
_: plushie_widget_sdk::protocol::EffectResponse,
) -> io::Result<()> {
Ok(())
}
fn emit_query_response(
&mut self,
_: &str,
_: &str,
_: &serde_json::Value,
) -> io::Result<()> {
Ok(())
}
fn emit_screenshot_response(
&mut self,
_: &str,
_: &str,
_: &str,
_: u32,
_: u32,
_: &[u8],
) -> io::Result<()> {
Ok(())
}
fn emit_hello(
&mut self,
_: &str,
_: &str,
_: &[&str],
_: &[&str],
_: &str,
) -> io::Result<()> {
Ok(())
}
fn emit_diagnostic(&mut self, _: DiagnosticMessage) -> io::Result<()> {
Ok(())
}
fn write_raw(&mut self, _: &[u8]) -> io::Result<()> {
Ok(())
}
}
#[test]
fn panic_payload_message_handles_str_and_string() {
let str_payload: Box<dyn std::any::Any + Send> = Box::new("static str panic");
assert_eq!(panic_payload_message(&*str_payload), "static str panic");
let string_payload: Box<dyn std::any::Any + Send> = Box::new("owned".to_string());
assert_eq!(panic_payload_message(&*string_payload), "owned");
let other_payload: Box<dyn std::any::Any + Send> = Box::new(42u32);
assert_eq!(panic_payload_message(&*other_payload), "(non-string panic)");
}
#[test]
fn emit_panic_events_writes_error_then_closed() {
let events: Arc<StdMutex<Vec<OutgoingEvent>>> = Arc::new(StdMutex::new(Vec::new()));
let recording = RecordingSink {
events: events.clone(),
};
let arc = Arc::new(PlMutex::new(Box::new(recording) as Box<dyn EventSink>));
let fresh_init = EVENT_SINK.set(arc).is_ok();
if !fresh_init {
eprintln!(
"skipping emit_panic_events test: global EVENT_SINK already set by another test"
);
return;
}
emit_panic_events("boom", "file.rs:1:1");
let ev = events.lock().unwrap();
assert_eq!(ev.len(), 2, "expected session_error then session_closed");
assert_eq!(ev[0].family, "session_error");
assert_eq!(ev[1].family, "session_closed");
let err_value = ev[0].value.as_ref().expect("session_error carries a value");
assert_eq!(
err_value.get("code").and_then(|v| v.as_str()),
Some("renderer_panic"),
);
assert_eq!(
err_value.get("error").and_then(|v| v.as_str()),
Some("boom")
);
assert_eq!(
err_value.get("location").and_then(|v| v.as_str()),
Some("file.rs:1:1"),
);
let closed_value = ev[1]
.value
.as_ref()
.expect("session_closed carries a value");
assert_eq!(
closed_value.get("reason").and_then(|v| v.as_str()),
Some("panic"),
);
}
#[test]
fn writer_sink_screenshot_response_json_includes_structured_fields_and_base64_rgba() {
let writer = SharedBuffer::default();
let output = writer.0.clone();
let mut sink = WriterSink::new(Box::new(writer), plushie_renderer_engine::Codec::Json);
sink.emit_screenshot_response("sc1", "homepage", "d4e5f6", 2, 3, &[0, 1, 2, 3])
.unwrap();
let bytes = output.lock().unwrap().clone();
let parsed: serde_json::Value = serde_json::from_slice(&bytes[..bytes.len() - 1]).unwrap();
assert_eq!(parsed["type"], "screenshot_response");
assert_eq!(parsed["session"], "");
assert_eq!(parsed["id"], "sc1");
assert_eq!(parsed["name"], "homepage");
assert_eq!(parsed["hash"], "d4e5f6");
assert_eq!(parsed["width"], 2);
assert_eq!(parsed["height"], 3);
let rgba = parsed["rgba"].as_str().expect("rgba base64 string");
let decoded = base64::engine::general_purpose::STANDARD
.decode(rgba)
.unwrap();
assert_eq!(decoded, vec![0, 1, 2, 3]);
}
#[test]
fn writer_sink_screenshot_response_omits_rgba_when_empty() {
let writer = SharedBuffer::default();
let output = writer.0.clone();
let mut sink = WriterSink::new(Box::new(writer), plushie_renderer_engine::Codec::Json);
sink.emit_screenshot_response("sc1", "mock", "", 0, 0, &[])
.unwrap();
let bytes = output.lock().unwrap().clone();
let parsed: serde_json::Value = serde_json::from_slice(&bytes[..bytes.len() - 1]).unwrap();
assert_eq!(parsed["type"], "screenshot_response");
assert_eq!(parsed["hash"], "");
assert_eq!(parsed["width"], 0);
assert_eq!(parsed["height"], 0);
assert!(parsed.get("rgba").is_none());
}
#[test]
fn writer_sink_screenshot_response_msgpack_round_trips_rgba() {
let writer = SharedBuffer::default();
let output = writer.0.clone();
let mut sink = WriterSink::new(Box::new(writer), plushie_renderer_engine::Codec::MsgPack);
sink.emit_screenshot_response("sc1", "homepage", "d4e5f6", 2, 3, &[0, 1, 2, 3])
.unwrap();
let bytes = output.lock().unwrap().clone();
let parsed: serde_json::Value = plushie_renderer_engine::Codec::MsgPack
.decode(&bytes[4..])
.unwrap();
assert_eq!(parsed["type"], "screenshot_response");
assert_eq!(parsed["session"], "");
assert_eq!(parsed["id"], "sc1");
assert_eq!(parsed["name"], "homepage");
assert_eq!(parsed["hash"], "d4e5f6");
assert_eq!(parsed["width"], 2);
assert_eq!(parsed["height"], 3);
let rgba = parsed["rgba"].as_str().expect("rgba base64 string");
let decoded = base64::engine::general_purpose::STANDARD
.decode(rgba)
.unwrap();
assert_eq!(decoded, vec![0, 1, 2, 3]);
}
}