mod layout; use layout::{ encode_error_tape, LayoutContext, LayoutDocument, LayoutTape, TapeIdentity, TapeOutputOptions, MAX_LAYOUT_DIMENSION, MIN_TAPE_BYTES, }; use serde::{Deserialize, Serialize}; use std::collections::{HashMap, HashSet, VecDeque}; use std::ffi::c_void; use std::os::raw::c_int; use std::panic::{catch_unwind, AssertUnwindSafe}; use std::ptr; use std::slice; use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; use std::sync::{Arc, Condvar, Mutex}; use std::thread::{self, JoinHandle}; use std::time::Duration; const CONTROL_VERSION: u32 = 1; const MAX_CONTROL_BYTES: usize = 64 * 1024 * 1024; const SYNC_RENDER_MAX_RESULT_BYTES: usize = 64 * 1024 * 1024; const SYNC_RENDER_SESSION_ID: u64 = u64::MAX; static NEXT_SESSION_ID: AtomicU64 = AtomicU64::new(1); fn default_true() -> bool { true } #[repr(C)] pub struct NativeBytes { data: *mut u8, len: usize, } impl NativeBytes { fn empty() -> Self { Self { data: ptr::null_mut(), len: 0, } } fn from_vec(bytes: Vec) -> Self { if bytes.is_empty() { return Self::empty(); } let boxed = bytes.into_boxed_slice(); let len = boxed.len(); let data = Box::into_raw(boxed) as *mut u8; Self { data, len } } fn from_error(message: impl Into) -> Self { Self::from_vec(message.into().into_bytes()) } } /// One already-normalized flex item crossing the pure geometry boundary. /// /// Ebox owns the node identity, style cascade, display measurement, and /// rendering. The native kernel receives only integer geometry facts and /// numeric distribution factors, then returns target main-axis sizes. #[repr(C)] #[derive(Clone, Copy, Debug)] pub struct NativeFlexItem { pub base: i64, pub hypothetical: i64, pub min_main: i64, pub max_main: i64, pub grow: f64, pub shrink: f64, pub max_known: bool, } const MAX_FLEX_ITEMS: usize = 8192; const MAX_FLEX_LINES: usize = 2048; fn flex_distribute(amount: i64, weights: &[f64]) -> Vec { let amount = amount.max(0); let total: f64 = weights.iter().copied().filter(|weight| *weight > 0.0).sum(); if amount == 0 || total <= 0.0 { return vec![0; weights.len()]; } let mut shares = weights .iter() .map(|weight| { if *weight > 0.0 { ((amount as f64 * *weight) / total).floor() as i64 } else { 0 } }) .collect::>(); let mut remaining = amount - shares.iter().sum::(); while remaining > 0 { for (share, weight) in shares.iter_mut().zip(weights.iter()) { if remaining == 0 { break; } if *weight > 0.0 { *share += 1; remaining -= 1; } } } shares } fn clamp_flex_main(size: i64, min_main: i64, max_main: Option) -> i64 { let size = size.max(0); let min_main = min_main.max(0); let upper = max_main.unwrap_or(999_999_999).max(0); min_main.max(upper.min(size)) } fn flex_size_line(items: &[NativeFlexItem], main_limit: i64, main_gap: i64) -> Vec { let count = items.len(); let gap_total = main_gap * (count.saturating_sub(1) as i64); let available = main_limit - gap_total; let mut targets = items .iter() .map(|item| item.hypothetical) .collect::>(); let hypothetical_total = items.iter().map(|item| item.hypothetical).sum::(); let grow = hypothetical_total < available; let mut frozen = vec![false; count]; for (index, item) in items.iter().enumerate() { let factor = if grow { item.grow } else { item.shrink }; targets[index] = item.base; if factor <= 0.0 || (grow && item.base > item.hypothetical) || (!grow && item.base < item.hypothetical) { targets[index] = item.hypothetical; frozen[index] = true; } } let initial_free = available - items .iter() .enumerate() .map(|(index, item)| { if frozen[index] { targets[index] } else { item.base } }) .sum::(); loop { let free = available - items .iter() .enumerate() .map(|(index, item)| { if frozen[index] { targets[index] } else { item.base } }) .sum::(); let active = (0..count) .filter(|index| !frozen[*index]) .collect::>(); let weights = active .iter() .map(|index| { let item = items[*index]; if grow { item.grow } else { item.base as f64 * item.shrink } }) .collect::>(); let weight_total: f64 = weights.iter().copied().sum(); let factor_total: f64 = active .iter() .map(|index| { let item = items[*index]; if grow { item.grow } else { item.shrink } }) .sum(); let effective_free = if factor_total > 0.0 && factor_total < 1.0 { let partial = initial_free as f64 * factor_total; if partial.abs() < (free as f64).abs() { partial } else { free as f64 } } else { free as f64 }; if active.is_empty() || weight_total <= 0.0 || (grow && effective_free < 0.0) || (!grow && effective_free > 0.0) { break; } let deltas = flex_distribute(effective_free.abs().floor() as i64, &weights); let mut min_violations = Vec::new(); let mut max_violations = Vec::new(); let mut total_violation = 0i64; for (active_offset, index) in active.iter().enumerate() { let item = items[*index]; let candidate = if grow { item.base + deltas[active_offset] } else { item.base - deltas[active_offset] }; let clamped = clamp_flex_main( candidate, item.min_main, item.max_known.then_some(item.max_main), ); let adjustment = clamped - candidate; targets[*index] = clamped; if adjustment > 0 { min_violations.push(*index); } else if adjustment < 0 { max_violations.push(*index); } total_violation += adjustment; } if total_violation == 0 { for index in active { frozen[index] = true; } } else if total_violation > 0 { let indices = if min_violations.is_empty() { active } else { min_violations }; for index in indices { frozen[index] = true; } } else { let indices = if max_violations.is_empty() { active } else { max_violations }; for index in indices { frozen[index] = true; } } } targets } /// Calculate all flex line target sizes without touching Ebox state. /// /// `line_offsets` contains `line_count + 1` offsets into ITEMS. The output /// is a little-endian i64 stream whose values follow the input line order. /// /// # Safety /// /// Nonempty ITEMS and LINE_OFFSETS must reference their declared readable /// lengths, and OUTPUT must point to writable `NativeBytes` storage. #[no_mangle] pub unsafe extern "C" fn ebox_native_flex_size_lines( items: *const NativeFlexItem, item_count: usize, line_offsets: *const usize, line_count: usize, main_limit: i64, main_gap: i64, output: *mut NativeBytes, ) -> bool { if output.is_null() || (item_count > 0 && items.is_null()) || (line_count > 0 && line_offsets.is_null()) || item_count > MAX_FLEX_ITEMS || line_count > MAX_FLEX_LINES || !(-MAX_LAYOUT_DIMENSION..=MAX_LAYOUT_DIMENSION).contains(&main_limit) || !(0..=MAX_LAYOUT_DIMENSION).contains(&main_gap) { set_error( output, "Native flex geometry input is outside the bounded contract", ); return false; } let items = slice::from_raw_parts(items, item_count); let offsets = slice::from_raw_parts(line_offsets, line_count.saturating_add(1)); if offsets.first().copied().unwrap_or(0) != 0 || offsets.last().copied().unwrap_or(0) != item_count || offsets.windows(2).any(|window| window[0] > window[1]) { set_error(output, "Native flex geometry line offsets are invalid"); return false; } if items.iter().any(|item| { !item.grow.is_finite() || !item.shrink.is_finite() || item.grow < 0.0 || item.shrink < 0.0 }) { set_error(output, "Native flex geometry factors are invalid"); return false; } let mut bytes = Vec::with_capacity(item_count.saturating_mul(8)); for window in offsets.windows(2).take(line_count) { for target in flex_size_line(&items[window[0]..window[1]], main_limit, main_gap) { bytes.extend_from_slice(&target.to_le_bytes()); } } unsafe { *output = NativeBytes::from_vec(bytes); } true } #[derive(Debug, Deserialize)] #[serde(deny_unknown_fields)] struct ControlBatch { version: u32, #[serde(default)] document: Option, frames: Vec, } #[derive(Debug, Deserialize)] #[serde(rename_all = "kebab-case", deny_unknown_fields)] struct ControlFrame { key: i64, #[serde(default)] payload: Option, #[serde(default)] viewport_width: Option, #[serde(default = "default_true")] viewport_width_known: bool, #[serde(default)] viewport_height: Option, #[serde(default)] root_width: Option, #[serde(default)] root_width_override: bool, #[serde(default)] patch: bool, #[serde(default)] base_viewport_width: Option, #[serde(default = "default_true")] base_viewport_width_known: bool, #[serde(default)] base_viewport_height: Option, #[serde(default)] base_root_width: Option, #[serde(default)] base_root_width_override: bool, #[serde(default)] runtime_revision: u64, #[serde(default)] context_hash: i64, #[serde(default = "default_true")] complete: bool, #[serde(default = "default_true")] root_metadata: bool, #[serde(default)] delay_ms: u64, } #[derive(Debug)] enum JobPayload { Echo(Vec), Layout { document: Arc, context: LayoutContext, root_width: i64, root_width_override: bool, base_context: Option, base_root_width: i64, base_root_width_override: bool, runtime_revision: u64, context_hash: i64, complete: bool, root_metadata: bool, }, } #[derive(Clone, Debug)] struct BaselineIdentity { context: LayoutContext, root_width: i64, root_width_override: bool, runtime_revision: u64, context_hash: i64, complete: bool, } #[derive(Clone, Debug)] struct ConfirmedBaseline { identity: BaselineIdentity, tape: LayoutTape, styles: Vec, } #[derive(Debug)] struct PendingBaseline { confirmed_identity: BaselineIdentity, tape: LayoutTape, styles: Vec, } #[derive(Debug)] struct ResultEntry { bytes: Vec, } #[derive(Debug)] struct RenderedJob { bytes: Vec, pending: Option, baseline_hit: bool, base_renders: u64, target_renders: u64, } #[derive(Debug)] struct Job { generation: u64, key: i64, payload: JobPayload, delay_ms: u64, } #[derive(Debug)] struct PreparedJob { key: i64, delay_ms: u64, payload: JobPayload, } #[derive(Debug)] struct RuntimeState { jobs: VecDeque, results: HashMap<(u64, i64), ResultEntry>, pending_baselines: HashMap<(u64, i64), PendingBaseline>, result_bytes: usize, stale_drops: u64, completed_jobs: u64, baseline_hits: u64, baseline_misses: u64, base_renders: u64, target_renders: u64, confirmed_baseline: Option, } #[derive(Clone, Copy, Debug)] struct ReadinessChannel { fd: c_int, signal: extern "C" fn(c_int) -> bool, close: extern "C" fn(c_int), } #[derive(Debug)] struct Shared { id: u64, alive: AtomicBool, generation: AtomicU64, max_jobs: usize, max_results: usize, max_result_bytes: usize, state: Mutex, readiness_channel: Mutex>, job_available: Condvar, result_available: Condvar, } #[derive(Debug)] struct Session { shared: Arc, workers: Mutex>>>, layout_document: Mutex>>, worker_count: usize, } #[derive(Debug, Serialize)] #[serde(rename_all = "kebab-case")] struct SessionStats { session_id: u64, workers: usize, queued_jobs: usize, ready_results: usize, result_bytes: usize, max_jobs: usize, max_results: usize, max_result_bytes: usize, generation: u64, stale_drops: u64, completed_jobs: u64, baseline_hits: u64, baseline_misses: u64, base_renders: u64, target_renders: u64, pending_baselines: usize, confirmed_baseline: bool, confirmed_baseline_bytes: usize, layout_registered: bool, alive: bool, } impl Shared { fn signal_readiness(&self) { let mut readiness_channel = self .readiness_channel .lock() .unwrap_or_else(|poison| poison.into_inner()); if let Some(channel) = *readiness_channel { if !(channel.signal)(channel.fd) { (channel.close)(channel.fd); *readiness_channel = None; } } } } impl Session { fn new( workers: usize, max_jobs: usize, max_results: usize, max_result_bytes: usize, ) -> Result, String> { if workers == 0 || max_jobs == 0 || max_results == 0 || max_result_bytes == 0 { return Err("Native reflow capacities must be positive".to_owned()); } let available = thread::available_parallelism().map_or(1, usize::from); let worker_count = workers.min(available).min(max_jobs).max(1); let shared = Arc::new(Shared { id: NEXT_SESSION_ID.fetch_add(1, Ordering::Relaxed), alive: AtomicBool::new(true), generation: AtomicU64::new(0), max_jobs, max_results, max_result_bytes, state: Mutex::new(RuntimeState { jobs: VecDeque::new(), results: HashMap::new(), pending_baselines: HashMap::new(), result_bytes: 0, stale_drops: 0, completed_jobs: 0, baseline_hits: 0, baseline_misses: 0, base_renders: 0, target_renders: 0, confirmed_baseline: None, }), readiness_channel: Mutex::new(None), job_available: Condvar::new(), result_available: Condvar::new(), }); let mut handles = Vec::with_capacity(worker_count); for index in 0..worker_count { let worker_shared = Arc::clone(&shared); let spawn = thread::Builder::new() .name(format!("ebox-native-reflow-{index}")) .spawn(move || worker_loop(worker_shared)); match spawn { Ok(handle) => handles.push(handle), Err(error) => { shared.alive.store(false, Ordering::Release); shared.job_available.notify_all(); for handle in handles { let _ = handle.join(); } return Err(format!("Native reflow worker spawn failed: {error}")); } } } Ok(Box::new(Self { shared, workers: Mutex::new(Some(handles)), layout_document: Mutex::new(None), worker_count, })) } fn submit(&self, generation: u64, payload: &[u8]) -> Result { let batch = parse_control_batch(payload)?; if !self.shared.alive.load(Ordering::Acquire) { return Err("Native reflow session is closed".to_owned()); } let layout_requested = batch.document.is_some() || batch.frames.iter().any(|frame| frame.payload.is_none()); let register_document = batch.document.is_some(); let document = match batch.document { Some(document) => { document.validate()?; Some(Arc::new(document)) } None if layout_requested => Some( self.layout_document .lock() .unwrap_or_else(|poison| poison.into_inner()) .clone() .ok_or_else(|| { "Native layout frames require a registered document".to_owned() })?, ), None => None, }; if document.is_some() && self.shared.max_result_bytes < MIN_TAPE_BYTES { return Err(format!( "Native layout results require at least {MIN_TAPE_BYTES} bytes" )); } let mut prepared = Vec::with_capacity(batch.frames.len()); for frame in batch.frames { let job = if let Some(document) = &document { prepare_layout_job(document, frame)? } else { let payload = frame .payload .ok_or_else(|| format!("Native echo frame {} requires payload", frame.key))?; if payload.len() > self.shared.max_result_bytes { return Err(format!( "Native reflow frame {} exceeds the result byte limit", frame.key )); } PreparedJob { key: frame.key, delay_ms: frame.delay_ms, payload: JobPayload::Echo(payload.into_bytes()), } }; prepared.push(job); } let mut state = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); let current = self.shared.generation.load(Ordering::Acquire); if generation < current { return Err(format!( "Native reflow generation {generation} is stale; current generation is {current}" )); } if generation > current { self.shared.generation.store(generation, Ordering::Release); let jobs_before = state.jobs.len(); state.jobs.retain(|job| job.generation >= generation); let removed_jobs = jobs_before - state.jobs.len(); let result_keys: Vec<_> = state .results .keys() .copied() .filter(|(result_generation, _)| *result_generation < generation) .collect(); let mut removed_results = 0; for key in result_keys { if let Some(entry) = state.results.remove(&key) { state.result_bytes = state.result_bytes.saturating_sub(entry.bytes.len()); removed_results += 1; } } state .pending_baselines .retain(|(pending_generation, _), _| *pending_generation >= generation); state.stale_drops += (removed_jobs + removed_results) as u64; self.shared.result_available.notify_all(); } if state.jobs.len() + prepared.len() > self.shared.max_jobs { return Err("Native reflow job queue is full".to_owned()); } let mut batch_keys = HashSet::with_capacity(prepared.len()); for job in &prepared { if !batch_keys.insert(job.key) { return Err(format!( "Duplicate native reflow frame key {} in generation {generation}", job.key )); } let duplicate_job = state .jobs .iter() .any(|queued| queued.generation == generation && queued.key == job.key); if duplicate_job || state.results.contains_key(&(generation, job.key)) { return Err(format!( "Duplicate native reflow frame key {} in generation {generation}", job.key )); } } let accepted = prepared.len(); for job in prepared { state.jobs.push_back(Job { generation, key: job.key, payload: job.payload, delay_ms: job.delay_ms, }); } drop(state); if register_document { *self .layout_document .lock() .unwrap_or_else(|poison| poison.into_inner()) = document; let mut state = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); state.confirmed_baseline = None; state.pending_baselines.clear(); } self.shared.job_available.notify_all(); Ok(accepted) } fn ready(&self, generation: u64, key: i64) -> bool { let state = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); state.results.contains_key(&(generation, key)) } fn take(&self, generation: u64, key: i64) -> Option> { let mut state = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); let result = state.results.remove(&(generation, key)); if let Some(entry) = &result { state.result_bytes = state.result_bytes.saturating_sub(entry.bytes.len()); self.shared.result_available.notify_all(); } result.map(|entry| entry.bytes) } fn render_sync(&self, generation: u64, payload: &[u8]) -> Result, String> { if !self.shared.alive.load(Ordering::Acquire) { return Err("Native reflow session is closed".to_owned()); } let batch = parse_control_batch(payload)?; let document = batch .document .ok_or_else(|| "Native retained render requires an inline document".to_owned())?; document.validate()?; if batch.frames.len() != 1 { return Err("Native retained render requires exactly one frame".to_owned()); } let frame = batch .frames .into_iter() .next() .expect("checked retained frame"); if !frame.complete || frame.delay_ms != 0 { return Err("Native retained render requires one immediate complete frame".to_owned()); } let document = Arc::new(document); let prepared = prepare_layout_job(&document, frame)?; let confirmed = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()) .confirmed_baseline .clone(); let output = render_layout_payload( prepared.payload, self.shared.id, generation, prepared.key, self.shared.max_result_bytes, confirmed, true, ); if output.bytes.len() > self.shared.max_result_bytes { return Err("Native retained render exceeds the result byte limit".to_owned()); } let mut state = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); if output.base_renders > 0 { state.baseline_misses += 1; } else if output.baseline_hit { state.baseline_hits += 1; } state.base_renders += output.base_renders; state.target_renders += output.target_renders; state.pending_baselines.clear(); if let Some(pending) = output.pending { state .pending_baselines .insert((generation, prepared.key), pending); } state.completed_jobs += 1; Ok(output.bytes) } fn confirm(&self, generation: u64, key: i64, confirmed_revision: u64) -> Result { let mut state = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); let Some(pending) = state.pending_baselines.remove(&(generation, key)) else { return Ok(state .confirmed_baseline .as_ref() .is_some_and(|baseline| baseline.identity.runtime_revision == confirmed_revision)); }; if confirmed_revision <= pending.confirmed_identity.runtime_revision { return Err("Native confirmed frame revision mismatch".to_owned()); } let mut identity = pending.confirmed_identity; identity.runtime_revision = confirmed_revision; state.confirmed_baseline = Some(ConfirmedBaseline { identity, tape: pending.tape, styles: pending.styles, }); Ok(true) } fn cancel(&self, generation: u64) { let next = generation.saturating_add(1); self.shared.generation.fetch_max(next, Ordering::AcqRel); let mut state = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); let jobs_before = state.jobs.len(); state.jobs.retain(|job| job.generation > generation); let removed_jobs = jobs_before - state.jobs.len(); let result_keys: Vec<_> = state .results .keys() .copied() .filter(|(result_generation, _)| *result_generation <= generation) .collect(); let mut removed_results = 0; for key in result_keys { if let Some(entry) = state.results.remove(&key) { state.result_bytes = state.result_bytes.saturating_sub(entry.bytes.len()); removed_results += 1; } } state.stale_drops += (removed_jobs + removed_results) as u64; state .pending_baselines .retain(|(pending_generation, _), _| *pending_generation > generation); drop(state); self.shared.job_available.notify_all(); self.shared.result_available.notify_all(); } fn attach_readiness_channel( &self, fd: c_int, signal: extern "C" fn(c_int) -> bool, close: extern "C" fn(c_int), ) -> bool { if fd < 0 || !self.shared.alive.load(Ordering::Acquire) { return false; } let mut channel = self .shared .readiness_channel .lock() .unwrap_or_else(|poison| poison.into_inner()); if channel.is_some() { return false; } *channel = Some(ReadinessChannel { fd, signal, close }); true } fn detach_readiness_channel(&self) { let mut channel = self .shared .readiness_channel .lock() .unwrap_or_else(|poison| poison.into_inner()); if let Some(channel) = channel.take() { (channel.close)(channel.fd); } } fn stats(&self) -> SessionStats { let state = self .shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); let pending_baselines = state.pending_baselines.len(); let confirmed_baseline_bytes = state .confirmed_baseline .as_ref() .map_or(0, |baseline| std::mem::size_of_val(&baseline.tape)); SessionStats { session_id: self.shared.id, workers: self.worker_count, queued_jobs: state.jobs.len(), ready_results: state.results.len(), result_bytes: state.result_bytes, max_jobs: self.shared.max_jobs, max_results: self.shared.max_results, max_result_bytes: self.shared.max_result_bytes, generation: self.shared.generation.load(Ordering::Acquire), stale_drops: state.stale_drops, completed_jobs: state.completed_jobs, baseline_hits: state.baseline_hits, baseline_misses: state.baseline_misses, base_renders: state.base_renders, target_renders: state.target_renders, pending_baselines, confirmed_baseline: state.confirmed_baseline.is_some(), confirmed_baseline_bytes, layout_registered: self .layout_document .lock() .unwrap_or_else(|poison| poison.into_inner()) .is_some(), alive: self.shared.alive.load(Ordering::Acquire), } } fn stop(&self, join: bool) { self.detach_readiness_channel(); self.shared.alive.store(false, Ordering::Release); self.shared.job_available.notify_all(); self.shared.result_available.notify_all(); if let Some(handles) = self .workers .lock() .unwrap_or_else(|poison| poison.into_inner()) .take() { if join { for handle in handles { let _ = handle.join(); } } } } } fn parse_control_batch(payload: &[u8]) -> Result { if payload.len() > MAX_CONTROL_BYTES { return Err("Native reflow control payload exceeds the hard limit".to_owned()); } let batch: ControlBatch = serde_json::from_slice(payload) .map_err(|error| format!("Invalid native reflow control JSON: {error}"))?; if batch.version != CONTROL_VERSION { return Err(format!( "Unsupported native reflow control version {}", batch.version )); } if batch.frames.is_empty() { return Err("Native reflow control batch has no frames".to_owned()); } Ok(batch) } fn checked_layout_context( document: &LayoutDocument, frame_key: i64, viewport_width: i64, viewport_width_known: bool, viewport_height: i64, label: &str, ) -> Result { if viewport_width <= 0 || viewport_height <= 0 || viewport_width > MAX_LAYOUT_DIMENSION || viewport_height > MAX_LAYOUT_DIMENSION { return Err(format!( "Native layout frame {frame_key} {label} viewport is outside the dimension limit" )); } let context = LayoutContext { viewport_width, viewport_width_known, viewport_height, inline_auto_width_intrinsic: false, }; document.validate_context(context)?; Ok(context) } fn checked_root_width(frame_key: i64, root_width: i64, label: &str) -> Result { if root_width <= 0 || root_width > MAX_LAYOUT_DIMENSION { return Err(format!( "Native layout frame {frame_key} {label} root width is outside the dimension limit" )); } Ok(root_width) } fn confirmed_baseline_matches( baseline: &ConfirmedBaseline, context: LayoutContext, root_width: i64, root_width_override: bool, runtime_revision: u64, context_hash: i64, complete: bool, ) -> bool { baseline.identity.context == context && baseline.identity.root_width == root_width && baseline.identity.root_width_override == root_width_override && baseline.identity.runtime_revision == runtime_revision && baseline.identity.context_hash == context_hash && baseline.identity.complete == complete } fn prepare_layout_job( document: &Arc, frame: ControlFrame, ) -> Result { if frame.payload.is_some() { return Err(format!( "Native layout frame {} cannot contain an echo payload", frame.key )); } let viewport_width = frame .viewport_width .ok_or_else(|| format!("Native layout frame {} requires viewport-width", frame.key))?; let viewport_height = frame .viewport_height .ok_or_else(|| format!("Native layout frame {} requires viewport-height", frame.key))?; let context = checked_layout_context( document.as_ref(), frame.key, viewport_width, frame.viewport_width_known, viewport_height, "target", )?; let root_width = checked_root_width( frame.key, frame.root_width.unwrap_or(viewport_width), "target", )?; let (base_context, base_root_width) = if frame.patch { let base_viewport_width = frame.base_viewport_width.ok_or_else(|| { format!( "Native layout patch frame {} requires base-viewport-width", frame.key ) })?; let base_viewport_height = frame.base_viewport_height.ok_or_else(|| { format!( "Native layout patch frame {} requires base-viewport-height", frame.key ) })?; let base_context = checked_layout_context( document.as_ref(), frame.key, base_viewport_width, frame.base_viewport_width_known, base_viewport_height, "base", )?; let base_root_width = checked_root_width( frame.key, frame.base_root_width.unwrap_or(base_viewport_width), "base", )?; (Some(base_context), base_root_width) } else { if frame.base_viewport_width.is_some() || frame.base_viewport_height.is_some() || frame.base_root_width.is_some() || !frame.base_viewport_width_known || frame.base_root_width_override { return Err(format!( "Native full layout frame {} cannot contain base geometry", frame.key )); } (None, root_width) }; Ok(PreparedJob { key: frame.key, delay_ms: frame.delay_ms, payload: JobPayload::Layout { document: Arc::clone(document), context, root_width, root_width_override: frame.root_width_override, base_context, base_root_width, base_root_width_override: frame.base_root_width_override, runtime_revision: frame.runtime_revision, context_hash: frame.context_hash, complete: frame.complete, root_metadata: frame.root_metadata, }, }) } fn render_layout_payload( payload: JobPayload, session_id: u64, generation: u64, key: i64, max_result_bytes: usize, confirmed_baseline: Option, require_confirmed_patch_base: bool, ) -> RenderedJob { match payload { JobPayload::Echo(bytes) => RenderedJob { bytes, pending: None, baseline_hit: false, base_renders: 0, target_renders: 0, }, JobPayload::Layout { document, context, root_width, root_width_override, base_context, base_root_width, base_root_width_override, runtime_revision, context_hash, complete, root_metadata, } => { let identity = TapeIdentity { session_id, generation, key, runtime_revision, context_hash, viewport_width: context.viewport_width, viewport_height: context.viewport_height, root_width, complete, }; let output = TapeOutputOptions { root_metadata, max_bytes: max_result_bytes, }; let pending_identity = BaselineIdentity { context, root_width, root_width_override, runtime_revision, context_hash, complete, }; // Worker threads have no other panic boundary: an unwinding // panic would silently kill the worker and leave the popped // job with neither result nor error tape, so the Elisp ready // watcher would poll forever. Convert panics into the same // error-tape channel ordinary layout failures use. type LayoutRenderOutcome = Result<(Vec, LayoutTape, bool, u64, u64), String>; let result = catch_unwind(AssertUnwindSafe(|| -> LayoutRenderOutcome { if let Some(base_context) = base_context { let base_hit = confirmed_baseline .as_ref() .filter(|baseline| { confirmed_baseline_matches( baseline, base_context, base_root_width, base_root_width_override, runtime_revision, context_hash, complete, ) && baseline.styles == document.styles }) .cloned(); if require_confirmed_patch_base && base_hit.is_none() { let target = document .layout_tape(context, root_width_override.then_some(root_width))?; let bytes = layout::encode_layout_tape( target.clone(), &document.styles, identity, output.root_metadata, output.max_bytes, )?; return Ok((bytes, target, false, 0, 1)); } let (old, baseline_hit, base_renders) = if let Some(baseline) = base_hit { (baseline.tape, true, 0) } else { ( document.layout_tape( base_context, base_root_width_override.then_some(base_root_width), )?, false, 1, ) }; let target = document.layout_tape(context, root_width_override.then_some(root_width))?; let bytes = layout::encode_layout_patch_tape( old, target.clone(), &document.styles, identity, output.root_metadata, output.max_bytes, )?; Ok((bytes, target, baseline_hit, base_renders, 1)) } else { let target = document.layout_tape(context, root_width_override.then_some(root_width))?; let bytes = layout::encode_layout_tape( target.clone(), &document.styles, identity, output.root_metadata, output.max_bytes, )?; Ok((bytes, target, false, 0, 1)) } })); match result { Ok(Ok((bytes, tape, baseline_hit, base_renders, target_renders))) => RenderedJob { bytes, pending: Some(PendingBaseline { confirmed_identity: pending_identity, tape, styles: document.styles.clone(), }), baseline_hit, base_renders, target_renders, }, Ok(Err(error)) => RenderedJob { bytes: encode_error_tape(identity, &error, max_result_bytes), pending: None, baseline_hit: false, base_renders: 0, target_renders: 0, }, Err(_) => RenderedJob { bytes: encode_error_tape(identity, "native layout panicked", max_result_bytes), pending: None, baseline_hit: false, base_renders: 0, target_renders: 0, }, } } } } fn render_proof(payload: &[u8]) -> Result, String> { let batch = parse_control_batch(payload)?; let document = batch .document .ok_or_else(|| "Native proof render requires an inline document".to_owned())?; document.validate()?; if batch.frames.len() != 1 { return Err("Native proof render requires exactly one layout frame".to_owned()); } if SYNC_RENDER_MAX_RESULT_BYTES < MIN_TAPE_BYTES { return Err(format!( "Native layout results require at least {MIN_TAPE_BYTES} bytes" )); } let frame = batch .frames .into_iter() .next() .expect("checked single proof frame"); if frame.patch { return Err("Native proof render supports full layout frames only".to_owned()); } if !frame.complete { return Err("Native proof render requires complete frames".to_owned()); } if frame.delay_ms != 0 { return Err("Native proof render cannot contain delay-ms".to_owned()); } let prepared = prepare_layout_job(&Arc::new(document), frame)?; let output = render_layout_payload( prepared.payload, SYNC_RENDER_SESSION_ID, 0, prepared.key, SYNC_RENDER_MAX_RESULT_BYTES, None, false, ); if output.bytes.len() > SYNC_RENDER_MAX_RESULT_BYTES { return Err("Native proof render exceeds the result byte limit".to_owned()); } Ok(output.bytes) } fn worker_loop(shared: Arc) { loop { let job = { let mut state = shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); while shared.alive.load(Ordering::Acquire) && state.jobs.is_empty() { state = shared .job_available .wait(state) .unwrap_or_else(|poison| poison.into_inner()); } if !shared.alive.load(Ordering::Acquire) { return; } state.jobs.pop_front() }; let Some(job) = job else { continue; }; let mut remaining_delay = job.delay_ms; while remaining_delay > 0 { let slice_ms = remaining_delay.min(2); thread::sleep(Duration::from_millis(slice_ms)); remaining_delay -= slice_ms; if !shared.alive.load(Ordering::Acquire) || shared.generation.load(Ordering::Acquire) != job.generation { break; } } if !shared.alive.load(Ordering::Acquire) || shared.generation.load(Ordering::Acquire) != job.generation { let mut state = shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); state.stale_drops += 1; continue; } let confirmed_baseline = { let state = shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); state.confirmed_baseline.clone() }; let output = render_layout_payload( job.payload, shared.id, job.generation, job.key, shared.max_result_bytes, confirmed_baseline, false, ); if output.bytes.len() > shared.max_result_bytes { let mut state = shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); state.stale_drops += 1; continue; } let output_len = output.bytes.len(); let mut state = shared .state .lock() .unwrap_or_else(|poison| poison.into_inner()); while shared.alive.load(Ordering::Acquire) && shared.generation.load(Ordering::Acquire) == job.generation && (state.results.len() >= shared.max_results || state.result_bytes + output_len > shared.max_result_bytes) { state = shared .result_available .wait(state) .unwrap_or_else(|poison| poison.into_inner()); } if !shared.alive.load(Ordering::Acquire) || shared.generation.load(Ordering::Acquire) != job.generation { state.stale_drops += 1; continue; } state.result_bytes += output_len; if output.base_renders > 0 { state.baseline_misses += 1; } else if output.baseline_hit { state.baseline_hits += 1; } state.base_renders += output.base_renders; state.target_renders += output.target_renders; state.results.insert( (job.generation, job.key), ResultEntry { bytes: output.bytes, }, ); if let Some(pending) = output.pending { state .pending_baselines .insert((job.generation, job.key), pending); } state.completed_jobs += 1; drop(state); shared.signal_readiness(); } } fn session_ref<'a>(pointer: *mut c_void) -> Option<&'a Session> { if pointer.is_null() { None } else { Some(unsafe { &*(pointer as *mut Session) }) } } fn set_error(output: *mut NativeBytes, message: impl Into) { if !output.is_null() { unsafe { *output = NativeBytes::from_error(message); } } } #[no_mangle] pub static plugin_is_GPL_compatible: i32 = 0; extern "C" { fn ebox_module_init_impl(runtime: *mut c_void) -> i32; } #[no_mangle] /// Initialize the module through Emacs's official runtime pointer. /// /// # Safety /// /// `runtime` must be the live `emacs_runtime` pointer supplied by Emacs for /// this module initialization call. pub unsafe extern "C" fn emacs_module_init(runtime: *mut c_void) -> i32 { catch_unwind(AssertUnwindSafe(|| unsafe { ebox_module_init_impl(runtime) })) .unwrap_or(99) } #[no_mangle] pub extern "C" fn ebox_native_layout_ready() -> bool { true } #[no_mangle] pub extern "C" fn ebox_native_session_create( workers: u32, max_jobs: u32, max_results: u32, max_result_bytes: u64, error: *mut NativeBytes, ) -> *mut c_void { match catch_unwind(AssertUnwindSafe(|| { Session::new( workers as usize, max_jobs as usize, max_results as usize, usize::try_from(max_result_bytes) .map_err(|_| "Native reflow byte limit exceeds this platform".to_owned())?, ) })) { Ok(Ok(session)) => Box::into_raw(session) as *mut c_void, Ok(Err(message)) => { set_error(error, message); ptr::null_mut() } Err(_) => { set_error(error, "Native reflow session creation panicked"); ptr::null_mut() } } } #[no_mangle] /// Copy and submit one versioned control payload to a live native session. /// /// # Safety /// /// `session` must come from `ebox_native_session_create`. When `payload_len` /// is nonzero, `payload` must reference that many readable bytes for the /// duration of this call. `error`, when non-null, must be writable. pub unsafe extern "C" fn ebox_native_session_submit( session: *mut c_void, generation: u64, payload: *const u8, payload_len: usize, error: *mut NativeBytes, ) -> i64 { let outcome = catch_unwind(AssertUnwindSafe(|| { let session = session_ref(session).ok_or_else(|| "Native reflow session is null".to_owned())?; if payload.is_null() && payload_len != 0 { return Err("Native reflow payload pointer is null".to_owned()); } let bytes = if payload_len == 0 { &[][..] } else { unsafe { slice::from_raw_parts(payload, payload_len) } }; session.submit(generation, bytes) })); match outcome { Ok(Ok(accepted)) => accepted as i64, Ok(Err(message)) => { set_error(error, message); -1 } Err(_) => { set_error(error, "Native reflow submission panicked"); -1 } } } #[no_mangle] /// Render one retained frame synchronously through a live native session. /// /// # Safety /// /// `session` must come from `ebox_native_session_create`. PAYLOAD and OUTPUT /// follow the same validity rules as `ebox_native_render_proof`. pub unsafe extern "C" fn ebox_native_session_render_sync( session: *mut c_void, generation: u64, payload: *const u8, payload_len: usize, output: *mut NativeBytes, ) -> bool { if output.is_null() { return false; } let outcome = catch_unwind(AssertUnwindSafe(|| { let session = session_ref(session).ok_or_else(|| "Native reflow session is null".to_owned())?; if payload.is_null() && payload_len != 0 { return Err("Native retained render payload pointer is null".to_owned()); } let bytes = if payload_len == 0 { &[][..] } else { unsafe { slice::from_raw_parts(payload, payload_len) } }; session.render_sync(generation, bytes) })); let result = match outcome { Ok(result) => result, Err(_) => Err("Native retained render panicked".to_owned()), }; let success = result.is_ok(); unsafe { *output = match result { Ok(bytes) => NativeBytes::from_vec(bytes), Err(message) => NativeBytes::from_error(message), }; } success } #[no_mangle] /// Render one complete layout proof synchronously from a versioned payload. /// /// # Safety /// /// When `payload_len` is nonzero, `payload` must reference that many readable /// bytes for the duration of this call. `output` must be writable. pub unsafe extern "C" fn ebox_native_render_proof( payload: *const u8, payload_len: usize, output: *mut NativeBytes, ) -> bool { if output.is_null() { return false; } let outcome = catch_unwind(AssertUnwindSafe(|| { if payload.is_null() && payload_len != 0 { return Err("Native proof render payload pointer is null".to_owned()); } let bytes = if payload_len == 0 { &[][..] } else { unsafe { slice::from_raw_parts(payload, payload_len) } }; render_proof(bytes) })); let result = match outcome { Ok(result) => result, Err(_) => Err("Native proof render panicked".to_owned()), }; let success = result.is_ok(); unsafe { *output = match result { Ok(bytes) => NativeBytes::from_vec(bytes), Err(message) => NativeBytes::from_error(message), }; } success } #[no_mangle] pub extern "C" fn ebox_native_session_ready( session: *mut c_void, generation: u64, key: i64, ) -> bool { catch_unwind(AssertUnwindSafe(|| { session_ref(session).is_some_and(|session| session.ready(generation, key)) })) .unwrap_or(false) } #[no_mangle] pub extern "C" fn ebox_native_session_take( session: *mut c_void, generation: u64, key: i64, ) -> NativeBytes { catch_unwind(AssertUnwindSafe(|| { session_ref(session) .and_then(|session| session.take(generation, key)) .map_or_else(NativeBytes::empty, NativeBytes::from_vec) })) .unwrap_or_else(|_| NativeBytes::empty()) } #[no_mangle] pub extern "C" fn ebox_native_session_confirm_frame( session: *mut c_void, generation: u64, key: i64, confirmed_revision: u64, error: *mut NativeBytes, ) -> bool { let outcome = catch_unwind(AssertUnwindSafe(|| { let session = session_ref(session).ok_or_else(|| "Native reflow session is null".to_owned())?; session.confirm(generation, key, confirmed_revision) })); match outcome { Ok(Ok(confirmed)) => confirmed, Ok(Err(message)) => { set_error(error, message); false } Err(_) => { set_error(error, "Native reflow confirmation panicked"); false } } } #[no_mangle] pub extern "C" fn ebox_native_session_attach_readiness_channel( session: *mut c_void, fd: c_int, signal: Option bool>, close: Option, ) -> bool { catch_unwind(AssertUnwindSafe(|| { let Some(session) = session_ref(session) else { return false; }; let Some(signal) = signal else { return false; }; let Some(close) = close else { return false; }; session.attach_readiness_channel(fd, signal, close) })) .unwrap_or(false) } #[no_mangle] pub extern "C" fn ebox_native_session_detach_readiness_channel(session: *mut c_void) { let _ = catch_unwind(AssertUnwindSafe(|| { if let Some(session) = session_ref(session) { session.detach_readiness_channel(); } })); } #[no_mangle] pub extern "C" fn ebox_native_session_cancel(session: *mut c_void, generation: u64) { let _ = catch_unwind(AssertUnwindSafe(|| { if let Some(session) = session_ref(session) { session.cancel(generation); } })); } #[no_mangle] pub extern "C" fn ebox_native_session_stats(session: *mut c_void) -> NativeBytes { catch_unwind(AssertUnwindSafe(|| { session_ref(session) .and_then(|session| serde_json::to_vec(&session.stats()).ok()) .map_or_else(NativeBytes::empty, NativeBytes::from_vec) })) .unwrap_or_else(|_| NativeBytes::empty()) } #[no_mangle] pub extern "C" fn ebox_native_session_release(session: *mut c_void) { if session.is_null() { return; } let _ = catch_unwind(AssertUnwindSafe(|| { let session = unsafe { Box::from_raw(session as *mut Session) }; session.stop(false); })); } #[no_mangle] pub extern "C" fn ebox_native_session_finalize(session: *mut c_void) { if session.is_null() { return; } let _ = catch_unwind(AssertUnwindSafe(|| { let session = unsafe { Box::from_raw(session as *mut Session) }; session.stop(false); })); } #[no_mangle] pub extern "C" fn ebox_native_bytes_free(bytes: NativeBytes) { if bytes.data.is_null() || bytes.len == 0 { return; } let pointer = ptr::slice_from_raw_parts_mut(bytes.data, bytes.len); unsafe { drop(Box::from_raw(pointer)); } } #[cfg(test)] mod tests { use super::*; use std::sync::atomic::AtomicUsize; use std::time::{Duration, Instant}; static READINESS_WAKE_COUNT: AtomicUsize = AtomicUsize::new(0); static READINESS_CLOSE_COUNT: AtomicUsize = AtomicUsize::new(0); static READINESS_FAILURE_CLOSE_COUNT: AtomicUsize = AtomicUsize::new(0); extern "C" fn test_readiness_signal(fd: c_int) -> bool { if fd == 7 { READINESS_WAKE_COUNT.fetch_add(1, Ordering::SeqCst); } true } extern "C" fn test_readiness_signal_failure(_fd: c_int) -> bool { false } extern "C" fn test_readiness_close(fd: c_int) { if fd == 7 { READINESS_CLOSE_COUNT.fetch_add(1, Ordering::SeqCst); } else { READINESS_FAILURE_CLOSE_COUNT.fetch_add(1, Ordering::SeqCst); } } fn batch(frames: &str) -> Vec { format!(r#"{{"version":1,"frames":{frames}}}"#).into_bytes() } fn proof_layout_payload(frames: &str) -> Vec { format!( r#"{{"version":1,"document":{{"version":2,"space-width":8,"style-count":0,"styles":[],"root":{{"type":"box","region-id":1,"content":{{"lines":[{{"clusters":[{{"text":"x","width":8,"cjk":false,"space":false}}]}}]}},"child":null,"content-width-exact":false,"width":{{"kind":"viewport"}},"min-width":{{"kind":"pixels","value":0}},"max-width":{{"kind":"none"}},"height":{{"kind":"auto"}},"min-height":{{"kind":"lines","value":0}},"max-height":{{"kind":"none"}},"box-sizing":"border-box","padding-left":0,"padding-right":0,"padding-top":0,"padding-bottom":0,"margin-left":0,"margin-right":0,"margin-top":0,"margin-bottom":0,"border-left":0,"border-right":0,"foreground-style":null,"background-style":null,"border-left-style":null,"border-right-style":null,"border-top-style":null,"border-bottom-style":null,"text-align":"left","vertical-align":"top","overflow":"scroll","wrap-mode":"word","scroll-offset":0}}}},"frames":{frames}}}"# ) .into_bytes() } fn column_proof_layout_payload(frames: &str) -> Vec { format!( r#"{{"version":1,"document":{{"version":2,"space-width":8,"style-count":0,"styles":[],"root":{{"type":"column","children":[{{"type":"box","region-id":1,"content":{{"lines":[{{"clusters":[{{"text":"x","width":8,"cjk":false,"space":false}}]}}]}},"child":null,"content-width-exact":false,"width":{{"kind":"content"}},"min-width":{{"kind":"pixels","value":0}},"max-width":{{"kind":"none"}},"height":{{"kind":"auto"}},"min-height":{{"kind":"lines","value":0}},"max-height":{{"kind":"none"}},"box-sizing":"border-box","padding-left":0,"padding-right":0,"padding-top":0,"padding-bottom":0,"margin-left":0,"margin-right":0,"margin-top":0,"margin-bottom":0,"border-left":0,"border-right":0,"foreground-style":null,"background-style":null,"border-left-style":null,"border-right-style":null,"border-top-style":null,"border-bottom-style":null,"text-align":"left","vertical-align":"top","overflow":"scroll","wrap-mode":"word","scroll-offset":0}},{{"type":"box","region-id":2,"content":{{"lines":[{{"clusters":[{{"text":"yy","width":16,"cjk":false,"space":false}}]}}]}},"child":null,"content-width-exact":false,"width":{{"kind":"content"}},"min-width":{{"kind":"pixels","value":0}},"max-width":{{"kind":"none"}},"height":{{"kind":"auto"}},"min-height":{{"kind":"lines","value":0}},"max-height":{{"kind":"none"}},"box-sizing":"border-box","padding-left":0,"padding-right":0,"padding-top":0,"padding-bottom":0,"margin-left":0,"margin-right":0,"margin-top":0,"margin-bottom":0,"border-left":0,"border-right":0,"foreground-style":null,"background-style":null,"border-left-style":null,"border-right-style":null,"border-top-style":null,"border-bottom-style":null,"text-align":"left","vertical-align":"top","overflow":"scroll","wrap-mode":"word","scroll-offset":0}}]}}}},"frames":{frames}}}"# ) .into_bytes() } fn read_le_u64(bytes: &[u8], offset: usize) -> u64 { u64::from_le_bytes(bytes[offset..offset + 8].try_into().unwrap()) } fn full_tape_character_count(bytes: &[u8]) -> u64 { read_le_u64(bytes, 92) } fn full_tape_line_widths(bytes: &[u8]) -> Vec { let line_count = u32::from_le_bytes(bytes[120..124].try_into().unwrap()) as usize; (0..line_count) .map(|index| read_le_u64(bytes, 152 + index * 8)) .collect() } fn wait_until(mut predicate: impl FnMut() -> bool) { let deadline = Instant::now() + Duration::from_secs(2); while !predicate() && Instant::now() < deadline { thread::sleep(Duration::from_millis(1)); } assert!(predicate()); } #[test] fn bounded_session_round_trips_owned_payloads() { let session = Session::new(2, 8, 8, 4096).unwrap(); let control = batch(r#"[{"key":1,"payload":"alpha"},{"key":2,"payload":"beta"}]"#); assert_eq!(session.submit(1, &control).unwrap(), 2); wait_until(|| session.ready(1, 1) && session.ready(1, 2)); assert_eq!(session.take(1, 1).unwrap(), b"alpha"); assert_eq!(session.take(1, 2).unwrap(), b"beta"); assert_eq!(session.stats().result_bytes, 0); session.stop(true); } #[test] fn readiness_channel_wakes_on_completion_and_detaches_cleanly() { READINESS_WAKE_COUNT.store(0, Ordering::SeqCst); READINESS_CLOSE_COUNT.store(0, Ordering::SeqCst); let session = Session::new(1, 8, 8, 4096).unwrap(); assert!(session.attach_readiness_channel(7, test_readiness_signal, test_readiness_close)); assert!(!session.attach_readiness_channel(8, test_readiness_signal, test_readiness_close)); let control = batch(r#"[{"key":1,"payload":"alpha"},{"key":2,"payload":"beta","delay-ms":20}]"#); assert_eq!(session.submit(1, &control).unwrap(), 2); wait_until(|| READINESS_WAKE_COUNT.load(Ordering::SeqCst) >= 1); assert!(session.ready(1, 1)); assert_eq!(session.take(1, 1).unwrap(), b"alpha"); session.cancel(1); let wakes_after_cancel = READINESS_WAKE_COUNT.load(Ordering::SeqCst); thread::sleep(Duration::from_millis(30)); assert_eq!( READINESS_WAKE_COUNT.load(Ordering::SeqCst), wakes_after_cancel ); session.detach_readiness_channel(); assert_eq!(READINESS_CLOSE_COUNT.load(Ordering::SeqCst), 1); let detached_wakes = READINESS_WAKE_COUNT.load(Ordering::SeqCst); let detached = batch(r#"[{"key":3,"payload":"gamma"}]"#); assert_eq!(session.submit(2, &detached).unwrap(), 1); wait_until(|| session.ready(2, 3)); thread::sleep(Duration::from_millis(5)); assert_eq!(READINESS_WAKE_COUNT.load(Ordering::SeqCst), detached_wakes); session.stop(true); let raw = ebox_native_session_create(1, 2, 2, 4096, ptr::null_mut()); assert!(!raw.is_null()); assert!(ebox_native_session_attach_readiness_channel( raw, 7, Some(test_readiness_signal), Some(test_readiness_close) )); ebox_native_session_release(raw); assert_eq!(READINESS_CLOSE_COUNT.load(Ordering::SeqCst), 2); } #[test] fn readiness_signal_failure_closes_and_releases_the_channel() { READINESS_FAILURE_CLOSE_COUNT.store(0, Ordering::SeqCst); let session = Session::new(1, 2, 2, 4096).unwrap(); assert!(session.attach_readiness_channel( 9, test_readiness_signal_failure, test_readiness_close )); let control = batch(r#"[{"key":1,"payload":"ready"}]"#); assert_eq!(session.submit(1, &control).unwrap(), 1); wait_until(|| session.ready(1, 1)); wait_until(|| READINESS_FAILURE_CLOSE_COUNT.load(Ordering::SeqCst) == 1); assert!(session.attach_readiness_channel(10, test_readiness_signal, test_readiness_close)); session.detach_readiness_channel(); assert_eq!(READINESS_FAILURE_CLOSE_COUNT.load(Ordering::SeqCst), 2); session.stop(true); } #[test] fn cancellation_drops_delayed_generation() { let session = Session::new(1, 4, 4, 4096).unwrap(); let control = batch(r#"[{"key":1,"payload":"stale","delay-ms":50}]"#); assert_eq!(session.submit(7, &control).unwrap(), 1); session.cancel(7); thread::sleep(Duration::from_millis(80)); assert!(!session.ready(7, 1)); assert!(session.stats().stale_drops > 0); session.stop(true); } #[test] fn malformed_and_over_capacity_batches_are_rejected() { let session = Session::new(1, 1, 1, 4).unwrap(); assert!(session.submit(1, b"not-json").is_err()); let oversized = batch(r#"[{"key":1,"payload":"12345"}]"#); assert!(session.submit(1, &oversized).is_err()); assert_eq!(session.stats().queued_jobs, 0); session.stop(true); } #[test] fn newer_generation_purges_old_capacity_and_duplicate_keys_are_rejected() { let session = Session::new(1, 2, 2, 4096).unwrap(); let old = batch(r#"[{"key":1,"payload":"old","delay-ms":50},{"key":2,"payload":"old"}]"#); assert_eq!(session.submit(3, &old).unwrap(), 2); let replacement = batch(r#"[{"key":7,"payload":"new"}]"#); assert_eq!(session.submit(4, &replacement).unwrap(), 1); wait_until(|| session.ready(4, 7)); assert_eq!(session.take(4, 7).unwrap(), b"new"); assert!(!session.ready(3, 1)); assert!(!session.ready(3, 2)); let duplicate = batch(r#"[{"key":8,"payload":"a"},{"key":8,"payload":"b"}]"#); assert!(session.submit(4, &duplicate).is_err()); session.stop(true); } #[test] fn registered_layout_is_reused_by_later_generations() { let session = Session::new(1, 2, 2, 4096).unwrap(); let first = br#"{"version":1,"document":{"version":2,"space-width":8,"style-count":0,"styles":[],"root":{"type":"box","region-id":1,"content":{"lines":[{"clusters":[{"text":"x","width":8,"cjk":false,"space":false}]}]},"child":null,"content-width-exact":false,"width":{"kind":"viewport"},"min-width":{"kind":"pixels","value":0},"max-width":{"kind":"none"},"height":{"kind":"auto"},"min-height":{"kind":"lines","value":0},"max-height":{"kind":"none"},"box-sizing":"border-box","padding-left":0,"padding-right":0,"padding-top":0,"padding-bottom":0,"margin-left":0,"margin-right":0,"margin-top":0,"margin-bottom":0,"border-left":0,"border-right":0,"foreground-style":null,"background-style":null,"border-left-style":null,"border-right-style":null,"border-top-style":null,"border-bottom-style":null,"text-align":"left","vertical-align":"top","overflow":"scroll","wrap-mode":"word","scroll-offset":0}},"frames":[{"key":1,"viewport-width":80,"viewport-height":10}]}"#; assert_eq!(session.submit(1, first).unwrap(), 1); wait_until(|| session.ready(1, 1)); assert!(!session.take(1, 1).unwrap().is_empty()); let reused = br#"{"version":1,"frames":[{"key":2,"viewport-width":120,"viewport-height":12}]}"#; assert_eq!(session.submit(2, reused).unwrap(), 1); wait_until(|| session.ready(2, 2)); assert!(!session.take(2, 2).unwrap().is_empty()); session.stop(true); } #[test] fn confirmed_baseline_is_promoted_only_after_explicit_confirmation() { let session = Session::new(1, 4, 4, 64 * 1024).unwrap(); let first = proof_layout_payload( r#"[{"key":1,"viewport-width":120,"viewport-height":10,"root-width":120,"patch":true,"base-viewport-width":80,"base-viewport-height":10,"base-root-width":80,"runtime-revision":0,"context-hash":77}]"#, ); assert_eq!(session.submit(1, &first).unwrap(), 1); wait_until(|| session.ready(1, 1)); assert!(!session.confirm(1, 99, 1).unwrap()); let first_stats = session.stats(); assert_eq!(first_stats.baseline_misses, 1); assert_eq!(first_stats.baseline_hits, 0); assert_eq!(first_stats.base_renders, 1); assert_eq!(first_stats.target_renders, 1); assert_eq!(first_stats.pending_baselines, 1); assert!(!first_stats.confirmed_baseline); assert!(!session.take(1, 1).unwrap().is_empty()); let taken_stats = session.stats(); assert_eq!(taken_stats.pending_baselines, 1); assert!(!taken_stats.confirmed_baseline); assert!(session.confirm(1, 1, 0).is_err()); assert!(!session.stats().confirmed_baseline); let first_again = proof_layout_payload( r#"[{"key":2,"viewport-width":120,"viewport-height":10,"root-width":120,"patch":true,"base-viewport-width":80,"base-viewport-height":10,"base-root-width":80,"runtime-revision":0,"context-hash":77}]"#, ); assert_eq!(session.submit(2, &first_again).unwrap(), 1); wait_until(|| session.ready(2, 2)); assert!(session.confirm(2, 2, 1).unwrap()); assert!(session.take(2, 2).unwrap().len() >= MIN_TAPE_BYTES); let second = batch( r#"[{"key":3,"viewport-width":80,"viewport-height":10,"root-width":80,"patch":true,"base-viewport-width":120,"base-viewport-height":10,"base-root-width":120,"runtime-revision":1,"context-hash":77}]"#, ); assert_eq!(session.submit(3, &second).unwrap(), 1); wait_until(|| session.ready(3, 3)); assert!(!session.take(3, 3).unwrap().is_empty()); let final_stats = session.stats(); assert_eq!(final_stats.baseline_misses, 2); assert_eq!(final_stats.baseline_hits, 1); assert_eq!(final_stats.base_renders, 2); assert_eq!(final_stats.target_renders, 3); session.stop(true); } #[test] fn synchronous_proof_render_returns_full_tape_without_session() { let payload = proof_layout_payload(r#"[{"key":9,"viewport-width":80,"viewport-height":10}]"#); let tape = render_proof(&payload).unwrap(); assert!(tape.len() >= MIN_TAPE_BYTES); assert_eq!(read_le_u64(&tape, 20), SYNC_RENDER_SESSION_ID); } #[test] fn synchronous_proof_render_distinguishes_known_and_unknown_width() { let known = render_proof(&column_proof_layout_payload( r#"[{"key":1,"viewport-width":80,"viewport-height":10}]"#, )) .unwrap(); let unknown = render_proof(&column_proof_layout_payload( r#"[{"key":1,"viewport-width":80,"viewport-width-known":false,"viewport-height":10}]"#, )) .unwrap(); assert_eq!(full_tape_character_count(&unknown), 5); assert_eq!(full_tape_character_count(&known), 6); assert_eq!(full_tape_line_widths(&unknown), vec![16, 16]); assert_eq!(full_tape_line_widths(&known), vec![80, 80]); } #[test] fn synchronous_proof_render_rejects_non_proof_batches() { let echo = batch(r#"[{"key":1,"payload":"not-layout"}]"#); assert!(render_proof(&echo) .unwrap_err() .contains("requires an inline document")); let multiple = proof_layout_payload( r#"[{"key":1,"viewport-width":80,"viewport-height":10},{"key":2,"viewport-width":80,"viewport-height":10}]"#, ); assert!(render_proof(&multiple) .unwrap_err() .contains("exactly one layout frame")); let patch = proof_layout_payload( r#"[{"key":1,"viewport-width":80,"viewport-height":10,"patch":true,"base-viewport-width":80,"base-viewport-height":10}]"#, ); assert!(render_proof(&patch) .unwrap_err() .contains("full layout frames only")); let incomplete = proof_layout_payload( r#"[{"key":1,"viewport-width":80,"viewport-height":10,"complete":false}]"#, ); assert!(render_proof(&incomplete) .unwrap_err() .contains("requires complete frames")); let full_with_base = proof_layout_payload( r#"[{"key":1,"viewport-width":80,"viewport-height":10,"base-viewport-width-known":false}]"#, ); assert!(render_proof(&full_with_base) .unwrap_err() .contains("cannot contain base geometry")); } }