use std::io::{self, BufRead, BufReader, Read};
use std::sync::mpsc;
use std::thread;
use iced::mouse;
use iced::{Event, Size, Theme};
use serde::Serialize;
use plushie_renderer_engine::Codec;
use plushie_renderer_engine::Core;
use plushie_widget_sdk::PlushieRenderer;
use plushie_widget_sdk::image_registry::ImageRegistry;
use plushie_widget_sdk::protocol::{
IncomingMessage, OutgoingEvent, ScreenshotResponse, SessionMessage,
};
use plushie_widget_sdk::render_ctx::RenderCtx;
use plushie_widget_sdk::runtime::Message;
use plushie_renderer_lib::scripting::{interaction_to_iced_events, resolve_widget_id};
fn log_hello_error(err: &io::Error) {
if err.kind() != io::ErrorKind::BrokenPipe {
log::error!("failed to emit hello: {err}");
}
}
const DEFAULT_SCREENSHOT_WIDTH: u32 = 1024;
const DEFAULT_SCREENSHOT_HEIGHT: u32 = 768;
const MAX_SCREENSHOT_DIMENSION: u32 = 16384;
#[derive(Clone, Copy)]
pub(crate) enum Mode {
Headless,
Mock,
}
struct WireWriter {
inner: WriterInner,
codec: Codec,
}
enum WriterInner {
Channel(mpsc::SyncSender<Vec<u8>>),
}
impl WireWriter {
fn channel(tx: mpsc::SyncSender<Vec<u8>>, codec: Codec) -> Self {
Self {
inner: WriterInner::Channel(tx),
codec,
}
}
fn emit<T: Serialize>(&self, value: &T) -> io::Result<()> {
let bytes = self.codec.encode(value).map_err(io::Error::other)?;
self.write_bytes(&bytes)
}
fn emit_binary(
&self,
map: serde_json::Map<String, serde_json::Value>,
binary: Option<(&str, &[u8])>,
) -> io::Result<()> {
let bytes = self
.codec
.encode_binary_message(map, binary)
.map_err(io::Error::other)?;
self.write_bytes(&bytes)
}
fn write_bytes(&self, bytes: &[u8]) -> io::Result<()> {
match &self.inner {
WriterInner::Channel(tx) => tx
.send(bytes.to_vec())
.map_err(|_| io::Error::new(io::ErrorKind::BrokenPipe, "writer channel closed")),
}
}
}
type UiCache = iced_test::runtime::user_interface::Cache;
struct UiState<R: PlushieRenderer> {
renderer: R,
ui_cache: UiCache,
viewport_size: Size,
cursor: mouse::Cursor,
}
struct Session<R: PlushieRenderer> {
core: Core,
theme: Theme,
theme_chrome: plushie_widget_sdk::runtime::ThemeChrome,
registry: plushie_widget_sdk::registry::WidgetRegistry<R>,
images: ImageRegistry,
writer: WireWriter,
ui: UiState<R>,
mode: Mode,
transition_manager: plushie_widget_sdk::animation::TransitionManager,
current_modifiers: iced::keyboard::Modifiers,
fonts_loaded: u32,
}
impl<R: PlushieRenderer> Session<R> {
fn new(mode: Mode, writer: WireWriter) -> Self {
let mut registry = plushie_widget_sdk::registry::WidgetRegistry::new();
registry.register_set(&plushie_widget_sdk::runtime::iced_widget_set());
Self::with_registry(mode, writer, registry)
}
fn with_registry(
mode: Mode,
writer: WireWriter,
registry: plushie_widget_sdk::registry::WidgetRegistry<R>,
) -> Self {
let renderer_settings = iced::advanced::renderer::Settings {
default_font: iced::Font::DEFAULT,
default_text_size: iced::Pixels(16.0),
};
let renderer = iced::futures::executor::block_on(R::new(renderer_settings, None))
.expect("renderer must be available");
let ui = UiState {
renderer,
ui_cache: UiCache::default(),
viewport_size: Size::new(
DEFAULT_SCREENSHOT_WIDTH as f32,
DEFAULT_SCREENSHOT_HEIGHT as f32,
),
cursor: mouse::Cursor::Unavailable,
};
Self {
core: Core::new(),
theme: Theme::Dark,
theme_chrome: plushie_widget_sdk::runtime::ThemeChrome::default(),
registry,
images: ImageRegistry::new(),
writer,
ui,
mode,
fonts_loaded: 0,
transition_manager: plushie_widget_sdk::animation::TransitionManager::new(),
current_modifiers: iced::keyboard::Modifiers::default(),
}
}
fn rebuild_renderer(&mut self) {
let renderer_settings = iced::advanced::renderer::Settings {
default_font: self.core.default_font.unwrap_or(iced::Font::DEFAULT),
default_text_size: iced::Pixels(self.core.default_text_size.unwrap_or(16.0)),
};
if let Some(r) = iced::futures::executor::block_on(R::new(renderer_settings, None)) {
self.ui.renderer = r;
self.ui.ui_cache = UiCache::default();
}
}
fn with_ui<Ret>(
&mut self,
f: impl FnOnce(
&mut iced_test::runtime::UserInterface<
'_,
plushie_widget_sdk::runtime::Message,
Theme,
R,
>,
&mut R,
mouse::Cursor,
) -> Ret,
) -> Option<Ret> {
let root = self.core.tree.root()?;
let ctx = RenderCtx {
caches: &self.core.caches,
images: &self.images,
theme: &self.theme,
theme_chrome: self.theme_chrome,
registry: &self.registry,
default_text_size: self.core.default_text_size,
default_font: self.core.default_font,
window_id: "",
scale_factor: 1.0,
validate_props: self.core.is_validate_props_enabled(),
};
let element = plushie_widget_sdk::runtime::render(root, ctx);
let cache = std::mem::take(&mut self.ui.ui_cache);
let mut ui = iced_test::runtime::UserInterface::build(
element,
self.ui.viewport_size,
cache,
&mut self.ui.renderer,
);
let result = f(&mut ui, &mut self.ui.renderer, self.ui.cursor);
self.ui.ui_cache = ui.into_cache();
Some(result)
}
fn settle_ui(&mut self, session_id: &str) -> Vec<OutgoingEvent> {
let messages = self
.with_ui(|ui, renderer, cursor| {
let mut messages = Vec::new();
let redraw = Event::Window(iced::window::Event::RedrawRequested(
iced_test::core::time::Instant::now(),
));
let _status = ui.update(&[redraw], cursor, renderer, &mut messages);
messages
})
.unwrap_or_default();
self.process_captured_messages(messages)
.into_iter()
.map(|e| e.with_session(session_id))
.collect()
}
fn inject_and_capture(
&mut self,
session_id: &str,
interact_id: &str,
events: &[Event],
read_next: &mut dyn FnMut() -> Option<IncomingMessage>,
) -> bool {
if events.is_empty() {
return false;
}
let mut emitted_steps = false;
for event in events {
if let Event::Mouse(mouse::Event::CursorMoved { position }) = event {
self.ui.cursor = mouse::Cursor::Available(*position);
}
if let Event::Keyboard(iced::keyboard::Event::ModifiersChanged(mods)) = event {
self.current_modifiers = *mods;
}
let messages = self
.with_ui(|ui, renderer, cursor| {
let mut messages = Vec::new();
let statuses =
ui.update(std::slice::from_ref(event), cursor, renderer, &mut messages);
let (_ui_state, event_statuses) = statuses;
if let Some(&status) = event_statuses.first() {
iced_test::runtime::keyboard::handle_tab(event, status, ui, renderer);
}
messages
})
.unwrap_or_default();
let step_events: Vec<OutgoingEvent> = self
.process_captured_messages(messages)
.into_iter()
.map(|e| e.with_session(session_id))
.collect();
if !step_events.is_empty() {
emitted_steps = true;
let step = plushie_widget_sdk::protocol::InteractResponse {
message_type: "interact_step",
session: session_id.to_string(),
id: interact_id.to_string(),
events: step_events,
};
if self.writer.emit(&step).is_err() {
break;
}
let next = read_next();
if let Some(msg) = next {
let is_tree_change = matches!(
msg,
IncomingMessage::Snapshot { .. } | IncomingMessage::Patch { .. }
);
if !is_tree_change {
let msg_type = match &msg {
IncomingMessage::Snapshot { .. } => "snapshot",
IncomingMessage::Patch { .. } => "patch",
IncomingMessage::Query { .. } => "query",
IncomingMessage::Interact { .. } => "interact",
IncomingMessage::Reset { .. } => "reset",
IncomingMessage::Settings { .. } => "settings",
IncomingMessage::Effect { .. } => "effect",
IncomingMessage::WidgetOp { .. } => "widget_op",
IncomingMessage::WindowOp { .. } => "window_op",
IncomingMessage::SystemOp { .. } => "system_op",
IncomingMessage::SystemQuery { .. } => "system_query",
IncomingMessage::ImageOp { .. } => "image_op",
IncomingMessage::LoadFont { .. } => "load_font",
IncomingMessage::Subscribe { .. } => "subscribe",
IncomingMessage::Unsubscribe { .. } => "unsubscribe",
IncomingMessage::TreeHash { .. } => "tree_hash",
IncomingMessage::Screenshot { .. } => "screenshot",
IncomingMessage::Command { .. } => "command",
IncomingMessage::Commands { .. } => "commands",
IncomingMessage::AdvanceFrame { .. } => "advance_frame",
IncomingMessage::RegisterEffectStub { .. } => "register_effect_stub",
IncomingMessage::UnregisterEffectStub { .. } => {
"unregister_effect_stub"
}
};
log::warn!(
"interact_step: expected snapshot or patch from host, \
got {msg_type}; tree state may be stale"
);
}
let effects = self.core.apply(msg);
for effect in effects {
use plushie_renderer_engine::{CoreEffect, StateChange};
match effect {
CoreEffect::StateChange(StateChange::ThemeChanged(t, chrome)) => {
self.theme = t;
self.theme_chrome = chrome;
}
CoreEffect::StateChange(StateChange::ThemeFollowsSystem) => {
self.theme = Theme::Dark;
self.theme_chrome =
plushie_widget_sdk::runtime::ThemeChrome::default();
}
CoreEffect::StateChange(StateChange::WidgetConfig(config)) => {
let ctx = plushie_widget_sdk::registry::InitCtx {
config: &config,
theme: &self.theme,
default_text_size: self.core.default_text_size,
default_font: self.core.default_font,
};
self.registry.init_all(&ctx);
}
_ => {}
}
}
if is_tree_change {
let validate_props = self.core.is_validate_props_enabled();
if let Some(root) = self.core.tree.root_mut() {
self.registry.prepare_walk_with_validation(
root,
&mut self.core.caches,
&self.theme,
validate_props,
);
}
}
} else {
log::warn!("stdin closed mid-interact, stopping event injection");
break;
}
}
let settle_events = self.settle_ui(session_id);
if !settle_events.is_empty() {
emitted_steps = true;
let step = plushie_widget_sdk::protocol::InteractResponse {
message_type: "interact_step",
session: session_id.to_string(),
id: interact_id.to_string(),
events: settle_events,
};
self.writer.emit(&step).ok();
if let Some(msg) = read_next() {
let _ = self.core.apply(msg);
}
}
}
emitted_steps
}
fn process_captured_messages(&mut self, messages: Vec<Message>) -> Vec<OutgoingEvent> {
let mut events = Vec::new();
for msg in messages {
if matches!(&msg, Message::Event { family, .. } if family == "status") {
continue;
}
events.extend(self.registry.process_message(&msg));
}
events
}
}
fn handle_message<R: PlushieRenderer>(
s: &mut Session<R>,
session_id: &str,
msg: IncomingMessage,
read_next: &mut dyn FnMut() -> Option<IncomingMessage>,
) -> io::Result<()> {
let is_snapshot = matches!(msg, IncomingMessage::Snapshot { .. });
let is_tree_change = is_snapshot || matches!(msg, IncomingMessage::Patch { .. });
let is_settings = matches!(msg, IncomingMessage::Settings { .. });
if let IncomingMessage::Settings { ref settings } = msg {
load_fonts_from_settings(settings);
}
match msg {
IncomingMessage::Snapshot { .. }
| IncomingMessage::Patch { .. }
| IncomingMessage::Effect { .. }
| IncomingMessage::WidgetOp { .. }
| IncomingMessage::Subscribe { .. }
| IncomingMessage::Unsubscribe { .. }
| IncomingMessage::WindowOp { .. }
| IncomingMessage::SystemOp { .. }
| IncomingMessage::SystemQuery { .. }
| IncomingMessage::Settings { .. }
| IncomingMessage::ImageOp { .. }
| IncomingMessage::LoadFont { .. }
| IncomingMessage::RegisterEffectStub { .. }
| IncomingMessage::UnregisterEffectStub { .. } => {
let effects = s.core.apply(msg);
for effect in effects {
use plushie_renderer_engine::{CoreEffect, Dispatch, Emit, StateChange};
match effect {
CoreEffect::Emit(Emit::Event(event)) => {
s.writer.emit(&event.with_session(session_id))?;
}
CoreEffect::Emit(Emit::EffectResponse(response)) => {
s.writer.emit(&response.with_session(session_id))?;
}
CoreEffect::Emit(Emit::StubAck(ack)) => {
s.writer.emit(&ack.with_session(session_id))?;
}
CoreEffect::Dispatch(Dispatch::Effect {
request_id,
kind,
payload,
}) => {
if plushie_renderer_lib::effects::native::is_async_effect(&kind) {
let mode = match s.mode {
Mode::Headless => "headless",
Mode::Mock => "mock",
};
log::debug!("{mode}: async effect {kind} unsupported (no display)");
s.writer.emit(
&plushie_widget_sdk::protocol::EffectResponse::unsupported(
request_id,
)
.with_session(session_id),
)?;
} else {
let response = plushie_renderer_lib::effects::native::handle_effect(
request_id, &kind, &payload,
);
s.writer.emit(&response.with_session(session_id))?;
}
}
CoreEffect::StateChange(StateChange::ThemeChanged(t, chrome)) => {
let mode_str = if t == iced::Theme::Light {
"light"
} else {
"dark"
};
for entry in s.core.matching_entries(
plushie_renderer_lib::constants::SUB_THEME_CHANGE,
None,
) {
s.writer.emit(
&plushie_widget_sdk::protocol::OutgoingEvent::theme_changed(
entry.tag.as_str(),
mode_str,
)
.with_session(session_id),
)?;
}
s.theme = t;
s.theme_chrome = chrome;
}
CoreEffect::Dispatch(Dispatch::Image {
op,
handle,
data,
pixels,
width,
height,
}) => {
let mode = match s.mode {
Mode::Headless => "headless",
Mode::Mock => "mock",
};
if let Err(e) = s.images.apply_op(&op, &handle, data, pixels, width, height)
{
log::warn!("{mode}: image_op {op} failed: {e}");
}
}
CoreEffect::StateChange(StateChange::WidgetConfig(config)) => {
let ctx = plushie_widget_sdk::registry::InitCtx {
config: &config,
theme: &s.theme,
default_text_size: s.core.default_text_size,
default_font: s.core.default_font,
};
s.registry.init_all(&ctx);
}
CoreEffect::StateChange(StateChange::SyncWindows) => {}
CoreEffect::Dispatch(Dispatch::WidgetOp {
ref op,
ref payload,
}) if op == "load_font" => {
load_font_from_payload(s, session_id, payload);
}
CoreEffect::Dispatch(Dispatch::WidgetOp {
ref op,
ref payload,
}) if op == "announce" => {
let announce_text = payload
.get("text")
.and_then(|v| v.as_str())
.unwrap_or_default()
.to_string();
let politeness = payload
.get("politeness")
.and_then(|v| v.as_str())
.unwrap_or("assertive")
.to_string();
let event = plushie_widget_sdk::protocol::OutgoingEvent::generic(
"announce",
"",
Some(serde_json::json!({
"text": announce_text,
"politeness": politeness,
})),
);
s.writer.emit(&event.with_session(session_id))?;
}
CoreEffect::Dispatch(Dispatch::WidgetOp {
ref op,
ref payload,
}) if op == "find_focused" => {
let tag = payload
.get("tag")
.and_then(|v| v.as_str())
.unwrap_or("find_focused")
.to_string();
let resp = serde_json::json!({
"type": "op_query_response",
"session": session_id,
"kind": "find_focused",
"tag": tag,
"data": {"focused": null}
});
s.writer.emit(&resp)?;
}
CoreEffect::Dispatch(Dispatch::WidgetOp {
ref op,
ref payload,
}) if op == "list_images" => {
let tag = payload
.get("tag")
.or_else(|| payload.get("target"))
.and_then(|v| v.as_str())
.unwrap_or("list_images")
.to_string();
let handles = s.images.handle_names();
let resp = serde_json::json!({
"type": "op_query_response",
"session": session_id,
"kind": "list_images",
"tag": tag,
"data": {"handles": handles}
});
s.writer.emit(&resp)?;
}
CoreEffect::Dispatch(Dispatch::WidgetOp { ref op, .. })
if op == "clear_images" =>
{
s.images = ImageRegistry::new();
}
CoreEffect::Dispatch(Dispatch::WidgetOp { .. }) => {}
CoreEffect::Dispatch(Dispatch::Window(_))
| CoreEffect::Dispatch(Dispatch::WindowQuery(_))
| CoreEffect::Dispatch(Dispatch::System(_))
| CoreEffect::Dispatch(Dispatch::SystemQuery(_)) => {}
CoreEffect::StateChange(StateChange::ThemeFollowsSystem) => {
s.theme = Theme::Dark;
s.theme_chrome = plushie_widget_sdk::runtime::ThemeChrome::default();
}
CoreEffect::StateChange(StateChange::ExitNodes(nodes)) => {
for (parent_id, index, node) in nodes {
s.transition_manager
.ghosts
.add_ghost(&parent_id, node, index);
}
}
}
}
if is_settings {
s.rebuild_renderer();
}
if is_tree_change {
if is_snapshot {
s.transition_manager.clear();
}
let validate_props = s.core.is_validate_props_enabled();
if let Some(root) = s.core.tree.root_mut() {
s.registry.prepare_and_scan_with_validation(
root,
&mut s.core.caches,
&s.theme,
&mut s.transition_manager,
validate_props,
);
}
let settle_events = s.settle_ui(session_id);
for event in settle_events {
s.writer.emit(&event).ok();
}
}
}
IncomingMessage::Query {
id,
target,
selector,
} => {
let resp = plushie_renderer_lib::scripting::build_query_response(
&s.core, id, target, selector,
)
.with_session(session_id);
s.writer.emit(&resp)?;
}
IncomingMessage::Interact {
id,
action,
selector,
payload,
} => {
let widget_id = resolve_widget_id(&s.core, &selector);
let use_synthetic = matches!(
action.as_str(),
"click" | "toggle" | "select" | "canvas_press" | "canvas_release" | "canvas_move"
) && widget_id.is_some();
let iced_events = if use_synthetic {
vec![]
} else {
let cursor = s.ui.cursor;
interaction_to_iced_events(&action, widget_id.as_deref(), &payload, cursor)
};
let events = if use_synthetic {
plushie_renderer_lib::scripting::build_interact_response(
&s.core,
id.clone(),
action,
selector,
payload,
)
.events
} else if !iced_events.is_empty() {
let had_steps = s.inject_and_capture(session_id, &id, &iced_events, read_next);
if had_steps {
vec![]
} else {
plushie_renderer_lib::scripting::build_interact_response(
&s.core,
id.clone(),
action,
selector,
payload,
)
.events
}
} else {
plushie_renderer_lib::scripting::build_interact_response(
&s.core,
id.clone(),
action,
selector,
payload,
)
.events
};
let resp = plushie_widget_sdk::protocol::InteractResponse::new(id, events)
.with_session(session_id);
s.writer.emit(&resp)?;
}
IncomingMessage::TreeHash { id, name, .. } => {
let resp = plushie_renderer_lib::scripting::build_tree_hash_response(&s.core, id, name)
.with_session(session_id);
s.writer.emit(&resp)?;
}
IncomingMessage::Screenshot {
id,
name,
width,
height,
} => {
let w = width
.unwrap_or(DEFAULT_SCREENSHOT_WIDTH)
.clamp(1, MAX_SCREENSHOT_DIMENSION);
let h = height
.unwrap_or(DEFAULT_SCREENSHOT_HEIGHT)
.clamp(1, MAX_SCREENSHOT_DIMENSION);
handle_screenshot(s, session_id, id, name, w, h)?;
}
IncomingMessage::Reset { id } => {
s.images = ImageRegistry::new();
s.theme = Theme::Dark;
s.transition_manager.clear();
s.ui.ui_cache = UiCache::default();
s.ui.cursor = mouse::Cursor::Unavailable;
s.rebuild_renderer();
let resp = plushie_renderer_lib::scripting::build_reset_response(&mut s.core, id)
.with_session(session_id);
s.writer.emit(&resp)?;
}
IncomingMessage::Command { id, family, value } => {
if let Some(events) = s.registry.handle_widget_op(&id, &family, &value) {
for event in events {
s.writer.emit(&event.with_session(session_id))?;
}
}
}
IncomingMessage::Commands { commands } => {
for cmd in commands {
if let Some(events) = s
.registry
.handle_widget_op(&cmd.id, &cmd.family, &cmd.value)
{
for event in events {
s.writer.emit(&event.with_session(session_id))?;
}
}
}
}
IncomingMessage::AdvanceFrame { timestamp } => {
let completions = s
.transition_manager
.advance_with_timestamp(timestamp, &mut s.core.caches.interpolated_props);
for c in completions {
let event = plushie_widget_sdk::protocol::OutgoingEvent::generic(
"transition_complete",
c.widget_id.clone(),
Some(serde_json::json!({
"tag": c.tag,
"prop": c.prop_name,
})),
);
s.writer.emit(&event.with_session(session_id))?;
}
for entry in s
.core
.matching_entries(plushie_renderer_lib::constants::SUB_ANIMATION_FRAME, None)
{
s.writer.emit(
&plushie_widget_sdk::protocol::OutgoingEvent::animation_frame(
entry.tag.as_str(),
timestamp,
)
.with_session(session_id),
)?;
}
}
}
Ok(())
}
fn handle_screenshot<R: PlushieRenderer>(
s: &mut Session<R>,
session_id: &str,
id: String,
name: String,
width: u32,
height: u32,
) -> io::Result<()> {
let emit_stub = |s: &Session<R>| {
let map = screenshot_response_map(session_id, &id, &name, "", 0, 0);
s.writer.emit_binary(map, None)
};
if matches!(s.mode, Mode::Mock) {
return emit_stub(s);
}
use iced_test::core::theme::Base;
use sha2::{Digest, Sha256};
s.ui.viewport_size = Size::new(width as f32, height as f32);
let root = match s.core.tree.root() {
Some(r) => r,
None => return emit_stub(s),
};
let ctx = RenderCtx {
caches: &s.core.caches,
images: &s.images,
theme: &s.theme,
theme_chrome: s.theme_chrome,
registry: &s.registry,
default_text_size: s.core.default_text_size,
default_font: s.core.default_font,
window_id: "",
scale_factor: 1.0,
validate_props: s.core.is_validate_props_enabled(),
};
let element: iced::Element<'_, plushie_widget_sdk::runtime::Message, Theme, R> =
plushie_widget_sdk::runtime::render(root, ctx);
let cache = std::mem::take(&mut s.ui.ui_cache);
let mut ui = iced_test::runtime::UserInterface::build(
element,
s.ui.viewport_size,
cache,
&mut s.ui.renderer,
);
{
let cursor = s.ui.cursor;
let mut messages = Vec::new();
let redraw = Event::Window(iced::window::Event::RedrawRequested(
iced_test::core::time::Instant::now(),
));
let _status = ui.update(&[redraw], cursor, &mut s.ui.renderer, &mut messages);
}
let base = s.theme.base();
ui.draw(
&mut s.ui.renderer,
&s.theme,
&iced_test::core::renderer::Style {
text_color: base.text_color,
},
s.ui.cursor,
);
s.ui.ui_cache = ui.into_cache();
let phys_size = iced::Size::new(width, height);
let rgba =
s.ui.renderer
.screenshot(phys_size, 1.0, base.background_color);
let hash = {
let mut hasher = Sha256::new();
hasher.update(&rgba);
format!("{:x}", hasher.finalize())
};
let map = screenshot_response_map(session_id, &id, &name, &hash, width, height);
let binary = if rgba.is_empty() {
None
} else {
Some(("rgba", rgba.as_slice()))
};
s.writer.emit_binary(map, binary)
}
fn screenshot_response_map(
session: &str,
id: &str,
name: &str,
hash: &str,
width: u32,
height: u32,
) -> serde_json::Map<String, serde_json::Value> {
let response = ScreenshotResponse::new(
id.to_string(),
name.to_string(),
hash.to_string(),
width,
height,
)
.with_session(session);
match serde_json::to_value(response).expect("ScreenshotResponse must serialize") {
serde_json::Value::Object(map) => map,
_ => unreachable!("ScreenshotResponse must serialize as an object"),
}
}
#[allow(clippy::too_many_arguments)]
pub(crate) fn run(
forced_codec: Option<Codec>,
mode: Mode,
max_sessions: usize,
ext_keys: &[String],
transport_name: &str,
mut reader: BufReader<Box<dyn Read + Send>>,
writer: Box<dyn std::io::Write + Send>,
expected_token: Option<&str>,
session_factory: Option<plushie_widget_sdk::app::SessionRegistryFactory<iced::Renderer>>,
) {
let codec = match crate::startup::detect_codec(forced_codec, &mut reader) {
Ok(c) => c,
Err(e) => {
log::error!("{e}");
return;
}
};
let sink = plushie_renderer_lib::WriterSink::new(writer, codec);
plushie_renderer_lib::emitters::init_sink(Box::new(sink));
plushie_renderer_lib::emitters::install_panic_hook();
let (mode_str, backend) = match mode {
Mode::Headless => ("headless", "tiny-skia"),
Mode::Mock => ("mock", "mock"),
};
let ext_key_refs: Vec<&str> = ext_keys.iter().map(|s| s.as_str()).collect();
if let Err(e) = plushie_renderer_lib::emitters::emit_hello(
mode_str,
backend,
&ext_key_refs,
&["iced"],
transport_name,
) {
log_hello_error(&e);
return;
}
let initial = match crate::startup::read_required_settings(&codec, &mut reader) {
Ok(v) => v,
Err(e) => {
crate::startup::emit_startup_error(&codec, &e);
return;
}
};
if let Err(e) =
crate::startup::validate_settings(&initial.settings, expected_token, &ext_key_refs)
{
crate::startup::emit_startup_error(&codec, &e);
return;
}
plushie_renderer_lib::settings::apply_validate_props(&initial.settings);
match mode {
Mode::Headless => {
if max_sessions <= 1 {
run_single::<iced::Renderer>(codec, mode, &mut reader, initial, None);
} else {
run_multiplexed::<iced::Renderer>(
codec,
mode,
max_sessions,
&mut reader,
initial,
session_factory,
);
}
}
Mode::Mock => {
if max_sessions <= 1 {
run_single::<()>(codec, mode, &mut reader, initial, None);
} else {
run_multiplexed::<()>(codec, mode, max_sessions, &mut reader, initial, None);
}
}
}
log::info!("stdin closed, exiting");
}
fn load_fonts_from_settings(settings: &serde_json::Value) {
for bytes in plushie_renderer_lib::settings::parse_inline_fonts(settings) {
load_font_bytes(bytes);
}
let Some(fonts) = settings.get("fonts").and_then(|v| v.as_array()) else {
return;
};
let max_bytes = plushie_renderer_lib::constants::MAX_FONT_BYTES;
let max_count = plushie_renderer_lib::constants::MAX_LOADED_FONTS;
for font_val in fonts {
if let Some(path) = font_val.as_str() {
match std::fs::read(path) {
Ok(bytes) if bytes.is_empty() => {
log::warn!("font file is empty, skipping: {path}");
}
Ok(bytes) if bytes.len() > max_bytes => {
log::warn!(
"font {path} ({} bytes) exceeds {max_bytes} byte limit, rejecting",
bytes.len(),
);
}
Ok(bytes) => {
if !plushie_renderer_lib::constants::try_reserve_font_slot() {
log::warn!(
"font {path} dropped: process-wide cap of {max_count} fonts reached",
);
continue;
}
load_font_bytes(bytes);
log::info!("loaded font: {path}");
}
Err(e) => {
log::error!("failed to load font {path}: {e}");
}
}
}
}
}
use plushie_renderer_lib::constants::MAX_FONT_BYTES;
fn load_font_from_payload<R: PlushieRenderer>(
session: &mut Session<R>,
session_id: &str,
payload: &serde_json::Value,
) {
let Some(data_val) = payload.get("data") else {
log::error!("[code=font_load_failed] load_font: missing 'data' field");
return;
};
let Some(bytes) = plushie_renderer_lib::settings::decode_font_data(data_val) else {
log::error!("[code=font_load_failed] load_font: failed to decode font data");
return;
};
if bytes.is_empty() {
log::error!("[code=font_load_failed] load_font: empty font data");
return;
}
if bytes.len() > MAX_FONT_BYTES {
log::error!(
"[code=font_load_failed] load_font: font data ({} bytes) exceeds {} byte limit, rejecting",
bytes.len(),
MAX_FONT_BYTES
);
return;
}
if !plushie_renderer_lib::constants::try_reserve_font_slot() {
let max = plushie_renderer_lib::constants::MAX_LOADED_FONTS;
let msg = format!(
"load_font rejected: process-wide cap of {max} fonts reached \
(this session has loaded {})",
session.fonts_loaded
);
log::error!("[code=font_cap_exceeded] session '{session_id}': {msg}");
let event = plushie_widget_sdk::protocol::OutgoingEvent::generic(
"session_error",
"",
Some(serde_json::json!({ "code": "font_cap_exceeded", "error": msg })),
);
if let Err(e) = session.writer.emit(&event.with_session(session_id)) {
log::error!("[code=font_load_failed] session '{session_id}': write failed: {e}");
}
return;
}
session.fonts_loaded = session.fonts_loaded.saturating_add(1);
let len = bytes.len();
load_font_bytes(bytes);
log::info!(
"[code=font_loaded] session '{session_id}': loaded font ({len} bytes, session total {})",
session.fonts_loaded
);
}
fn load_font_bytes(bytes: Vec<u8>) {
let fs = iced::advanced::graphics::text::font_system();
let mut guard = fs.write().unwrap_or_else(|e| e.into_inner());
guard.load_font(std::borrow::Cow::Owned(bytes));
}
fn read_message(codec: Codec, reader: &mut impl BufRead) -> Option<SessionMessage> {
loop {
match codec.read_message(reader) {
Ok(None) => return None,
Ok(Some(bytes)) => {
let value: serde_json::Value = match codec.decode(&bytes) {
Ok(v) => v,
Err(e) => {
log::error!("decode error: {e}");
continue;
}
};
match SessionMessage::from_value(value) {
Ok(sm) => return Some(sm),
Err(e) => {
log::error!("decode error: {e}");
continue;
}
}
}
Err(e) => {
log::error!("read error: {e}");
return None;
}
}
}
}
fn run_single<R: PlushieRenderer>(
codec: Codec,
mode: Mode,
reader: &mut impl BufRead,
initial: crate::startup::InitialSettings,
session_factory: Option<plushie_widget_sdk::app::SessionRegistryFactory<R>>,
) {
let (writer_tx, writer_rx) = mpsc::sync_channel::<Vec<u8>>(256);
let writer_handle = thread::spawn(move || {
for bytes in writer_rx {
if plushie_renderer_lib::emitters::write_output(&bytes).is_err() {
break;
}
}
});
let writer = WireWriter::channel(writer_tx.clone(), codec);
let mut session = match session_factory {
Some(factory) => Session::<R>::with_registry(mode, writer, factory()),
None => Session::<R>::new(mode, writer),
};
{
let (session_id, msg) = initial.into_parts();
let mut read_next = || read_message(codec, reader).map(|sm| sm.message);
if let Err(e) = handle_message(&mut session, &session_id, msg, &mut read_next) {
log::error!("write error processing initial settings: {e}");
drop(session);
drop(writer_tx);
let _ = writer_handle.join();
return;
}
}
while let Some(sm) = read_message(codec, reader) {
let mut read_next = || read_message(codec, reader).map(|sm| sm.message);
if let Err(e) = handle_message(&mut session, &sm.session, sm.message, &mut read_next) {
log::error!("write error: {e}");
break;
}
}
drop(session);
drop(writer_tx);
let _ = writer_handle.join();
}
const SESSION_CHANNEL_CAP: usize = 512;
const PENDING_BUFFER_CAP: usize = 512;
enum SessionSignal {
Closed(String),
Panicked(String),
}
#[allow(clippy::too_many_lines)]
fn run_multiplexed<R: PlushieRenderer>(
codec: Codec,
mode: Mode,
max_sessions: usize,
reader: &mut impl BufRead,
initial: crate::startup::InitialSettings,
session_factory: Option<plushie_widget_sdk::app::SessionRegistryFactory<R>>,
) {
use std::collections::{HashMap, HashSet, VecDeque};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
let writer_alive = Arc::new(AtomicBool::new(true));
let (writer_tx, writer_rx) = mpsc::sync_channel::<Vec<u8>>(256);
let writer_alive_for_thread = writer_alive.clone();
let writer_handle = thread::spawn(move || {
for bytes in writer_rx {
if plushie_renderer_lib::emitters::write_output(&bytes).is_err() {
break;
}
}
writer_alive_for_thread.store(false, Ordering::SeqCst);
});
let (signal_tx, signal_rx) = mpsc::channel::<SessionSignal>();
struct Dispatch {
tx: mpsc::SyncSender<IncomingMessage>,
pending: VecDeque<IncomingMessage>,
}
let mut sessions: HashMap<String, Dispatch> = HashMap::new();
let mut session_handles: Vec<thread::JoinHandle<()>> = Vec::new();
let mut closing_sessions: HashSet<String> = HashSet::new();
let mut pending_initial_settings = Some(initial.into_incoming_message());
let emit_session_error =
|writer_tx: &mpsc::SyncSender<Vec<u8>>, session_id: &str, code: &str, message: &str| {
let event = serde_json::json!({
"type": "event",
"session": session_id,
"family": "session_error",
"id": "",
"data": { "code": code, "error": message }
});
if let Ok(bytes) = codec.encode(&event) {
let _ = writer_tx.try_send(bytes);
}
};
loop {
while let Ok(signal) = signal_rx.try_recv() {
match signal {
SessionSignal::Closed(sid) => {
closing_sessions.remove(&sid);
sessions.remove(&sid);
}
SessionSignal::Panicked(sid) => {
sessions.remove(&sid);
closing_sessions.remove(&sid);
}
}
}
if !writer_alive.load(Ordering::SeqCst) {
for sid in sessions.keys().cloned().collect::<Vec<_>>() {
emit_session_error(
&writer_tx,
&sid,
"writer_dead",
"renderer stdout writer thread exited unexpectedly",
);
}
log::error!("writer thread exited; dispatcher stopping");
break;
}
match codec.read_message(reader) {
Ok(None) => break,
Ok(Some(bytes)) => {
let value: serde_json::Value = match codec.decode(&bytes) {
Ok(v) => v,
Err(e) => {
log::error!("decode error: {e}");
continue;
}
};
let sm = match SessionMessage::from_value(value) {
Ok(sm) => sm,
Err(e) => {
log::error!("decode error: {e}");
continue;
}
};
let session_id = sm.session.clone();
while let Ok(signal) = signal_rx.try_recv() {
match signal {
SessionSignal::Closed(sid) => {
closing_sessions.remove(&sid);
sessions.remove(&sid);
}
SessionSignal::Panicked(sid) => {
sessions.remove(&sid);
closing_sessions.remove(&sid);
}
}
}
if closing_sessions.contains(&session_id) {
log::debug!("session '{session_id}': message rejected during reset teardown");
emit_session_error(
&writer_tx,
&session_id,
"session_reset_in_progress",
"session is closing; wait for session_closed before reusing this ID",
);
continue;
}
let is_reset = matches!(sm.message, IncomingMessage::Reset { .. });
let session_existed = sessions.contains_key(&session_id);
if !session_existed {
if sessions.len() >= max_sessions {
log::error!(
"max sessions ({max_sessions}) reached; \
rejecting session '{session_id}'"
);
emit_session_error(
&writer_tx,
&session_id,
"max_sessions_reached",
&format!(
"max sessions ({max_sessions}) reached; session \
'{session_id}' rejected"
),
);
continue;
}
let (tx, rx) = mpsc::sync_channel::<IncomingMessage>(SESSION_CHANNEL_CAP);
let writer = WireWriter::channel(writer_tx.clone(), codec);
let sid = session_id.clone();
let closed_writer_tx = writer_tx.clone();
let thread_factory = session_factory.clone();
let signal_tx_for_thread = signal_tx.clone();
let handle = thread::spawn(move || {
let panicked = {
let result =
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let mut session = match thread_factory {
Some(factory) => {
Session::<R>::with_registry(mode, writer, factory())
}
None => Session::<R>::new(mode, writer),
};
const READ_TIMEOUT: std::time::Duration =
std::time::Duration::from_secs(30);
let mut session_disconnected = false;
for msg in &rx {
let mut disconnected_during_read = false;
let mut read_next = || -> Option<IncomingMessage> {
match rx.recv_timeout(READ_TIMEOUT) {
Ok(m) => Some(m),
Err(_) => {
disconnected_during_read = true;
None
}
}
};
let res =
handle_message(&mut session, &sid, msg, &mut read_next);
if disconnected_during_read {
log::warn!(
"session '{sid}': read timed out or \
channel disconnected during interact"
);
let error = serde_json::json!({
"type": "event",
"session": sid,
"family": "session_error",
"id": "",
"data": {
"code": "host_disconnect",
"error": "host stopped delivering messages \
during interact"
}
});
if let Ok(bytes) = codec.encode(&error) {
let _ = closed_writer_tx.try_send(bytes);
}
session_disconnected = true;
}
if let Err(e) = res {
log::error!("session '{sid}': write error: {e}");
break;
}
if session_disconnected {
break;
}
}
log::debug!("session '{sid}' thread exiting");
}));
if let Err(payload) = result {
let msg = payload
.downcast_ref::<&str>()
.copied()
.or_else(|| {
payload.downcast_ref::<String>().map(|s| s.as_str())
})
.unwrap_or("(non-string panic)");
log::error!("session '{sid}' thread panicked: {msg}");
let error = serde_json::json!({
"type": "event",
"session": sid,
"family": "session_error",
"id": "",
"data": { "code": "session_panic", "error": msg }
});
if let Ok(bytes) = codec.encode(&error) {
let _ = closed_writer_tx.send(bytes);
}
true
} else {
false
}
};
let signal = if panicked {
SessionSignal::Panicked(sid.clone())
} else {
SessionSignal::Closed(sid.clone())
};
let _ = signal_tx_for_thread.send(signal);
let closed = serde_json::json!({
"type": "event",
"session": sid,
"family": "session_closed",
"id": "",
"data": {}
});
match codec.encode(&closed) {
Ok(bytes) => {
if closed_writer_tx.try_send(bytes).is_err() {
log::info!(
"session '{sid}': session_closed send failed \
(writer likely gone or full)"
);
}
}
Err(e) => {
log::info!("session '{sid}': session_closed encode failed: {e}");
}
}
});
sessions.insert(
session_id.clone(),
Dispatch {
tx,
pending: VecDeque::new(),
},
);
session_handles.push(handle);
log::info!(
"session '{}' created (active: {})",
session_id,
sessions.len()
);
if let Some(settings_msg) = pending_initial_settings.take()
&& let Some(d) = sessions.get_mut(&session_id)
&& try_enqueue(&mut d.tx, &mut d.pending, settings_msg).is_err()
{
log::error!("session '{session_id}': failed to queue initial settings");
}
}
let (ejected, dispatch_error) = if let Some(d) = sessions.get_mut(&session_id) {
match try_enqueue(&mut d.tx, &mut d.pending, sm.message) {
Ok(()) => (false, None),
Err(EnqueueError::Overflow) => (
true,
Some((
"session_backpressure_overflow",
format!(
"session '{session_id}' send queue saturated \
(channel + pending = {}); ejecting",
SESSION_CHANNEL_CAP + PENDING_BUFFER_CAP
),
)),
),
Err(EnqueueError::Disconnected) => (
true,
Some((
"session_channel_closed",
format!("session '{session_id}' channel closed unexpectedly"),
)),
),
}
} else {
(false, None)
};
if let Some((code, msg)) = dispatch_error {
emit_session_error(&writer_tx, &session_id, code, &msg);
}
if ejected {
sessions.remove(&session_id);
continue;
}
if is_reset {
closing_sessions.insert(session_id.clone());
sessions.remove(&session_id);
log::info!(
"session '{session_id}' reset in progress (active: {})",
sessions.len()
);
}
}
Err(e) => {
log::error!("read error: {e}");
break;
}
}
}
sessions.clear();
drop(writer_tx);
drop(signal_tx);
for handle in session_handles {
if let Err(payload) = handle.join() {
let msg = payload
.downcast_ref::<&str>()
.copied()
.or_else(|| payload.downcast_ref::<String>().map(|s| s.as_str()))
.unwrap_or("(non-string panic)");
log::error!("session thread panicked: {msg}");
}
}
if let Err(payload) = writer_handle.join() {
let msg = payload
.downcast_ref::<&str>()
.copied()
.or_else(|| payload.downcast_ref::<String>().map(|s| s.as_str()))
.unwrap_or("(non-string panic)");
log::error!("writer thread panicked: {msg}");
}
}
enum EnqueueError {
Overflow,
Disconnected,
}
fn try_enqueue(
tx: &mut mpsc::SyncSender<IncomingMessage>,
pending: &mut std::collections::VecDeque<IncomingMessage>,
msg: IncomingMessage,
) -> Result<(), EnqueueError> {
while let Some(buffered) = pending.pop_front() {
match tx.try_send(buffered) {
Ok(()) => {}
Err(mpsc::TrySendError::Full(back)) => {
pending.push_front(back);
break;
}
Err(mpsc::TrySendError::Disconnected(_)) => return Err(EnqueueError::Disconnected),
}
}
if pending.is_empty() {
match tx.try_send(msg) {
Ok(()) => Ok(()),
Err(mpsc::TrySendError::Full(back)) => {
if pending.len() >= PENDING_BUFFER_CAP {
Err(EnqueueError::Overflow)
} else {
pending.push_back(back);
Ok(())
}
}
Err(mpsc::TrySendError::Disconnected(_)) => Err(EnqueueError::Disconnected),
}
} else if pending.len() >= PENDING_BUFFER_CAP {
Err(EnqueueError::Overflow)
} else {
pending.push_back(msg);
Ok(())
}
}