use anyhow::Context;
pub(crate) const ENCODER_RESIDENT_MULTIPLIER: u64 = 2;
pub(crate) const POOL_RAM_FRACTION_DENOM: u64 = 2;
#[cfg(any(test, target_os = "linux"))]
pub(crate) fn parse_cgroup_memory_limit(raw: &str) -> Option<u64> {
let s = raw.trim();
if s.is_empty() || s.eq_ignore_ascii_case("max") {
return None;
}
let bytes: u64 = s.parse().ok()?;
if bytes == 0 || bytes >= (1u64 << 62) {
return None;
}
Some(bytes)
}
pub(crate) fn cgroup_memory_limit_bytes() -> Option<u64> {
#[cfg(target_os = "linux")]
{
const CANDIDATES: &[&str] = &[
"/sys/fs/cgroup/memory.max", "/sys/fs/cgroup/memory/memory.limit_in_bytes", ];
for path in CANDIDATES {
if let Ok(raw) = std::fs::read_to_string(path)
&& let Some(bytes) = parse_cgroup_memory_limit(&raw)
{
return Some(bytes);
}
}
None
}
#[cfg(not(target_os = "linux"))]
{
None
}
}
pub(crate) fn total_ram_bytes() -> u64 {
#[cfg(target_os = "macos")]
{
let mut mem: u64 = 0;
let mut len = std::mem::size_of::<u64>();
let mib = [libc::CTL_HW, libc::HW_MEMSIZE];
let rc = unsafe {
libc::sysctl(
mib.as_ptr() as *mut libc::c_int,
mib.len() as libc::c_uint,
&mut mem as *mut u64 as *mut libc::c_void,
&mut len,
std::ptr::null_mut(),
0,
)
};
if rc == 0 { mem } else { 0 }
}
#[cfg(all(unix, not(target_os = "macos")))]
{
let pages = unsafe { libc::sysconf(libc::_SC_PHYS_PAGES) };
let page_size = unsafe { libc::sysconf(libc::_SC_PAGESIZE) };
if pages > 0 && page_size > 0 {
(pages as u64).saturating_mul(page_size as u64)
} else {
0
}
}
#[cfg(not(unix))]
{
0
}
}
pub(crate) fn effective_ram_bytes() -> u64 {
let host = total_ram_bytes();
match cgroup_memory_limit_bytes() {
Some(limit) if limit > 0 => {
if host == 0 {
limit
} else {
host.min(limit)
}
}
_ => host,
}
}
pub(crate) fn batch_split_count(n: usize, batch_pool_size: usize) -> usize {
batch_pool_size.min(n.saturating_sub(1))
}
pub(crate) fn cap_pool_size_for_ram(requested: usize, encoder_bytes: u64, total_ram: u64) -> usize {
if requested <= 1 || encoder_bytes == 0 || total_ram == 0 {
return requested.max(1);
}
let per_triplet = encoder_bytes.saturating_mul(ENCODER_RESIDENT_MULTIPLIER);
let budget = total_ram / POOL_RAM_FRACTION_DENOM;
let max_slots = (budget / per_triplet.max(1)).max(1) as usize;
if max_slots < requested {
tracing::warn!(
"Capping pool size {requested} -> {max_slots}: \
{requested} encoder slots (~{} MiB each) would exceed half of \
{} MiB total RAM. Concurrency is reduced; add RAM or lower \
--pool-size to silence this.",
per_triplet / (1024 * 1024),
total_ram / (1024 * 1024),
);
max_slots
} else {
requested
}
}
pub(crate) fn clamp_encoder_intra_threads(
pool_size: usize,
requested: usize,
logical_cpus: usize,
) -> usize {
let requested = requested.max(1);
let pool_size = pool_size.max(1);
let logical_cpus = logical_cpus.max(1);
let max_per_encoder = (logical_cpus / pool_size).max(1);
if requested > max_per_encoder {
tracing::warn!(
"Capping encoder intra-op threads {requested} -> {max_per_encoder}: \
{pool_size} pooled encoder(s) x {requested} threads would exceed \
the {logical_cpus} logical CPU(s) available. Lower --pool-size or \
--encoder-intra-threads to silence this."
);
max_per_encoder
} else {
requested
}
}
pub(crate) fn split_pool_items<T: Send>(
mut items: Vec<T>,
batch_pool_size: usize,
) -> (Vec<T>, Option<Vec<T>>) {
let n = items.len();
let batch = batch_split_count(n, batch_pool_size);
if batch == 0 {
return (items, None);
}
let batch_items = items.split_off(n - batch);
(items, Some(batch_items))
}
#[cfg_attr(not(any(test, feature = "coreml")), allow(dead_code))]
pub(crate) fn probe_or_rebuild<S>(
state: S,
probe: impl Fn(&S) -> anyhow::Result<()>,
rebuild: impl FnOnce(S, anyhow::Error) -> anyhow::Result<S>,
) -> anyhow::Result<S> {
match probe(&state) {
Ok(()) => Ok(state),
Err(probe_err) => {
let rebuilt = rebuild(state, probe_err)?;
probe(&rebuilt).context("state failed probe even after rebuild")?;
Ok(rebuilt)
}
}
}
pub(crate) fn finalize_pool_load<T>(
results: Vec<anyhow::Result<T>>,
pool_size: usize,
min_size: usize,
) -> anyhow::Result<Vec<T>> {
let min_size = min_size.clamp(1, pool_size.max(1));
let mut loaded = Vec::with_capacity(results.len());
let mut first_err: Option<anyhow::Error> = None;
for r in results {
match r {
Ok(t) => loaded.push(t),
Err(e) => {
if first_err.is_none() {
first_err = Some(e);
}
}
}
}
let n = loaded.len();
if n >= min_size {
if n < pool_size {
let detail = first_err
.map(|e| format!("; first error: {e:#}"))
.unwrap_or_default();
tracing::warn!(
"degraded pool: loaded {n}/{pool_size} session triplets ({} failed){detail}",
pool_size - n
);
}
Ok(loaded)
} else {
let detail = first_err.map(|e| format!(": {e:#}")).unwrap_or_default();
Err(anyhow::anyhow!(
"loaded only {n}/{pool_size} session triplets, need at least {min_size}{detail}"
))
}
}
#[cfg(test)]
mod tests;