use std::sync::Arc;
use std::sync::OnceLock;
use std::time::SystemTime;
use arrow_schema::DataType;
use datafusion::logical_expr::{ColumnarValue, Volatility};
use uni_plugin::scheduler::Scheduler;
use uni_plugin::traits::background::Schedule;
use uni_plugin::traits::scalar::{ArgType, FnSignature, ScalarPluginFn};
use uni_plugin::{
Capability, CapabilitySet, FnError, PluginId, PluginRegistrar, PluginRegistry, QName,
};
struct NoopScalar;
impl ScalarPluginFn for NoopScalar {
fn signature(&self) -> &FnSignature {
static SIG: OnceLock<FnSignature> = OnceLock::new();
SIG.get_or_init(|| {
FnSignature::new(
vec![ArgType::Primitive(DataType::Float64)],
ArgType::Primitive(DataType::Float64),
Volatility::Immutable,
)
})
}
fn invoke(&self, args: &[ColumnarValue], _rows: usize) -> Result<ColumnarValue, FnError> {
Ok(args[0].clone())
}
}
fn scalar_sig() -> FnSignature {
FnSignature::new(
vec![ArgType::Primitive(DataType::Float64)],
ArgType::Primitive(DataType::Float64),
Volatility::Immutable,
)
}
#[test]
fn repro_apply_pending_overwrites_ownership_record() {
let registry = PluginRegistry::new();
let caps = CapabilitySet::from_iter_of([Capability::ScalarFn]);
let pid = PluginId::new("mycorp");
{
let mut r = PluginRegistrar::new(pid.clone(), &caps, ®istry);
r.scalar_fn(
QName::new("mycorp", "f1"),
scalar_sig(),
Arc::new(NoopScalar),
)
.unwrap();
r.commit_to_registry().unwrap();
}
{
let mut r = PluginRegistrar::new(pid.clone(), &caps, ®istry);
r.scalar_fn(
QName::new("mycorp", "f2"),
scalar_sig(),
Arc::new(NoopScalar),
)
.unwrap();
r.commit_to_registry().unwrap();
}
assert!(registry.scalar_fn(&QName::new("mycorp", "f1")).is_some());
assert!(registry.scalar_fn(&QName::new("mycorp", "f2")).is_some());
let snap = registry
.iter_for_plugin(&pid)
.expect("plugin record exists");
assert_eq!(
snap.scalars.len(),
2,
"ownership record must contain both commits' scalars (merged, not overwritten)"
);
assert!(snap.scalars.contains(&QName::new("mycorp", "f1")));
assert!(snap.scalars.contains(&QName::new("mycorp", "f2")));
registry.remove_plugin(&pid);
assert!(
registry.scalar_fn(&QName::new("mycorp", "f2")).is_none(),
"f2 removed"
);
assert!(
registry.scalar_fn(&QName::new("mycorp", "f1")).is_none(),
"f1 must also be removed — merged record tracks it, so no orphan leaks"
);
}
#[test]
fn repro_abi_range_upper_bounded_minor_reports_unsupported() {
use uni_plugin::AbiRange;
assert!(
AbiRange::parse("~1.2").unwrap().matches(1),
"~1.2 must report host major 1 as supported"
);
assert!(
AbiRange::parse("=1.2.3").unwrap().matches(1),
"=1.2.3 must report host major 1 as supported"
);
assert!(
AbiRange::parse(">=1.2, <1.6").unwrap().matches(1),
">=1.2,<1.6 must report host major 1 as supported"
);
assert!(
!AbiRange::parse("~1.2").unwrap().matches(2),
"~1.2 must reject host major 2"
);
assert!(AbiRange::parse("^1.2").unwrap().matches(1));
assert!(AbiRange::parse("^1").unwrap().matches(1));
}
#[test]
fn repro_intra_batch_duplicate_bypasses_preflight() {
let registry = PluginRegistry::new();
let caps = CapabilitySet::from_iter_of([Capability::ScalarFn]);
let pid = PluginId::new("dupco");
let mut r = PluginRegistrar::new(pid.clone(), &caps, ®istry);
r.scalar_fn(
QName::new("dupco", "myfn"),
scalar_sig(),
Arc::new(NoopScalar),
)
.unwrap();
r.scalar_fn(
QName::new("dupco", "myfn"),
scalar_sig(),
Arc::new(NoopScalar),
)
.unwrap();
let result = r.commit_to_registry();
assert!(
matches!(&result, Err(uni_plugin::PluginError::DuplicateRegistration(q)) if *q == QName::new("dupco", "myfn")),
"intra-batch duplicate qname must be rejected as DuplicateRegistration, got {result:?}"
);
assert!(registry.scalar_fn(&QName::new("dupco", "myfn")).is_none());
assert!(registry.iter_for_plugin(&pid).is_none());
}
#[test]
fn repro_unparseable_cron_dispatched_once() {
let s = Scheduler::new();
s.resume();
s.add_scheduled_job(
QName::builtin("bad_cron"),
Schedule::Cron(smol_str::SmolStr::new("this is not a valid cron")),
);
let jobs = s.list();
assert_eq!(jobs.len(), 1);
assert!(
jobs[0].next_fire_at.is_none(),
"unparseable cron has no computed fire time"
);
let due = s.tick_at(SystemTime::now());
assert!(
due.is_empty(),
"unparseable-cron job must be skipped (never due), got {due:?}"
);
}