use std::sync::OnceLock;
use std::sync::mpsc::{Receiver, SyncSender};
use cljrs_reader::Form;
use std::sync::Arc;
pub(crate) struct LowerArityRequest {
pub arity_id: u64,
pub params: Vec<Arc<str>>,
pub rest_param: Option<Arc<str>>,
pub destructure_params: Vec<(usize, Form)>,
pub destructure_rest: Option<Form>,
pub expanded_body: Vec<Form>,
pub param_hints: Vec<Option<cljrs_value::TypeHint>>,
}
pub(crate) struct LowerRequest {
pub tiers: std::sync::Weak<crate::tiered::tiers::Tiers>,
pub name: Option<Arc<str>>,
pub ns: Arc<str>,
pub is_async: bool,
pub arities: Vec<LowerArityRequest>,
}
const MAX_LOWER_ATTEMPTS: u32 = 3;
static SENDER: OnceLock<SyncSender<LowerRequest>> = OnceLock::new();
pub(crate) fn enqueue(req: LowerRequest) -> bool {
let tx = SENDER.get_or_init(|| {
let (tx, rx) = std::sync::mpsc::sync_channel::<LowerRequest>(256);
crate::env::gc_roots::set_stw_reclaim_hook(|| {
crate::tiered::tiers::sweep_idle(
crate::tiered::ir_cache::now_secs(),
crate::tiered::ir_cache::ir_cache_ttl_secs(),
);
});
crate::tiered::defn_registry::install_invalidation_hook();
std::thread::Builder::new()
.name("cljrs-ir-lower".into())
.spawn(move || worker_loop(rx))
.expect("failed to spawn IR lowering worker thread");
tx
});
tx.try_send(req).is_ok()
}
fn worker_loop(rx: Receiver<LowerRequest>) {
for req in &rx {
let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
process_request(&req);
}));
if result.is_err() {
cljrs_logging::feat_debug!(
"ir",
"background lower panicked for {:?}; staying at Tier 0",
req.name
);
}
}
}
fn process_request(req: &LowerRequest) {
let Some(tiers) = req.tiers.upgrade() else {
cljrs_logging::feat_debug!("ir", "dropping lower request for dead runtime");
return;
};
let ir_cache = tiers.ir_cache();
let globals_id = tiers.globals_id();
let mut registered: Vec<(usize, bool, Arc<cljrs_ir::IrFunction>)> = Vec::new();
for arity in &req.arities {
let id = arity.arity_id;
if !ir_cache.should_attempt(id) && !crate::tiered::defn_registry::relower_marked(id) {
if let Some(ir) = ir_cache.get(id) {
registered.push((arity.params.len(), arity.rest_param.is_some(), ir));
}
continue;
}
let mut terminal = false;
for _ in 0..MAX_LOWER_ATTEMPTS {
crate::tiered::defn_registry::take_relower(id);
ir_cache.invalidate(id);
match crate::tiered::lower::lower_expanded_arity(
req.name.as_deref(),
&arity.params,
arity.rest_param.as_ref(),
&arity.destructure_params,
arity.destructure_rest.as_ref(),
&arity.expanded_body,
&req.ns,
globals_id,
Some(id),
true,
req.is_async,
) {
Ok((mut ir, _used)) => {
ir.seed_reprs = crate::tiered::lower::seed_reprs_from_hints(&arity.param_hints);
let ir = Arc::new(ir);
ir_cache.store(id, ir.clone());
if crate::tiered::defn_registry::relower_marked(id) {
ir_cache.invalidate(id);
continue;
}
tiers.jit().on_ir_published(id);
registered.push((arity.params.len(), arity.rest_param.is_some(), ir));
cljrs_logging::feat_debug!(
"ir",
"background lower published arity_id={} ({:?})",
id,
req.name
);
terminal = true;
break;
}
Err(e) => {
ir_cache.store_unsupported(id);
cljrs_logging::feat_debug!(
"ir",
"background lower unsupported arity_id={} ({:?}): {}",
id,
req.name,
e
);
terminal = true;
break;
}
}
}
if !terminal {
ir_cache.invalidate(id);
tiers.jit().clear_lower_queued(id);
cljrs_logging::feat_debug!(
"ir",
"background lower abandoned after {} rebind retries arity_id={} ({:?})",
MAX_LOWER_ATTEMPTS,
id,
req.name
);
}
}
if !req.is_async
&& !registered.is_empty()
&& let Some(name) = &req.name
{
crate::tiered::defn_registry::register_defn(globals_id, &req.ns, name, registered);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn span() -> cljrs_types::span::Span {
cljrs_types::span::Span::new(Arc::new("<test>".to_string()), 0, 0, 1, 1)
}
fn sym(name: &str) -> Form {
Form::new(cljrs_reader::form::FormKind::Symbol(name.into()), span())
}
fn identity_request(
arity_id: u64,
fn_name: &str,
ns: &str,
tiers: &Arc<crate::tiered::tiers::Tiers>,
) -> LowerRequest {
LowerRequest {
tiers: tiers.handle(),
name: Some(Arc::from(fn_name)),
ns: Arc::from(ns),
is_async: false,
arities: vec![LowerArityRequest {
arity_id,
params: vec![Arc::from("x")],
rest_param: None,
destructure_params: Vec::new(),
destructure_rest: None,
expanded_body: vec![sym("x")],
param_hints: Vec::new(),
}],
}
}
#[test]
fn process_request_publishes_and_registers() {
let _sweep_guard = crate::tiered::tiers::SWEEP_TEST_LOCK
.lock()
.unwrap_or_else(|p| p.into_inner());
let id = 0xC500_0001u64;
let gid = 0xC500_0001u64;
let tiers = crate::tiered::tiers::Tiers::new(gid);
let ir_cache = tiers.ir_cache();
let ns = "test.worker-ns";
let req = identity_request(id, "ident", ns, &tiers);
assert!(ir_cache.should_attempt(id));
process_request(&req);
assert!(ir_cache.get(id).is_some());
let mut referenced = std::collections::HashSet::new();
referenced.insert((Arc::<str>::from(ns), Arc::<str>::from("ident")));
let externals = crate::tiered::defn_registry::externals_for(gid, &referenced);
assert_eq!(externals.len(), 1);
let before = ir_cache.get(id).unwrap();
process_request(&req);
let after = ir_cache.get(id).unwrap();
assert!(Arc::ptr_eq(&before, &after));
}
#[test]
fn process_request_relowers_marked_arity_despite_cache_hit() {
let _sweep_guard = crate::tiered::tiers::SWEEP_TEST_LOCK
.lock()
.unwrap_or_else(|p| p.into_inner());
let id = 0xC500_0011u64;
let gid = 0xC500_0011u64;
let tiers = crate::tiered::tiers::Tiers::new(gid);
let ir_cache = tiers.ir_cache();
let callee_ns: Arc<str> = Arc::from("test.worker-relower-ns");
let callee: Arc<str> = Arc::from("callee");
let stale = Arc::new(cljrs_ir::IrFunction::new(None, None));
ir_cache.store(id, stale.clone());
crate::tiered::defn_registry::record_dependents(
id,
vec![(callee_ns.clone(), callee.clone())],
);
crate::tiered::defn_registry::on_redefined(&callee_ns, &callee);
assert!(crate::tiered::defn_registry::relower_marked(id));
process_request(&identity_request(
id,
"dependent",
"test.worker-relower-ns",
&tiers,
));
let fresh = ir_cache.get(id).expect("re-lowered IR");
assert!(!Arc::ptr_eq(&fresh, &stale));
assert!(!crate::tiered::defn_registry::relower_marked(id));
}
}