#[cfg(all(target_os = "linux", target_arch = "x86_64"))]
use named_sem::NamedSemaphore;
use zisk_common::{stats_begin, stats_end, stats_mark, ExecutorStatsHandle, Plan};
use std::ffi::c_void;
use std::sync::atomic::{fence, Ordering};
use tracing::{error, warn};
use crate::SEM_CHUNK_DONE_WAIT_DURATION;
use crate::TRACE_DELTA_SIZE;
use crate::TRACE_INITIAL_SIZE;
use crate::TRACE_MAX_SIZE;
use crate::{
sem_chunk_done_name, shmem_output_name, AsmMOChunk, AsmMOHeader, AsmMultiShmem, AsmRunError,
AsmService, AsmServices, GpuBufferSource,
};
#[cfg(gpu)]
use proofman_util::{timer_start_info, timer_stop_and_log_info};
#[cfg(gpu)]
use zisk_sm_mem_planner::GpuCountAndPlan;
use zisk_sm_mem_planner::MemPlanner;
use anyhow::{Context, Result};
#[cfg(feature = "save_mem_plans")]
use zisk_sm_mem_common::save_plans;
#[cfg(gpu)]
fn register_mo_shmem_pinned(
gpu_count_and_plan: &GpuCountAndPlan,
shmem: &AsmMultiShmem<AsmMOHeader>,
registered: &mut usize,
) {
if *registered == usize::MAX {
return; }
let total = shmem.total_mapped_size();
if total <= *registered {
return;
}
let new_ptr = unsafe { (shmem.mapped_ptr() as *const c_void).add(*registered) };
if gpu_count_and_plan.register_input_pinned(new_ptr, total - *registered) {
*registered = total;
} else {
*registered = usize::MAX; }
}
#[cfg(gpu)]
fn setup_gpu_count_and_plan(gpu_buffer: GpuBufferSource) -> Option<GpuCountAndPlan> {
let (d_buf, bytes, gpu_id): (*mut c_void, usize, i32) = match gpu_buffer {
GpuBufferSource::Cpu => {
tracing::info!("[gpu] no GPU buffer requested; using CPU mem_planner path");
return None;
}
GpuBufferSource::Borrowed { ptr, size, .. } if ptr == 0 || size == 0 => {
tracing::info!(
"[gpu] borrowed buffer is empty (--gpu not set at runtime); using CPU mem_planner path"
);
return None;
}
GpuBufferSource::Borrowed { ptr, size, gpu_id } => {
let Ok(gpu_id) = i32::try_from(gpu_id) else {
tracing::error!(
"[gpu] gpu_id {gpu_id} exceeds i32::MAX; using CPU mem_planner path"
);
return None;
};
(ptr as *mut c_void, size, gpu_id)
}
GpuBufferSource::SelfAllocated => (std::ptr::null_mut(), 0, -1),
};
let gpu_count_and_plan = GpuCountAndPlan::new();
if !unsafe { gpu_count_and_plan.setup(d_buf, bytes, 1, 0, gpu_id) } {
tracing::error!("[gpu] GpuCountAndPlan::setup returned false; falling back to CPU");
return None;
}
match gpu_buffer {
GpuBufferSource::SelfAllocated => {
tracing::info!("[gpu] GpuCountAndPlan set up (self-allocated device buffer)");
}
_ => {
tracing::info!(
"[gpu] GpuCountAndPlan set up (borrowed {:.3} GB)",
bytes as f64 / (1024.0 * 1024.0 * 1024.0),
);
}
}
Some(gpu_count_and_plan)
}
pub struct MOShmemReader {
pub(crate) output_shmem: AsmMultiShmem<AsmMOHeader>,
mem_planner: Option<MemPlanner>,
handle_mo: Option<std::thread::JoinHandle<MemPlanner>>,
#[cfg(gpu)]
gpu_count_and_plan: Option<GpuCountAndPlan>,
#[cfg(gpu)]
registered_bytes: usize,
}
impl MOShmemReader {
pub fn new(
shm_prefix: &str,
unlock_mapped_memory: bool,
buffer_source: GpuBufferSource,
) -> Result<Self> {
let output_name = shmem_output_name(shm_prefix, AsmService::MO, None);
let output_shared_memory = AsmMultiShmem::<AsmMOHeader>::open_and_map(
&output_name,
TRACE_INITIAL_SIZE,
TRACE_DELTA_SIZE,
TRACE_MAX_SIZE,
unlock_mapped_memory,
cfg!(gpu),
)?;
#[cfg(gpu)]
let gpu_count_and_plan = setup_gpu_count_and_plan(buffer_source);
#[cfg(not(gpu))]
let _ = buffer_source;
Ok(Self {
output_shmem: output_shared_memory,
mem_planner: Some(MemPlanner::new()),
handle_mo: None,
#[cfg(gpu)]
gpu_count_and_plan,
#[cfg(gpu)]
registered_bytes: usize::MAX,
})
}
}
impl Drop for MOShmemReader {
fn drop(&mut self) {
if let Some(handle_mo) = self.handle_mo.take() {
match handle_mo.join() {
Ok(mem_planner) => {
drop(mem_planner);
}
Err(e) => {
eprintln!("Warning: background thread panicked in MOShmemReader: {e:?}");
}
}
}
}
}
pub struct AsmRunnerMO {
pub plans: Vec<Plan>,
pub gpu_mops_used_bytes: Option<u64>,
}
impl AsmRunnerMO {
pub fn new(plans: Vec<Plan>) -> Self {
Self { plans, gpu_mops_used_bytes: None }
}
#[allow(clippy::too_many_arguments)]
pub fn run<R>(
preloaded: &mut MOShmemReader,
max_steps: u64,
chunk_size: u64,
on_runner_failure: R,
asm_services: AsmServices,
_stats: ExecutorStatsHandle,
) -> Result<Self>
where
R: FnOnce() -> Result<()>,
{
stats_begin!(_stats, 0, _runner_scope, "ASM_MO_RUNNER", 0);
let sem_chunk_done_name = sem_chunk_done_name(asm_services.sem_prefix(), AsmService::MO);
let mut sem_chunk_done = NamedSemaphore::create(sem_chunk_done_name.clone(), 0)
.map_err(|e| AsmRunError::SemaphoreError(sem_chunk_done_name.clone(), e))?;
let stale = crate::drain_chunk_done(&mut sem_chunk_done);
if stale > 0 {
warn!(
"MO semaphore '{sem_chunk_done_name}' had {stale} stale chunk_done post(s) at run start; a prior run skipped its end-side cleanup"
);
}
let _parent_id = _runner_scope.id();
let _thread_stats = _stats.clone();
let handle = std::thread::spawn(move || {
stats_begin!(_thread_stats, _parent_id, _mo_scope, "ASM_MO", 0);
#[allow(clippy::let_and_return)]
let result = asm_services.send_memory_ops_request(max_steps, chunk_size);
stats_end!(_thread_stats, &_mo_scope);
result
});
let mem_planner = match preloaded.mem_planner.take() {
Some(p) => p,
None => preloaded
.handle_mo
.take()
.ok_or_else(|| {
anyhow::anyhow!("MOShmemReader: both mem_planner and handle_mo are None")
})?
.join()
.map_err(|_| anyhow::anyhow!("MO preload background thread panicked"))?,
};
#[cfg(gpu)]
let gpu_count_and_plan_opt: Option<GpuCountAndPlan> = preloaded.gpu_count_and_plan.take();
let mut data_ptr = preloaded.output_shmem.data_ptr() as *const AsmMOChunk;
#[cfg(gpu)]
if let Some(ref gpu_count_and_plan) = gpu_count_and_plan_opt {
gpu_count_and_plan.reset();
register_mo_shmem_pinned(
gpu_count_and_plan,
&preloaded.output_shmem,
&mut preloaded.registered_bytes,
);
} else {
mem_planner.execute();
}
#[cfg(not(gpu))]
mem_planner.execute();
stats_begin!(_stats, &_runner_scope, _process_scope, "MO_PROCESS_CHUNKS", 0);
const MAX_MTRACE_REGS_ACCESS_SIZE: usize = (2 + 2 + 3) * 8;
const MAX_BYTES_DIRECT_MTRACE: usize = 256;
const MAX_BYTES_MTRACE_STEP: usize = MAX_BYTES_DIRECT_MTRACE + MAX_MTRACE_REGS_ACCESS_SIZE;
const MAX_TRACE_CHUNK_INFO: usize = (44 * 8) + 32;
let threshold_bytes = (chunk_size as usize * MAX_BYTES_MTRACE_STEP) + MAX_TRACE_CHUNK_INFO;
let mut threshold = unsafe {
preloaded
.output_shmem
.mapped_ptr()
.add(preloaded.output_shmem.total_mapped_size() - threshold_bytes)
as *const AsmMOChunk
};
let mut on_runner_failure = Some(on_runner_failure);
let mut signal_runner_failure = || {
if let Some(on_failure) = on_runner_failure.take() {
if let Err(reset_err) = on_failure() {
error!("MO on_runner_failure failed: {reset_err:#}");
}
}
};
let loop_result: Result<u64> = loop {
match sem_chunk_done.timed_wait(SEM_CHUNK_DONE_WAIT_DURATION) {
Ok(()) => {
fence(Ordering::Acquire);
if data_ptr >= threshold {
match preloaded.output_shmem.check_size_changed() {
Ok(true) => {
threshold = unsafe {
preloaded.output_shmem.mapped_ptr().add(
preloaded.output_shmem.total_mapped_size()
- threshold_bytes,
) as *const AsmMOChunk
};
#[cfg(gpu)]
if let Some(ref gpu_count_and_plan) = gpu_count_and_plan_opt {
register_mo_shmem_pinned(
gpu_count_and_plan,
&preloaded.output_shmem,
&mut preloaded.registered_bytes,
);
}
}
Ok(false) => {}
Err(e) => {
signal_runner_failure();
break Err(e).context(
"Failed to check and map new shared memory files for MO trace",
);
}
}
}
let chunk = unsafe { std::ptr::read(data_ptr) };
data_ptr = unsafe { data_ptr.add(1) };
stats_mark!(_stats, &_runner_scope, "MO_CHUNK_DONE", 0);
#[cfg(gpu)]
if let Some(ref gpu_count_and_plan) = gpu_count_and_plan_opt {
if !gpu_count_and_plan
.add_chunk(chunk.mem_ops_size, data_ptr as *const c_void)
{
tracing::error!("[gpu] add_chunk failed (n={})", chunk.mem_ops_size);
}
} else {
mem_planner.add_chunk(chunk.mem_ops_size, data_ptr as *const c_void);
}
#[cfg(not(gpu))]
mem_planner.add_chunk(chunk.mem_ops_size, data_ptr as *const c_void);
if chunk.end == 1 {
break Ok(0);
}
data_ptr = unsafe {
(data_ptr as *mut u64).add(chunk.mem_ops_size as usize) as *const AsmMOChunk
};
}
Err(named_sem::Error::WaitFailed(e))
if e.kind() == std::io::ErrorKind::Interrupted =>
{
continue
}
Err(e) => {
error!("Semaphore '{}' error: {:?}", sem_chunk_done_name, e);
signal_runner_failure();
break Ok(preloaded.output_shmem.map_header().exit_code);
}
}
};
mem_planner.set_completed();
#[cfg(gpu)]
if gpu_count_and_plan_opt.is_some() {
proofman_starks_lib_c::stream_commit_pause_c();
}
#[cfg(gpu)]
timer_start_info!(GPU_MOPS_TIME);
#[cfg(gpu)]
let gpu_metas_view: Option<(*const c_void, u32)> = gpu_count_and_plan_opt
.as_ref()
.and_then(|gpu_count_and_plan| gpu_count_and_plan.run())
.map(|metas| (metas.as_ptr() as *const c_void, metas.len() as u32));
#[cfg(gpu)]
timer_stop_and_log_info!(GPU_MOPS_TIME);
#[cfg(gpu)]
let gpu_mops_used_bytes: Option<u64> =
gpu_count_and_plan_opt.as_ref().map(|gcp| gcp.max_used_bytes() as u64);
#[cfg(not(gpu))]
let gpu_mops_used_bytes: Option<u64> = None;
mem_planner.wait();
let joined = handle.join();
crate::drain_chunk_done(&mut sem_chunk_done);
#[cfg(gpu)]
let inject_ok = match gpu_metas_view {
Some((metas_ptr, n)) => unsafe {
mem_planner.inject_gpu_metas_from_pointers(metas_ptr, n)
},
None => gpu_count_and_plan_opt.is_none(),
};
let result: Result<Vec<Plan>> = (|| -> Result<Vec<Plan>> {
let exit_code = loop_result?;
if exit_code != 0 {
return Err(AsmRunError::ExitCode(exit_code as u32))
.context("Child process returned error");
}
let response =
joined.map_err(|_| AsmRunError::JoinPanic)?.map_err(AsmRunError::ServiceError)?;
if response.result != 0 {
return Err(anyhow::anyhow!(
"ASM MO service returned non-zero result: {}",
response.result
));
}
if response.trace_len == 0 {
return Err(anyhow::anyhow!("ASM MO service returned empty trace"));
}
if response.trace_len > response.allocated_len {
return Err(anyhow::anyhow!(
"ASM MO service trace_len ({}) exceeds allocated_len ({})",
response.trace_len,
response.allocated_len
));
}
#[cfg(gpu)]
if !inject_ok {
return Err(anyhow::anyhow!(
"[gpu] GPU count-and-plan produced no metas or they were rejected; \
segment table is unpopulated, aborting MO run"
));
}
#[cfg(gpu)]
let mut mem_align_plans = gpu_count_and_plan_opt
.as_ref()
.map(|gpu_count_and_plan| gpu_count_and_plan.build_align_plans())
.unwrap_or_else(|| mem_planner.wait_mem_align_plans());
#[cfg(not(gpu))]
let mut mem_align_plans = mem_planner.wait_mem_align_plans();
stats_end!(_stats, &_process_scope);
stats_begin!(_stats, &_runner_scope, _collect_scope, "MO_COLLECT_PLANS", 0);
let plans = mem_planner.collect_plans(&mut mem_align_plans);
stats_end!(_stats, &_collect_scope);
Ok(plans)
})();
preloaded.handle_mo = Some(std::thread::spawn(move || {
drop(mem_planner);
MemPlanner::new()
}));
#[cfg(gpu)]
if let Some(gpu_count_and_plan) = gpu_count_and_plan_opt {
preloaded.gpu_count_and_plan = Some(gpu_count_and_plan);
}
let plans = result?;
#[cfg(feature = "save_mem_plans")]
save_plans(&plans, "mem_plans_cpp.txt");
stats_end!(_stats, &_runner_scope);
Ok(AsmRunnerMO { plans, gpu_mops_used_bytes })
}
}