2856 lines
104 KiB
Rust

#![allow(clippy::collapsible_if)]
use crate::memory_banks::{GlyphBankPoolInstaller, SceneBankPoolInstaller, SoundBankPoolInstaller};
use prometeu_hal::AssetBridge;
use prometeu_hal::asset::{
AssetBacklogInfo, AssetBacklogPosition, AssetCodec, AssetEntry, AssetId, AssetLoadError,
AssetOpStatus, AssetTargetStatus, BankTelemetry, BankType, HandleId, LoadStatus, PreloadEntry,
SCENE_DECODED_LAYER_OVERHEAD_BYTES_V1, SCENE_HEADER_BYTES_V1, SCENE_LAYER_COUNT_V1,
SCENE_LAYER_HEADER_BYTES_V1, SCENE_PAYLOAD_MAGIC_V1, SCENE_PAYLOAD_VERSION_V1,
SCENE_TILE_RECORD_BYTES_V1, SlotRef, SlotStats,
};
use prometeu_hal::cartridge::AssetsPayloadSource;
use prometeu_hal::color::Color;
use prometeu_hal::glyph::Glyph;
use prometeu_hal::glyph_bank::{GlyphBank, TileSize};
use prometeu_hal::sample::Sample;
use prometeu_hal::scene_bank::SceneBank;
use prometeu_hal::scene_layer::{ParallaxFactor, SceneLayer};
use prometeu_hal::sound_bank::SoundBank;
use prometeu_hal::tile::Tile;
use prometeu_hal::tilemap::TileMap;
use std::collections::{HashMap, VecDeque};
use std::io::Read;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Condvar, Mutex, RwLock};
use std::thread;
use std::time::Instant;
type ResidentMap<T> = HashMap<AssetId, ResidentEntry<T>>;
type StagedValue<T> = (Arc<T>, usize);
type StagingMap<T> = HashMap<HandleId, StagedValue<T>>;
type AssetTable = HashMap<AssetId, AssetEntry>;
type HandleTable = HashMap<HandleId, LoadHandleInfo>;
type TargetHandleTable = HashMap<SlotRef, HandleId>;
type TargetGenerationTable = HashMap<SlotRef, u64>;
const ASSET_PROGRESS_IDLE: u16 = 0;
const ASSET_PROGRESS_LOADING: u16 = 1_000;
const ASSET_PROGRESS_DONE: u16 = 10_000;
const ASSET_TELEMETRY_SAMPLE_WINDOW: usize = 32;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct AssetLatencyPercentiles {
pub p50_micros: u64,
pub p95_micros: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct AssetPipelineTelemetry {
pub submitted_requests: u64,
pub started_jobs: u64,
pub completed_jobs: u64,
pub failed_jobs: u64,
pub canceled_requests: u64,
pub superseded_requests: u64,
pub stale_results_discarded: u64,
pub current_backlog_depth: usize,
pub max_backlog_depth: usize,
pub active_progress: u16,
pub glyph_latency: AssetLatencyPercentiles,
pub sound_latency: AssetLatencyPercentiles,
pub scene_latency: AssetLatencyPercentiles,
pub last_asset_id: Option<AssetId>,
pub last_asset_latency: AssetLatencyPercentiles,
}
#[derive(Default)]
struct AssetPipelineTelemetryState {
snapshot: AssetPipelineTelemetry,
bank_samples: HashMap<BankType, VecDeque<u64>>,
asset_samples: HashMap<AssetId, VecDeque<u64>>,
}
impl AssetPipelineTelemetryState {
fn record_submitted(&mut self, backlog_depth: usize) {
self.snapshot.submitted_requests += 1;
self.snapshot.current_backlog_depth = backlog_depth;
self.snapshot.max_backlog_depth = self.snapshot.max_backlog_depth.max(backlog_depth);
}
fn record_started(&mut self, backlog_depth: usize) {
self.snapshot.started_jobs += 1;
self.snapshot.current_backlog_depth = backlog_depth;
self.snapshot.active_progress = ASSET_PROGRESS_LOADING;
}
fn record_completed(&mut self, asset_id: AssetId, bank_type: BankType, duration_micros: u64) {
self.snapshot.completed_jobs += 1;
self.snapshot.active_progress = ASSET_PROGRESS_IDLE;
self.record_latency_sample(asset_id, bank_type, duration_micros);
}
fn record_failed(&mut self, asset_id: AssetId, bank_type: BankType, duration_micros: u64) {
self.snapshot.failed_jobs += 1;
self.snapshot.active_progress = ASSET_PROGRESS_IDLE;
self.record_latency_sample(asset_id, bank_type, duration_micros);
}
fn record_canceled(&mut self) {
self.snapshot.canceled_requests += 1;
}
fn record_superseded(&mut self, count: usize) {
self.snapshot.superseded_requests += count as u64;
}
fn record_stale_discard(&mut self) {
self.snapshot.stale_results_discarded += 1;
self.snapshot.superseded_requests += 1;
self.snapshot.active_progress = ASSET_PROGRESS_IDLE;
}
fn record_latency_sample(
&mut self,
asset_id: AssetId,
bank_type: BankType,
duration_micros: u64,
) {
push_sample(
self.bank_samples.entry(bank_type).or_default(),
duration_micros,
ASSET_TELEMETRY_SAMPLE_WINDOW,
);
push_sample(
self.asset_samples.entry(asset_id).or_default(),
duration_micros,
ASSET_TELEMETRY_SAMPLE_WINDOW,
);
self.snapshot.glyph_latency = percentile_snapshot(self.bank_samples.get(&BankType::GLYPH));
self.snapshot.sound_latency = percentile_snapshot(self.bank_samples.get(&BankType::SOUNDS));
self.snapshot.scene_latency = percentile_snapshot(self.bank_samples.get(&BankType::SCENE));
self.snapshot.last_asset_id = Some(asset_id);
self.snapshot.last_asset_latency = percentile_snapshot(self.asset_samples.get(&asset_id));
}
}
fn push_sample(samples: &mut VecDeque<u64>, value: u64, window: usize) {
samples.push_back(value);
while samples.len() > window {
samples.pop_front();
}
}
fn percentile_snapshot(samples: Option<&VecDeque<u64>>) -> AssetLatencyPercentiles {
let Some(samples) = samples else {
return AssetLatencyPercentiles::default();
};
if samples.is_empty() {
return AssetLatencyPercentiles::default();
}
let mut sorted = samples.iter().copied().collect::<Vec<_>>();
sorted.sort_unstable();
AssetLatencyPercentiles {
p50_micros: percentile_value(&sorted, 50),
p95_micros: percentile_value(&sorted, 95),
}
}
fn percentile_value(sorted: &[u64], percentile: usize) -> u64 {
let index = ((sorted.len() - 1) * percentile).div_ceil(100);
sorted[index]
}
#[derive(Clone, Default)]
pub struct GlyphAssetSlotIndex {
slots_by_asset_id: Arc<RwLock<HashMap<AssetId, usize>>>,
}
impl GlyphAssetSlotIndex {
pub fn new() -> Self {
Self { slots_by_asset_id: Arc::new(RwLock::new(HashMap::new())) }
}
pub fn glyph_slot_for_asset(&self, asset_id: AssetId) -> Option<usize> {
self.slots_by_asset_id.read().unwrap().get(&asset_id).copied()
}
pub fn rebuild_from_slots(&self, slots: &[Option<AssetId>; 16]) {
let mut map = self.slots_by_asset_id.write().unwrap();
map.clear();
for (slot, asset_id) in slots.iter().enumerate() {
if let Some(asset_id) = asset_id {
map.insert(*asset_id, slot);
}
}
}
fn clear(&self) {
self.slots_by_asset_id.write().unwrap().clear();
}
}
const GLYPH_BANK_PALETTE_COUNT_V1: usize = 64;
const GLYPH_BANK_COLORS_PER_PALETTE: usize = 16;
const GLYPH_BANK_PALETTE_BYTES_V1: usize =
GLYPH_BANK_PALETTE_COUNT_V1 * GLYPH_BANK_COLORS_PER_PALETTE * size_of::<u32>();
/// Resident metadata for a decoded/materialized asset inside a BankPolicy.
#[derive(Debug)]
pub struct ResidentEntry<T> {
/// The resident, materialized object.
pub value: Arc<T>,
/// Resident size in bytes (post-decode). Used for telemetry/budgets.
pub bytes: usize,
// /// Pin count (optional): if > 0, entry should not be evicted by policy.
// pub pins: u32,
/// Telemetry / profiling fields (optional but useful).
pub loads: u64,
pub last_used: Instant,
}
impl<T> ResidentEntry<T> {
pub fn new(value: Arc<T>, bytes: usize) -> Self {
Self {
value,
bytes,
// pins: 0,
loads: 1,
last_used: Instant::now(),
}
}
}
/// Encapsulates the residency and staging policy for a specific type of asset.
/// This is internal to the AssetManager and not visible to peripherals.
pub struct BankPolicy<T> {
/// Dedup table: asset_id -> resident entry (value + telemetry).
pub resident: Arc<RwLock<ResidentMap<T>>>,
/// Staging area: handle -> value ready to commit.
pub staging: Arc<RwLock<StagingMap<T>>>,
/// Total bytes currently in resident storage.
pub used_bytes: Arc<AtomicUsize>,
/// Bytes in staging awaiting commit.
pub inflight_bytes: Arc<AtomicUsize>,
}
impl<T> Clone for BankPolicy<T> {
fn clone(&self) -> Self {
Self {
resident: Arc::clone(&self.resident),
staging: Arc::clone(&self.staging),
used_bytes: Arc::clone(&self.used_bytes),
inflight_bytes: Arc::clone(&self.inflight_bytes),
}
}
}
impl<T> BankPolicy<T> {
pub fn new() -> Self {
Self {
resident: Arc::new(RwLock::new(HashMap::new())),
staging: Arc::new(RwLock::new(HashMap::new())),
used_bytes: Arc::new(AtomicUsize::new(0)),
inflight_bytes: Arc::new(AtomicUsize::new(0)),
}
}
/// Try get a resident value by asset_id (dedupe path).
pub fn get_resident(&self, asset_id: AssetId) -> Option<Arc<T>> {
let mut map = self.resident.write().unwrap();
let entry = map.get_mut(&asset_id)?;
entry.last_used = Instant::now();
Some(Arc::clone(&entry.value))
}
/// Insert or reuse a resident entry. Returns the resident Arc<T>.
pub fn put_resident(&self, asset_id: AssetId, value: Arc<T>, bytes: usize) -> Arc<T> {
let mut map = self.resident.write().unwrap();
match map.get_mut(&asset_id) {
Some(existing) => {
existing.last_used = Instant::now();
existing.loads += 1;
Arc::clone(&existing.value)
}
None => {
let entry = ResidentEntry::new(Arc::clone(&value), bytes);
map.insert(asset_id, entry);
self.used_bytes.fetch_add(bytes, Ordering::Relaxed);
value
}
}
}
/// Place a value into staging for a given handle.
pub fn stage(&self, handle: HandleId, value: Arc<T>, bytes: usize) {
self.staging.write().unwrap().insert(handle, (value, bytes));
self.inflight_bytes.fetch_add(bytes, Ordering::Relaxed);
}
/// Take staged value (used by commit path).
pub fn take_staging(&self, handle: HandleId) -> Option<StagedValue<T>> {
let entry = self.staging.write().unwrap().remove(&handle);
if let Some((_, bytes)) = entry.as_ref() {
self.inflight_bytes.fetch_sub(*bytes, Ordering::Relaxed);
}
entry
}
pub fn clear(&self) {
self.resident.write().unwrap().clear();
self.staging.write().unwrap().clear();
self.used_bytes.store(0, Ordering::Relaxed);
self.inflight_bytes.store(0, Ordering::Relaxed);
}
}
pub struct AssetManager {
assets: Arc<RwLock<AssetTable>>,
handles: Arc<RwLock<HandleTable>>,
target_handles: Arc<RwLock<TargetHandleTable>>,
target_generations: Arc<RwLock<TargetGenerationTable>>,
next_handle_id: Mutex<HandleId>,
assets_data: Arc<RwLock<AssetsPayloadSource>>,
pipeline_telemetry: Arc<Mutex<AssetPipelineTelemetryState>>,
/// Narrow hardware interfaces
gfx_installer: Arc<dyn GlyphBankPoolInstaller>,
sound_installer: Arc<dyn SoundBankPoolInstaller>,
scene_installer: Arc<dyn SceneBankPoolInstaller>,
/// Track what is installed in each hardware slot (for stats/info).
gfx_slots: Arc<RwLock<[Option<AssetId>; 16]>>,
glyph_slot_index: GlyphAssetSlotIndex,
sound_slots: Arc<RwLock<[Option<AssetId>; 16]>>,
scene_slots: Arc<RwLock<[Option<AssetId>; 16]>>,
/// Residency policy for GFX glyph banks.
gfx_policy: BankPolicy<GlyphBank>,
/// Residency policy for sound banks.
sound_policy: BankPolicy<SoundBank>,
/// Residency policy for scene banks.
scene_policy: BankPolicy<SceneBank>,
load_worker: AssetLoadWorker,
// Commits that are ready to be applied at the next frame boundary.
pending_commits: Mutex<Vec<(HandleId, u64)>>,
}
struct LoadHandleInfo {
_asset_id: AssetId,
slot: SlotRef,
status: LoadStatus,
request_generation: u64,
progress: u16,
}
#[derive(Clone)]
struct AssetLoadJob {
handle_id: HandleId,
asset_id: AssetId,
slot: SlotRef,
request_generation: u64,
entry: AssetEntry,
}
struct AssetLoadJobResources<'a> {
handles: &'a Arc<RwLock<HandleTable>>,
target_generations: &'a Arc<RwLock<TargetGenerationTable>>,
assets_data: &'a Arc<RwLock<AssetsPayloadSource>>,
pipeline_telemetry: &'a Arc<Mutex<AssetPipelineTelemetryState>>,
gfx_policy: &'a BankPolicy<GlyphBank>,
sound_policy: &'a BankPolicy<SoundBank>,
scene_policy: &'a BankPolicy<SceneBank>,
}
#[derive(Default)]
struct AssetLoadQueueState {
pending: VecDeque<AssetLoadJob>,
active: Option<AssetLoadJob>,
shutdown: bool,
}
#[derive(Default)]
struct AssetLoadQueue {
state: Mutex<AssetLoadQueueState>,
ready: Condvar,
}
impl AssetLoadQueue {
fn submit(&self, job: AssetLoadJob) -> usize {
let mut state = self.state.lock().unwrap();
if state.shutdown {
return 0;
}
let before_len = state.pending.len();
state.pending.retain(|pending| pending.slot != job.slot);
let removed = before_len - state.pending.len();
state.pending.push_back(job);
self.ready.notify_one();
removed
}
fn next_job(&self) -> Option<AssetLoadJob> {
let mut state = self.state.lock().unwrap();
loop {
if let Some(job) = state.pending.pop_front() {
state.active = Some(job.clone());
return Some(job);
}
if state.shutdown {
return None;
}
state = self.ready.wait(state).unwrap();
}
}
fn clear_pending(&self) {
self.state.lock().unwrap().pending.clear();
}
fn complete_active(&self, handle_id: HandleId, request_generation: u64) {
let mut state = self.state.lock().unwrap();
if state.active.as_ref().is_some_and(|job| {
job.handle_id == handle_id && job.request_generation == request_generation
}) {
state.active = None;
}
}
fn pending_count(&self) -> usize {
self.state.lock().unwrap().pending.len()
}
fn active_job(&self) -> Option<AssetLoadJob> {
self.state.lock().unwrap().active.clone()
}
fn pending_position(&self, handle_id: HandleId) -> Option<usize> {
self.state
.lock()
.unwrap()
.pending
.iter()
.position(|job| job.handle_id == handle_id)
.map(|index| index + 1)
}
fn move_pending(&self, handle_id: HandleId, new_position: usize) -> bool {
if new_position == 0 {
return false;
}
let mut state = self.state.lock().unwrap();
let Some(current_index) = state.pending.iter().position(|job| job.handle_id == handle_id)
else {
return false;
};
let Some(job) = state.pending.remove(current_index) else {
return false;
};
let insert_index = new_position.saturating_sub(1).min(state.pending.len());
state.pending.insert(insert_index, job);
true
}
fn shutdown(&self) {
let mut state = self.state.lock().unwrap();
state.pending.clear();
state.shutdown = true;
self.ready.notify_all();
}
}
struct AssetLoadWorker {
queue: Arc<AssetLoadQueue>,
join_handle: Option<thread::JoinHandle<()>>,
}
impl AssetLoadWorker {
fn new(
handles: Arc<RwLock<HandleTable>>,
target_generations: Arc<RwLock<TargetGenerationTable>>,
assets_data: Arc<RwLock<AssetsPayloadSource>>,
pipeline_telemetry: Arc<Mutex<AssetPipelineTelemetryState>>,
gfx_policy: BankPolicy<GlyphBank>,
sound_policy: BankPolicy<SoundBank>,
scene_policy: BankPolicy<SceneBank>,
) -> Self {
let queue = Arc::new(AssetLoadQueue::default());
let worker_queue = Arc::clone(&queue);
let join_handle = thread::spawn(move || {
let resources = AssetLoadJobResources {
handles: &handles,
target_generations: &target_generations,
assets_data: &assets_data,
pipeline_telemetry: &pipeline_telemetry,
gfx_policy: &gfx_policy,
sound_policy: &sound_policy,
scene_policy: &scene_policy,
};
while let Some(job) = worker_queue.next_job() {
AssetManager::process_load_job(job.clone(), &resources);
worker_queue.complete_active(job.handle_id, job.request_generation);
}
});
Self { queue, join_handle: Some(join_handle) }
}
fn submit_counting_superseded(&self, job: AssetLoadJob) -> usize {
self.queue.submit(job)
}
fn clear_pending(&self) {
self.queue.clear_pending();
}
fn pending_count(&self) -> usize {
self.queue.pending_count()
}
fn active_job(&self) -> Option<AssetLoadJob> {
self.queue.active_job()
}
fn pending_position(&self, handle_id: HandleId) -> Option<usize> {
self.queue.pending_position(handle_id)
}
fn move_pending(&self, handle_id: HandleId, new_position: usize) -> bool {
self.queue.move_pending(handle_id, new_position)
}
}
impl Drop for AssetLoadWorker {
fn drop(&mut self) {
self.queue.shutdown();
if let Some(join_handle) = self.join_handle.take() {
let _ = join_handle.join();
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum AssetOpMode {
DirectFromSlice,
StageInMemory,
}
impl AssetBridge for AssetManager {
fn initialize_for_cartridge(
&self,
assets: Vec<AssetEntry>,
preload: Vec<PreloadEntry>,
assets_data: AssetsPayloadSource,
) {
self.initialize_for_cartridge(assets, preload, assets_data)
}
fn load(&self, asset_id: AssetId, slot_index: usize) -> Result<HandleId, AssetLoadError> {
self.load(asset_id, slot_index)
}
fn status(&self, handle: HandleId) -> LoadStatus {
self.status(handle)
}
fn commit(&self, handle: HandleId) -> AssetOpStatus {
self.commit(handle)
}
fn cancel(&self, handle: HandleId) -> AssetOpStatus {
self.cancel(handle)
}
fn backlog_info(&self) -> AssetBacklogInfo {
self.backlog_info()
}
fn backlog_position(&self, handle: HandleId) -> AssetBacklogPosition {
self.backlog_position(handle)
}
fn backlog_move(&self, handle: HandleId, new_position: usize) -> AssetOpStatus {
self.backlog_move(handle, new_position)
}
fn backlog_promote(&self, handle: HandleId) -> AssetOpStatus {
self.backlog_promote(handle)
}
fn backlog_demote(&self, handle: HandleId) -> AssetOpStatus {
self.backlog_demote(handle)
}
fn target_status(&self, bank_type: BankType, slot: usize) -> AssetTargetStatus {
self.target_status(bank_type, slot)
}
fn apply_commits(&self) {
self.apply_commits()
}
fn bank_telemetry(&self) -> Vec<BankTelemetry> {
self.bank_telemetry()
}
fn slot_info(&self, slot: SlotRef) -> SlotStats {
self.slot_info(slot)
}
fn shutdown(&self) {
self.shutdown()
}
}
impl AssetManager {
fn decode_glyph_bank_layout(
entry: &AssetEntry,
) -> Result<(TileSize, usize, usize, usize), String> {
let meta = entry.metadata_as_glyph_bank()?;
let tile_size = match meta.tile_size {
8 => TileSize::Size8,
16 => TileSize::Size16,
32 => TileSize::Size32,
_ => return Err(format!("Invalid tile_size: {}", meta.tile_size)),
};
if meta.palette_count as usize != GLYPH_BANK_PALETTE_COUNT_V1 {
return Err(format!("Invalid palette_count: {}", meta.palette_count));
}
let width = meta.width as usize;
let height = meta.height as usize;
let logical_pixels = width.checked_mul(height).ok_or("GlyphBank dimensions overflow")?;
let serialized_pixel_bytes = logical_pixels.div_ceil(2);
let serialized_size = serialized_pixel_bytes
.checked_add(GLYPH_BANK_PALETTE_BYTES_V1)
.ok_or("GlyphBank serialized size overflow")?;
let decoded_size = logical_pixels
.checked_add(GLYPH_BANK_PALETTE_BYTES_V1)
.ok_or("GlyphBank decoded size overflow")?;
if entry.size != serialized_size as u64 {
return Err(format!(
"Invalid GLYPHBANK serialized size: expected {}, got {}",
serialized_size, entry.size
));
}
if entry.decoded_size != decoded_size as u64 {
return Err(format!(
"Invalid GLYPHBANK decoded_size: expected {}, got {}",
decoded_size, entry.decoded_size
));
}
Ok((tile_size, width, height, serialized_pixel_bytes))
}
fn unpack_glyph_bank_pixels(packed_pixels: &[u8], logical_pixels: usize) -> Vec<u8> {
let mut pixel_indices = Vec::with_capacity(logical_pixels);
for &packed in packed_pixels {
if pixel_indices.len() < logical_pixels {
pixel_indices.push(packed >> 4);
}
if pixel_indices.len() < logical_pixels {
pixel_indices.push(packed & 0x0f);
}
}
pixel_indices
}
fn op_mode_for(entry: &AssetEntry) -> Result<AssetOpMode, String> {
match (entry.bank_type, entry.codec) {
(BankType::GLYPH, AssetCodec::None) => Ok(AssetOpMode::StageInMemory),
(BankType::SOUNDS, AssetCodec::None) => Ok(AssetOpMode::DirectFromSlice),
(BankType::SCENE, AssetCodec::None) => Ok(AssetOpMode::DirectFromSlice),
}
}
pub fn new(
assets: Vec<AssetEntry>,
assets_data: AssetsPayloadSource,
gfx_installer: Arc<dyn GlyphBankPoolInstaller>,
sound_installer: Arc<dyn SoundBankPoolInstaller>,
scene_installer: Arc<dyn SceneBankPoolInstaller>,
) -> Self {
let mut asset_map = HashMap::new();
for entry in assets {
asset_map.insert(entry.asset_id, entry);
}
let handles = Arc::new(RwLock::new(HashMap::new()));
let target_generations = Arc::new(RwLock::new(HashMap::new()));
let assets_data = Arc::new(RwLock::new(assets_data));
let pipeline_telemetry = Arc::new(Mutex::new(AssetPipelineTelemetryState::default()));
let gfx_policy = BankPolicy::new();
let sound_policy = BankPolicy::new();
let scene_policy = BankPolicy::new();
let load_worker = AssetLoadWorker::new(
Arc::clone(&handles),
Arc::clone(&target_generations),
Arc::clone(&assets_data),
Arc::clone(&pipeline_telemetry),
gfx_policy.clone(),
sound_policy.clone(),
scene_policy.clone(),
);
Self {
assets: Arc::new(RwLock::new(asset_map)),
gfx_installer,
sound_installer,
scene_installer,
gfx_slots: Arc::new(RwLock::new(std::array::from_fn(|_| None))),
glyph_slot_index: GlyphAssetSlotIndex::new(),
sound_slots: Arc::new(RwLock::new(std::array::from_fn(|_| None))),
scene_slots: Arc::new(RwLock::new(std::array::from_fn(|_| None))),
gfx_policy,
sound_policy,
scene_policy,
handles,
target_handles: Arc::new(RwLock::new(HashMap::new())),
target_generations,
next_handle_id: Mutex::new(1),
assets_data,
pipeline_telemetry,
load_worker,
pending_commits: Mutex::new(Vec::new()),
}
}
pub fn initialize_for_cartridge(
&self,
assets: Vec<AssetEntry>,
preload: Vec<PreloadEntry>,
assets_data: AssetsPayloadSource,
) {
self.shutdown();
{
let mut asset_map = self.assets.write().unwrap();
asset_map.clear();
for entry in assets.iter() {
asset_map.insert(entry.asset_id, entry.clone());
}
}
*self.assets_data.write().unwrap() = assets_data;
// Perform Preload for assets in the preload list
for item in preload {
let entry_opt = {
let assets = self.assets.read().unwrap();
assets.get(&item.asset_id).cloned()
};
if let Some(entry) = entry_opt {
let slot_index = item.slot;
match entry.bank_type {
BankType::GLYPH => {
if let Ok(bank) =
Self::perform_load_glyph_bank(&entry, self.assets_data.clone())
{
let bank_arc = Arc::new(bank);
self.gfx_policy.put_resident(
entry.asset_id,
Arc::clone(&bank_arc),
entry.decoded_size as usize,
);
self.gfx_installer.install_glyph_bank(slot_index, bank_arc);
let mut slots = self.gfx_slots.write().unwrap();
if slot_index < slots.len() {
slots[slot_index] = Some(entry.asset_id);
self.glyph_slot_index.rebuild_from_slots(&slots);
}
}
}
BankType::SOUNDS => {
if let Ok(bank) =
Self::perform_load_sound_bank(&entry, self.assets_data.clone())
{
let bank_arc = Arc::new(bank);
self.sound_policy.put_resident(
entry.asset_id,
Arc::clone(&bank_arc),
entry.decoded_size as usize,
);
self.sound_installer.install_sound_bank(slot_index, bank_arc);
let mut slots = self.sound_slots.write().unwrap();
if slot_index < slots.len() {
slots[slot_index] = Some(entry.asset_id);
}
}
}
BankType::SCENE => {
if let Ok(bank) =
Self::perform_load_scene_bank(&entry, self.assets_data.clone())
{
let bank_arc = Arc::new(bank);
self.scene_policy.put_resident(
entry.asset_id,
Arc::clone(&bank_arc),
entry.decoded_size as usize,
);
self.scene_installer.install_scene_bank(slot_index, bank_arc);
let mut slots = self.scene_slots.write().unwrap();
if slot_index < slots.len() {
slots[slot_index] = Some(entry.asset_id);
}
}
}
}
}
}
}
fn handle_for_slot(&self, slot: SlotRef) -> HandleId {
if let Some(handle_id) = self.target_handles.read().unwrap().get(&slot).copied() {
return handle_id;
}
let mut target_handles = self.target_handles.write().unwrap();
if let Some(handle_id) = target_handles.get(&slot).copied() {
return handle_id;
}
let mut next_id = self.next_handle_id.lock().unwrap();
let handle_id = *next_id;
*next_id += 1;
target_handles.insert(slot, handle_id);
handle_id
}
fn bump_target_generation(&self, slot: SlotRef) -> u64 {
let mut generations = self.target_generations.write().unwrap();
let generation = generations.entry(slot).or_insert(0);
*generation += 1;
*generation
}
fn current_target_generation(
target_generations: &Arc<RwLock<TargetGenerationTable>>,
slot: SlotRef,
) -> u64 {
target_generations.read().unwrap().get(&slot).copied().unwrap_or(0)
}
fn is_current_target_generation(
target_generations: &Arc<RwLock<TargetGenerationTable>>,
slot: SlotRef,
request_generation: u64,
) -> bool {
Self::current_target_generation(target_generations, slot) == request_generation
}
fn clear_staging_for_handle(&self, handle_id: HandleId) {
self.gfx_policy.take_staging(handle_id);
self.sound_policy.take_staging(handle_id);
self.scene_policy.take_staging(handle_id);
}
fn target_for_handle(&self, handle_id: HandleId) -> Option<SlotRef> {
self.target_handles
.read()
.unwrap()
.iter()
.find_map(|(slot, target_handle)| (*target_handle == handle_id).then_some(*slot))
}
fn slot_asset_id(&self, slot: SlotRef) -> Option<AssetId> {
match slot.asset_type {
BankType::GLYPH => self.gfx_slots.read().unwrap().get(slot.index).and_then(|s| *s),
BankType::SOUNDS => self.sound_slots.read().unwrap().get(slot.index).and_then(|s| *s),
BankType::SCENE => self.scene_slots.read().unwrap().get(slot.index).and_then(|s| *s),
}
}
fn idle_target_state(&self, slot: SlotRef) -> LoadStatus {
if self.slot_asset_id(slot).is_some() { LoadStatus::COMMITTED } else { LoadStatus::EMPTY }
}
pub fn load(&self, asset_id: AssetId, slot_index: usize) -> Result<HandleId, AssetLoadError> {
if slot_index >= 16 {
return Err(AssetLoadError::SlotIndexInvalid);
}
let entry = {
let assets = self.assets.read().unwrap();
assets.get(&asset_id).ok_or(AssetLoadError::AssetNotFound)?.clone()
};
let slot = match entry.bank_type {
BankType::GLYPH => SlotRef::gfx(slot_index),
BankType::SOUNDS => SlotRef::audio(slot_index),
BankType::SCENE => SlotRef::scene(slot_index),
};
let handle_id = self.handle_for_slot(slot);
let request_generation = self.bump_target_generation(slot);
self.clear_staging_for_handle(handle_id);
// Check if already resident (Dedup)
let already_resident = match entry.bank_type {
BankType::GLYPH => {
if let Some(bank) = self.gfx_policy.get_resident(asset_id) {
self.gfx_policy.stage(handle_id, bank, entry.decoded_size as usize);
true
} else {
false
}
}
BankType::SOUNDS => {
if let Some(bank) = self.sound_policy.get_resident(asset_id) {
self.sound_policy.stage(handle_id, bank, entry.decoded_size as usize);
true
} else {
false
}
}
BankType::SCENE => {
if let Some(bank) = self.scene_policy.get_resident(asset_id) {
self.scene_policy.stage(handle_id, bank, entry.decoded_size as usize);
true
} else {
false
}
}
};
if already_resident {
self.handles.write().unwrap().insert(
handle_id,
LoadHandleInfo {
_asset_id: asset_id,
slot,
status: LoadStatus::READY,
request_generation,
progress: ASSET_PROGRESS_DONE,
},
);
self.pipeline_telemetry.lock().unwrap().record_submitted(0);
self.pipeline_telemetry.lock().unwrap().record_completed(asset_id, entry.bank_type, 0);
return Ok(handle_id);
}
// Not resident, start loading
self.handles.write().unwrap().insert(
handle_id,
LoadHandleInfo {
_asset_id: asset_id,
slot,
status: LoadStatus::PENDING,
request_generation,
progress: ASSET_PROGRESS_IDLE,
},
);
let superseded = self.load_worker.submit_counting_superseded(AssetLoadJob {
handle_id,
asset_id,
slot,
request_generation,
entry,
});
let backlog_depth = self.load_worker.pending_count();
let mut telemetry = self.pipeline_telemetry.lock().unwrap();
telemetry.record_submitted(backlog_depth);
telemetry.record_superseded(superseded);
Ok(handle_id)
}
fn process_load_job(job: AssetLoadJob, resources: &AssetLoadJobResources<'_>) {
if !Self::is_current_target_generation(
resources.target_generations,
job.slot,
job.request_generation,
) {
resources.pipeline_telemetry.lock().unwrap().record_stale_discard();
return;
}
{
let mut handles_map = resources.handles.write().unwrap();
let Some(handle) = handles_map.get_mut(&job.handle_id) else {
return;
};
if handle.request_generation != job.request_generation
|| handle.status != LoadStatus::PENDING
{
return;
}
handle.status = LoadStatus::LOADING;
handle.progress = ASSET_PROGRESS_LOADING;
}
let started_at = Instant::now();
resources.pipeline_telemetry.lock().unwrap().record_started(0);
match job.entry.bank_type {
BankType::GLYPH => {
let result =
Self::perform_load_glyph_bank(&job.entry, Arc::clone(resources.assets_data));
match result {
Ok(tilebank) => {
if !Self::is_current_target_generation(
resources.target_generations,
job.slot,
job.request_generation,
) {
resources.pipeline_telemetry.lock().unwrap().record_stale_discard();
return;
}
let bank_arc = Arc::new(tilebank);
let resident_arc = resources.gfx_policy.put_resident(
job.asset_id,
bank_arc,
job.entry.decoded_size as usize,
);
if !Self::is_current_target_generation(
resources.target_generations,
job.slot,
job.request_generation,
) {
resources.pipeline_telemetry.lock().unwrap().record_stale_discard();
return;
}
resources.gfx_policy.stage(
job.handle_id,
resident_arc,
job.entry.decoded_size as usize,
);
Self::complete_load_job(resources.handles, &job, LoadStatus::READY);
Self::record_load_job_finished(
resources.pipeline_telemetry,
&job,
started_at,
LoadStatus::READY,
);
}
Err(_) => {
Self::complete_load_job(resources.handles, &job, LoadStatus::ERROR);
Self::record_load_job_finished(
resources.pipeline_telemetry,
&job,
started_at,
LoadStatus::ERROR,
);
}
}
}
BankType::SOUNDS => {
let result =
Self::perform_load_sound_bank(&job.entry, Arc::clone(resources.assets_data));
match result {
Ok(soundbank) => {
if !Self::is_current_target_generation(
resources.target_generations,
job.slot,
job.request_generation,
) {
resources.pipeline_telemetry.lock().unwrap().record_stale_discard();
return;
}
let bank_arc = Arc::new(soundbank);
let resident_arc = resources.sound_policy.put_resident(
job.asset_id,
bank_arc,
job.entry.decoded_size as usize,
);
if !Self::is_current_target_generation(
resources.target_generations,
job.slot,
job.request_generation,
) {
resources.pipeline_telemetry.lock().unwrap().record_stale_discard();
return;
}
resources.sound_policy.stage(
job.handle_id,
resident_arc,
job.entry.decoded_size as usize,
);
Self::complete_load_job(resources.handles, &job, LoadStatus::READY);
Self::record_load_job_finished(
resources.pipeline_telemetry,
&job,
started_at,
LoadStatus::READY,
);
}
Err(_) => {
Self::complete_load_job(resources.handles, &job, LoadStatus::ERROR);
Self::record_load_job_finished(
resources.pipeline_telemetry,
&job,
started_at,
LoadStatus::ERROR,
);
}
}
}
BankType::SCENE => {
let result =
Self::perform_load_scene_bank(&job.entry, Arc::clone(resources.assets_data));
match result {
Ok(scenebank) => {
if !Self::is_current_target_generation(
resources.target_generations,
job.slot,
job.request_generation,
) {
resources.pipeline_telemetry.lock().unwrap().record_stale_discard();
return;
}
let bank_arc = Arc::new(scenebank);
let resident_arc = resources.scene_policy.put_resident(
job.asset_id,
bank_arc,
job.entry.decoded_size as usize,
);
if !Self::is_current_target_generation(
resources.target_generations,
job.slot,
job.request_generation,
) {
resources.pipeline_telemetry.lock().unwrap().record_stale_discard();
return;
}
resources.scene_policy.stage(
job.handle_id,
resident_arc,
job.entry.decoded_size as usize,
);
Self::complete_load_job(resources.handles, &job, LoadStatus::READY);
Self::record_load_job_finished(
resources.pipeline_telemetry,
&job,
started_at,
LoadStatus::READY,
);
}
Err(_) => {
Self::complete_load_job(resources.handles, &job, LoadStatus::ERROR);
Self::record_load_job_finished(
resources.pipeline_telemetry,
&job,
started_at,
LoadStatus::ERROR,
);
}
}
}
}
}
fn record_load_job_finished(
pipeline_telemetry: &Arc<Mutex<AssetPipelineTelemetryState>>,
job: &AssetLoadJob,
started_at: Instant,
status: LoadStatus,
) {
let duration_micros = started_at.elapsed().as_micros() as u64;
let mut telemetry = pipeline_telemetry.lock().unwrap();
match status {
LoadStatus::READY => {
telemetry.record_completed(job.asset_id, job.entry.bank_type, duration_micros)
}
LoadStatus::ERROR => {
telemetry.record_failed(job.asset_id, job.entry.bank_type, duration_micros)
}
_ => {}
}
}
fn complete_load_job(
handles: &Arc<RwLock<HandleTable>>,
job: &AssetLoadJob,
status: LoadStatus,
) {
let mut handles_map = handles.write().unwrap();
if let Some(handle) = handles_map.get_mut(&job.handle_id) {
if handle.request_generation == job.request_generation
&& handle.status == LoadStatus::LOADING
{
handle.status = status;
handle.progress = ASSET_PROGRESS_DONE;
}
}
}
fn perform_load_glyph_bank(
entry: &AssetEntry,
assets_data: Arc<RwLock<AssetsPayloadSource>>,
) -> Result<GlyphBank, String> {
let op_mode = Self::op_mode_for(entry)?;
let slice = {
let assets_data = assets_data.read().unwrap();
assets_data
.open_slice(entry.offset, entry.size)
.map_err(|_| "Asset offset/size out of bounds".to_string())?
};
match op_mode {
AssetOpMode::StageInMemory => {
let buffer =
slice.read_all().map_err(|_| "Asset payload read failed".to_string())?;
Self::decode_glyph_bank_from_buffer(entry, &buffer)
}
AssetOpMode::DirectFromSlice => {
let mut reader =
slice.open_reader().map_err(|_| "Asset payload read failed".to_string())?;
Self::decode_glyph_bank_from_reader(entry, &mut reader)
}
}
}
fn decode_glyph_bank_from_buffer(
entry: &AssetEntry,
buffer: &[u8],
) -> Result<GlyphBank, String> {
let (tile_size, width, height, packed_pixel_bytes) = Self::decode_glyph_bank_layout(entry)?;
if buffer.len() < packed_pixel_bytes + GLYPH_BANK_PALETTE_BYTES_V1 {
return Err("Buffer too small for GLYPHBANK".to_string());
}
let logical_pixels = width * height;
let packed_pixels = &buffer[0..packed_pixel_bytes];
let pixel_indices = Self::unpack_glyph_bank_pixels(packed_pixels, logical_pixels);
let palette_data =
&buffer[packed_pixel_bytes..packed_pixel_bytes + GLYPH_BANK_PALETTE_BYTES_V1];
let mut palettes =
vec![[Color::BLACK; GLYPH_BANK_COLORS_PER_PALETTE]; GLYPH_BANK_PALETTE_COUNT_V1];
for (p, pal) in palettes.iter_mut().enumerate() {
for (c, slot) in pal.iter_mut().enumerate() {
let offset = (p * 16 + c) * 4;
*slot = Color::rgba(
palette_data[offset],
palette_data[offset + 1],
palette_data[offset + 2],
palette_data[offset + 3],
);
}
}
Ok(GlyphBank { tile_size, width, height, pixel_indices, palettes })
}
fn decode_glyph_bank_from_reader(
entry: &AssetEntry,
reader: &mut impl Read,
) -> Result<GlyphBank, String> {
let (tile_size, width, height, packed_pixel_bytes) = Self::decode_glyph_bank_layout(entry)?;
let logical_pixels = width * height;
let mut packed_pixels = vec![0_u8; packed_pixel_bytes];
reader
.read_exact(&mut packed_pixels)
.map_err(|_| "Buffer too small for GLYPHBANK".to_string())?;
let pixel_indices = Self::unpack_glyph_bank_pixels(&packed_pixels, logical_pixels);
let mut palette_data = [0_u8; GLYPH_BANK_PALETTE_BYTES_V1];
reader
.read_exact(&mut palette_data)
.map_err(|_| "Buffer too small for GLYPHBANK".to_string())?;
let mut palettes =
vec![[Color::BLACK; GLYPH_BANK_COLORS_PER_PALETTE]; GLYPH_BANK_PALETTE_COUNT_V1];
for (p, pal) in palettes.iter_mut().enumerate() {
for (c, slot) in pal.iter_mut().enumerate() {
let offset = (p * 16 + c) * 4;
*slot = Color::rgba(
palette_data[offset],
palette_data[offset + 1],
palette_data[offset + 2],
palette_data[offset + 3],
);
}
}
Ok(GlyphBank { tile_size, width, height, pixel_indices, palettes })
}
fn perform_load_sound_bank(
entry: &AssetEntry,
assets_data: Arc<RwLock<AssetsPayloadSource>>,
) -> Result<SoundBank, String> {
let op_mode = Self::op_mode_for(entry)?;
let slice = {
let assets_data = assets_data.read().unwrap();
assets_data
.open_slice(entry.offset, entry.size)
.map_err(|_| "Asset offset/size out of bounds".to_string())?
};
match op_mode {
AssetOpMode::DirectFromSlice => {
let mut reader =
slice.open_reader().map_err(|_| "Asset payload read failed".to_string())?;
Self::decode_sound_bank_from_reader(entry, &mut reader)
}
AssetOpMode::StageInMemory => {
let buffer =
slice.read_all().map_err(|_| "Asset payload read failed".to_string())?;
Self::decode_sound_bank_from_buffer(entry, &buffer)
}
}
}
fn perform_load_scene_bank(
entry: &AssetEntry,
assets_data: Arc<RwLock<AssetsPayloadSource>>,
) -> Result<SceneBank, String> {
let _ = entry.metadata_as_scene_bank()?;
let op_mode = Self::op_mode_for(entry)?;
let slice = {
let assets_data = assets_data.read().unwrap();
assets_data
.open_slice(entry.offset, entry.size)
.map_err(|_| "Asset offset/size out of bounds".to_string())?
};
match op_mode {
AssetOpMode::DirectFromSlice => {
let mut reader =
slice.open_reader().map_err(|_| "Asset payload read failed".to_string())?;
Self::decode_scene_bank_from_reader(entry, &mut reader)
}
AssetOpMode::StageInMemory => {
let buffer =
slice.read_all().map_err(|_| "Asset payload read failed".to_string())?;
Self::decode_scene_bank_from_buffer(entry, &buffer)
}
}
}
fn decode_scene_bank_from_buffer(
entry: &AssetEntry,
buffer: &[u8],
) -> Result<SceneBank, String> {
let _ = entry.metadata_as_scene_bank()?;
if buffer.len() < SCENE_HEADER_BYTES_V1 {
return Err("Buffer too small for SCENE".to_string());
}
if buffer[0..4] != SCENE_PAYLOAD_MAGIC_V1 {
return Err("Invalid SCENE magic".to_string());
}
let version = u16::from_le_bytes([buffer[4], buffer[5]]);
if version != SCENE_PAYLOAD_VERSION_V1 {
return Err(format!("Unsupported SCENE version: {}", version));
}
let layer_count = u16::from_le_bytes([buffer[6], buffer[7]]) as usize;
if layer_count != SCENE_LAYER_COUNT_V1 {
return Err(format!("Invalid SCENE layer count: {}", layer_count));
}
let mut offset = SCENE_HEADER_BYTES_V1;
let mut decoded_size = 0_usize;
let layers = std::array::from_fn(|_| SceneLayer {
active: false,
glyph_asset_id: 0,
tile_size: TileSize::Size8,
parallax_factor: ParallaxFactor { x: 1.0, y: 1.0 },
tilemap: TileMap { width: 0, height: 0, tiles: Vec::new() },
});
let mut layers = layers;
for layer in &mut layers {
let header_end = offset
.checked_add(SCENE_LAYER_HEADER_BYTES_V1)
.ok_or("SCENE layer header offset overflow")?;
if header_end > buffer.len() {
return Err("Buffer too small for SCENE layer header".to_string());
}
let flags = buffer[offset];
let glyph_asset_id = i32::from_le_bytes([
buffer[offset + 1],
buffer[offset + 2],
buffer[offset + 3],
buffer[offset + 4],
]);
let tile_size_raw = buffer[offset + 5];
let tile_size = match tile_size_raw {
8 => TileSize::Size8,
16 => TileSize::Size16,
32 => TileSize::Size32,
other => return Err(format!("Invalid SCENE tile size: {}", other)),
};
let parallax_factor_x = f32::from_le_bytes([
buffer[offset + 8],
buffer[offset + 9],
buffer[offset + 10],
buffer[offset + 11],
]);
let parallax_factor_y = f32::from_le_bytes([
buffer[offset + 12],
buffer[offset + 13],
buffer[offset + 14],
buffer[offset + 15],
]);
if !parallax_factor_x.is_finite() || !parallax_factor_y.is_finite() {
return Err("Invalid SCENE parallax_factor".to_string());
}
let width = u32::from_le_bytes([
buffer[offset + 16],
buffer[offset + 17],
buffer[offset + 18],
buffer[offset + 19],
]) as usize;
let height = u32::from_le_bytes([
buffer[offset + 20],
buffer[offset + 21],
buffer[offset + 22],
buffer[offset + 23],
]) as usize;
let tile_count = u32::from_le_bytes([
buffer[offset + 24],
buffer[offset + 25],
buffer[offset + 26],
buffer[offset + 27],
]) as usize;
let expected_tile_count =
width.checked_mul(height).ok_or("SCENE tile count overflow")?;
if tile_count != expected_tile_count {
return Err(format!(
"Invalid SCENE tile count for layer: expected {}, got {}",
expected_tile_count, tile_count
));
}
offset = header_end;
let tile_bytes = tile_count
.checked_mul(SCENE_TILE_RECORD_BYTES_V1)
.ok_or("SCENE tile payload overflow")?;
let tiles_end = offset.checked_add(tile_bytes).ok_or("SCENE payload overflow")?;
if tiles_end > buffer.len() {
return Err("Buffer too small for SCENE tile data".to_string());
}
let mut tiles = Vec::with_capacity(tile_count);
for _ in 0..tile_count {
let tile_flags = buffer[offset];
let palette_id = buffer[offset + 1];
let glyph_id = u16::from_le_bytes([buffer[offset + 2], buffer[offset + 3]]);
tiles.push(Tile {
active: (tile_flags & 0b0000_0001) != 0,
glyph: Glyph { glyph_id, palette_id },
flip_x: (tile_flags & 0b0000_0010) != 0,
flip_y: (tile_flags & 0b0000_0100) != 0,
});
offset += SCENE_TILE_RECORD_BYTES_V1;
}
decoded_size = decoded_size
.checked_add(SCENE_DECODED_LAYER_OVERHEAD_BYTES_V1)
.and_then(|size| size.checked_add(tile_count * size_of::<Tile>()))
.ok_or("SCENE decoded_size overflow")?;
*layer = SceneLayer {
active: (flags & 0b0000_0001) != 0,
glyph_asset_id,
tile_size,
parallax_factor: ParallaxFactor { x: parallax_factor_x, y: parallax_factor_y },
tilemap: TileMap { width, height, tiles },
};
}
if offset != buffer.len() {
return Err("Trailing bytes in SCENE payload".to_string());
}
if entry.decoded_size != decoded_size as u64 {
return Err(format!(
"Invalid SCENE decoded_size: expected {}, got {}",
decoded_size, entry.decoded_size
));
}
Ok(SceneBank { layers })
}
fn decode_scene_bank_from_reader(
entry: &AssetEntry,
reader: &mut impl Read,
) -> Result<SceneBank, String> {
let mut raw = Vec::new();
reader.read_to_end(&mut raw).map_err(|_| "Asset payload read failed".to_string())?;
Self::decode_scene_bank_from_buffer(entry, &raw)
}
fn decode_sound_bank_from_buffer(
entry: &AssetEntry,
buffer: &[u8],
) -> Result<SoundBank, String> {
let meta = entry.metadata_as_sound_bank()?;
let sample_rate = meta.sample_rate;
let mut data = Vec::with_capacity(buffer.len() / 2);
for i in (0..buffer.len()).step_by(2) {
if i + 1 < buffer.len() {
data.push(i16::from_le_bytes([buffer[i], buffer[i + 1]]));
}
}
let sample = Arc::new(Sample::new(sample_rate, data));
Ok(SoundBank::new(vec![sample]))
}
fn decode_sound_bank_from_reader(
entry: &AssetEntry,
reader: &mut impl Read,
) -> Result<SoundBank, String> {
let mut raw = Vec::new();
reader.read_to_end(&mut raw).map_err(|_| "Asset payload read failed".to_string())?;
Self::decode_sound_bank_from_buffer(entry, &raw)
}
pub fn status(&self, handle: HandleId) -> LoadStatus {
if let Some(status) = self.handles.read().unwrap().get(&handle).map(|h| h.status) {
return status;
}
self.target_for_handle(handle)
.map(|slot| self.idle_target_state(slot))
.unwrap_or(LoadStatus::UnknownHandle)
}
pub fn backlog_info(&self) -> AssetBacklogInfo {
let active = self.load_worker.active_job();
let active_progress = active
.as_ref()
.and_then(|job| self.handles.read().unwrap().get(&job.handle_id).map(|h| h.progress))
.unwrap_or(ASSET_PROGRESS_IDLE);
AssetBacklogInfo {
status: AssetOpStatus::Ok,
pending_count: self.load_worker.pending_count() as u32,
active_handle: active.as_ref().map(|job| job.handle_id).unwrap_or(0),
active_asset_id: active.as_ref().map(|job| job.asset_id),
active_bank_type: active.as_ref().map(|job| job.slot.asset_type),
active_slot: active.as_ref().map(|job| job.slot.index),
active_progress,
}
}
pub fn backlog_position(&self, handle: HandleId) -> AssetBacklogPosition {
let Some(slot) = self.target_for_handle(handle) else {
return AssetBacklogPosition {
status: AssetOpStatus::UnknownHandle,
state: LoadStatus::UnknownHandle,
position: 0,
progress: 0,
};
};
if let Some(position) = self.load_worker.pending_position(handle) {
return AssetBacklogPosition {
status: AssetOpStatus::Ok,
state: LoadStatus::QUEUED,
position: position as u32,
progress: self
.handles
.read()
.unwrap()
.get(&handle)
.map(|h| h.progress)
.unwrap_or(ASSET_PROGRESS_IDLE),
};
}
if self.load_worker.active_job().as_ref().is_some_and(|job| job.handle_id == handle) {
return AssetBacklogPosition {
status: AssetOpStatus::Ok,
state: LoadStatus::ACTIVE,
position: 0,
progress: self
.handles
.read()
.unwrap()
.get(&handle)
.map(|h| h.progress)
.unwrap_or(ASSET_PROGRESS_LOADING),
};
}
AssetBacklogPosition {
status: AssetOpStatus::Ok,
state: self
.handles
.read()
.unwrap()
.get(&handle)
.map(|h| h.status)
.unwrap_or_else(|| self.idle_target_state(slot)),
position: 0,
progress: self
.handles
.read()
.unwrap()
.get(&handle)
.map(|h| h.progress)
.unwrap_or(ASSET_PROGRESS_IDLE),
}
}
pub fn backlog_move(&self, handle: HandleId, new_position: usize) -> AssetOpStatus {
if self.target_for_handle(handle).is_none() {
return AssetOpStatus::UnknownHandle;
}
if self.load_worker.move_pending(handle, new_position) {
AssetOpStatus::Ok
} else {
AssetOpStatus::InvalidState
}
}
pub fn backlog_promote(&self, handle: HandleId) -> AssetOpStatus {
self.backlog_move(handle, 1)
}
pub fn backlog_demote(&self, handle: HandleId) -> AssetOpStatus {
if self.target_for_handle(handle).is_none() {
return AssetOpStatus::UnknownHandle;
}
let new_position = self.load_worker.pending_count();
if new_position == 0 {
return AssetOpStatus::InvalidState;
}
self.backlog_move(handle, new_position)
}
pub fn target_status(&self, bank_type: BankType, slot: usize) -> AssetTargetStatus {
if slot >= 16 {
return AssetTargetStatus {
status: AssetOpStatus::InvalidState,
asset_id: None,
handle: 0,
state: LoadStatus::INVALID,
position: 0,
progress: 0,
};
}
let slot_ref = match bank_type {
BankType::GLYPH => SlotRef::gfx(slot),
BankType::SOUNDS => SlotRef::audio(slot),
BankType::SCENE => SlotRef::scene(slot),
};
let handle = self.handle_for_slot(slot_ref);
let position = self.backlog_position(handle);
let asset_id = self
.handles
.read()
.unwrap()
.get(&handle)
.map(|h| h._asset_id)
.or_else(|| self.slot_asset_id(slot_ref));
AssetTargetStatus {
status: AssetOpStatus::Ok,
asset_id,
handle,
state: position.state,
position: position.position,
progress: position.progress,
}
}
pub fn pipeline_telemetry(&self) -> AssetPipelineTelemetry {
let mut snapshot = self.pipeline_telemetry.lock().unwrap().snapshot.clone();
snapshot.current_backlog_depth = self.load_worker.pending_count();
snapshot.active_progress = self.backlog_info().active_progress;
snapshot
}
pub fn commit(&self, handle: HandleId) -> AssetOpStatus {
let mut handles_map = self.handles.write().unwrap();
let Some(h) = handles_map.get_mut(&handle) else {
return AssetOpStatus::UnknownHandle;
};
if h.status == LoadStatus::READY {
self.pending_commits.lock().unwrap().push((handle, h.request_generation));
AssetOpStatus::Ok
} else {
AssetOpStatus::InvalidState
}
}
pub fn cancel(&self, handle: HandleId) -> AssetOpStatus {
let mut final_status = AssetOpStatus::UnknownHandle;
let mut canceled_slot = None;
{
let mut handles_map = self.handles.write().unwrap();
if let Some(h) = handles_map.get_mut(&handle) {
final_status = match h.status {
LoadStatus::PENDING | LoadStatus::LOADING | LoadStatus::READY => {
AssetOpStatus::Ok
}
LoadStatus::CANCELED => AssetOpStatus::Ok,
_ => AssetOpStatus::InvalidState,
};
match h.status {
LoadStatus::PENDING | LoadStatus::LOADING | LoadStatus::READY => {
h.status = LoadStatus::CANCELED;
h.progress = ASSET_PROGRESS_DONE;
canceled_slot = Some(h.slot);
}
_ => {}
}
}
}
if let Some(slot) = canceled_slot {
self.bump_target_generation(slot);
self.pipeline_telemetry.lock().unwrap().record_canceled();
}
self.clear_staging_for_handle(handle);
final_status
}
pub fn apply_commits(&self) {
let mut pending = self.pending_commits.lock().unwrap();
let mut handles = self.handles.write().unwrap();
for (handle_id, request_generation) in pending.drain(..) {
if let Some(h) = handles.get_mut(&handle_id) {
if h.status == LoadStatus::READY && h.request_generation == request_generation {
match h.slot.asset_type {
BankType::GLYPH => {
if let Some((bank, _)) = self.gfx_policy.take_staging(handle_id) {
self.gfx_installer.install_glyph_bank(h.slot.index, bank);
let mut slots = self.gfx_slots.write().unwrap();
if h.slot.index < slots.len() {
slots[h.slot.index] = Some(h._asset_id);
self.glyph_slot_index.rebuild_from_slots(&slots);
}
h.status = LoadStatus::COMMITTED;
}
}
BankType::SOUNDS => {
if let Some((bank, _)) = self.sound_policy.take_staging(handle_id) {
self.sound_installer.install_sound_bank(h.slot.index, bank);
let mut slots = self.sound_slots.write().unwrap();
if h.slot.index < slots.len() {
slots[h.slot.index] = Some(h._asset_id);
}
h.status = LoadStatus::COMMITTED;
}
}
BankType::SCENE => {
if let Some((bank, _)) = self.scene_policy.take_staging(handle_id) {
self.scene_installer.install_scene_bank(h.slot.index, bank);
let mut slots = self.scene_slots.write().unwrap();
if h.slot.index < slots.len() {
slots[h.slot.index] = Some(h._asset_id);
}
h.status = LoadStatus::COMMITTED;
}
}
}
}
}
}
}
pub fn bank_telemetry(&self) -> Vec<BankTelemetry> {
vec![
self.bank_telemetry_for(BankType::GLYPH),
self.bank_telemetry_for(BankType::SOUNDS),
self.bank_telemetry_for(BankType::SCENE),
]
}
fn bank_telemetry_for(&self, kind: BankType) -> BankTelemetry {
let used_slots = match kind {
BankType::GLYPH => {
self.gfx_slots.read().unwrap().iter().filter(|slot| slot.is_some()).count()
}
BankType::SOUNDS => {
self.sound_slots.read().unwrap().iter().filter(|slot| slot.is_some()).count()
}
BankType::SCENE => {
self.scene_slots.read().unwrap().iter().filter(|slot| slot.is_some()).count()
}
};
BankTelemetry { bank_type: kind, used_slots, total_slots: 16 }
}
pub fn slot_info(&self, slot: SlotRef) -> SlotStats {
match slot.asset_type {
BankType::GLYPH => {
let slots = self.gfx_slots.read().unwrap();
let asset_id = slots.get(slot.index).and_then(|s| *s);
let (bytes, asset_name) = if let Some(id) = &asset_id {
let bytes = self
.gfx_policy
.resident
.read()
.unwrap()
.get(id)
.map(|entry| entry.bytes)
.unwrap_or(0);
let name = self.assets.read().unwrap().get(id).map(|e| e.asset_name.clone());
(bytes, name)
} else {
(0, None)
};
SlotStats { asset_id, asset_name, generation: 0, resident_bytes: bytes }
}
BankType::SOUNDS => {
let slots = self.sound_slots.read().unwrap();
let asset_id = slots.get(slot.index).and_then(|s| *s);
let (bytes, asset_name) = if let Some(id) = &asset_id {
let bytes = self
.sound_policy
.resident
.read()
.unwrap()
.get(id)
.map(|entry| entry.bytes)
.unwrap_or(0);
let name = self.assets.read().unwrap().get(id).map(|e| e.asset_name.clone());
(bytes, name)
} else {
(0, None)
};
SlotStats { asset_id, asset_name, generation: 0, resident_bytes: bytes }
}
BankType::SCENE => {
let slots = self.scene_slots.read().unwrap();
let asset_id = slots.get(slot.index).and_then(|s| *s);
let (bytes, asset_name) = if let Some(id) = &asset_id {
let bytes = self
.scene_policy
.resident
.read()
.unwrap()
.get(id)
.map(|entry| entry.bytes)
.unwrap_or(0);
let name = self.assets.read().unwrap().get(id).map(|e| e.asset_name.clone());
(bytes, name)
} else {
(0, None)
};
SlotStats { asset_id, asset_name, generation: 0, resident_bytes: bytes }
}
}
}
pub fn shutdown(&self) {
self.gfx_policy.clear();
self.sound_policy.clear();
self.scene_policy.clear();
self.handles.write().unwrap().clear();
self.target_handles.write().unwrap().clear();
for generation in self.target_generations.write().unwrap().values_mut() {
*generation += 1;
}
self.load_worker.clear_pending();
self.pending_commits.lock().unwrap().clear();
self.gfx_slots.write().unwrap().fill(None);
self.glyph_slot_index.clear();
self.sound_slots.write().unwrap().fill(None);
self.scene_slots.write().unwrap().fill(None);
self.gfx_installer.clear_glyph_banks();
self.sound_installer.clear_sound_banks();
self.scene_installer.clear_scene_banks();
}
pub fn glyph_asset_slot_index(&self) -> GlyphAssetSlotIndex {
self.glyph_slot_index.clone()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::memory_banks::{
GlyphBankPoolAccess, MemoryBanks, SceneBankPoolAccess, SoundBankPoolAccess,
};
use prometeu_hal::asset::{
AssetCodec, SCENE_DECODED_LAYER_OVERHEAD_BYTES_V1, SCENE_LAYER_COUNT_V1,
SCENE_PAYLOAD_MAGIC_V1, SCENE_PAYLOAD_VERSION_V1,
};
use prometeu_hal::glyph::Glyph;
use prometeu_hal::scene_layer::{ParallaxFactor, SceneLayer};
use prometeu_hal::tile::Tile;
use prometeu_hal::tilemap::TileMap;
fn expected_glyph_payload_size(width: usize, height: usize) -> usize {
(width * height).div_ceil(2) + GLYPH_BANK_PALETTE_BYTES_V1
}
fn expected_glyph_decoded_size(width: usize, height: usize) -> usize {
width * height + GLYPH_BANK_PALETTE_BYTES_V1
}
fn test_glyph_asset_data() -> Vec<u8> {
let mut data = vec![0x11u8; 128];
data.extend_from_slice(&[0u8; GLYPH_BANK_PALETTE_BYTES_V1]);
data
}
fn test_glyph_asset_entry(asset_name: &str, width: usize, height: usize) -> AssetEntry {
AssetEntry {
asset_id: 0,
asset_name: asset_name.to_string(),
bank_type: BankType::GLYPH,
offset: 0,
size: expected_glyph_payload_size(width, height) as u64,
decoded_size: expected_glyph_decoded_size(width, height) as u64,
codec: AssetCodec::None,
metadata: serde_json::json!({
"tile_size": 16,
"width": width,
"height": height,
"palette_count": GLYPH_BANK_PALETTE_COUNT_V1,
"palette_authored": GLYPH_BANK_PALETTE_COUNT_V1
}),
}
}
fn test_scene() -> SceneBank {
let make_layer = |glyph_asset_id: AssetId,
active: bool,
parallax_x: f32,
parallax_y: f32,
tile_size: TileSize| SceneLayer {
active,
glyph_asset_id,
tile_size,
parallax_factor: ParallaxFactor { x: parallax_x, y: parallax_y },
tilemap: TileMap {
width: 2,
height: 2,
tiles: vec![
Tile {
active: true,
glyph: Glyph { glyph_id: 10, palette_id: 1 },
flip_x: false,
flip_y: false,
},
Tile {
active: true,
glyph: Glyph { glyph_id: 20, palette_id: 2 },
flip_x: true,
flip_y: false,
},
Tile {
active: true,
glyph: Glyph { glyph_id: 30, palette_id: 3 },
flip_x: false,
flip_y: true,
},
Tile {
active,
glyph: Glyph { glyph_id: 40, palette_id: 4 },
flip_x: true,
flip_y: true,
},
],
},
};
SceneBank {
layers: [
make_layer(100, true, 1.0, 1.0, TileSize::Size16),
make_layer(101, true, 0.5, 0.75, TileSize::Size8),
make_layer(102, true, 1.0, 0.5, TileSize::Size32),
make_layer(103, false, 0.25, 0.25, TileSize::Size16),
],
}
}
fn expected_scene_decoded_size(scene: &SceneBank) -> usize {
scene
.layers
.iter()
.map(|layer| {
SCENE_DECODED_LAYER_OVERHEAD_BYTES_V1
+ layer.tilemap.tiles.len() * size_of::<Tile>()
})
.sum()
}
fn encode_scene_payload(scene: &SceneBank) -> Vec<u8> {
let mut data = Vec::new();
data.extend_from_slice(&SCENE_PAYLOAD_MAGIC_V1);
data.extend_from_slice(&SCENE_PAYLOAD_VERSION_V1.to_le_bytes());
data.extend_from_slice(&(SCENE_LAYER_COUNT_V1 as u16).to_le_bytes());
data.extend_from_slice(&0_u32.to_le_bytes());
for layer in &scene.layers {
let layer_flags = if layer.active { 0b0000_0001 } else { 0 };
data.push(layer_flags);
data.extend_from_slice(&layer.glyph_asset_id.to_le_bytes());
data.push(layer.tile_size as u8);
data.extend_from_slice(&0_u16.to_le_bytes());
data.extend_from_slice(&layer.parallax_factor.x.to_le_bytes());
data.extend_from_slice(&layer.parallax_factor.y.to_le_bytes());
data.extend_from_slice(&(layer.tilemap.width as u32).to_le_bytes());
data.extend_from_slice(&(layer.tilemap.height as u32).to_le_bytes());
data.extend_from_slice(&(layer.tilemap.tiles.len() as u32).to_le_bytes());
data.extend_from_slice(&0_u32.to_le_bytes());
for tile in &layer.tilemap.tiles {
let mut tile_flags = 0_u8;
if tile.active {
tile_flags |= 0b0000_0001;
}
if tile.flip_x {
tile_flags |= 0b0000_0010;
}
if tile.flip_y {
tile_flags |= 0b0000_0100;
}
data.push(tile_flags);
data.push(tile.glyph.palette_id);
data.extend_from_slice(&tile.glyph.glyph_id.to_le_bytes());
}
}
data
}
fn test_scene_asset_entry(asset_name: &str, data: &[u8], scene: &SceneBank) -> AssetEntry {
AssetEntry {
asset_id: 2,
asset_name: asset_name.to_string(),
bank_type: BankType::SCENE,
offset: 0,
size: data.len() as u64,
decoded_size: expected_scene_decoded_size(scene) as u64,
codec: AssetCodec::None,
metadata: serde_json::json!({}),
}
}
#[test]
fn test_decode_glyph_bank_unpacks_packed_pixels_and_reads_palette_colors() {
let entry = test_glyph_asset_entry("glyphs", 2, 2);
let mut data = vec![0x10, 0x23];
data.extend_from_slice(&[0u8; GLYPH_BANK_PALETTE_BYTES_V1]);
data[2..6].copy_from_slice(&[0x12, 0x34, 0x56, 0x78]);
let bank =
AssetManager::decode_glyph_bank_from_buffer(&entry, &data).expect("glyph decode");
assert_eq!(bank.pixel_indices, vec![1, 0, 2, 3]);
assert_eq!(bank.palettes[0][0], Color::from_raw(0x12345678));
}
#[test]
fn test_decode_glyph_bank_rejects_short_packed_buffer() {
let entry = test_glyph_asset_entry("glyphs", 16, 16);
let data = vec![0u8; expected_glyph_payload_size(16, 16) - 1];
let err = match AssetManager::decode_glyph_bank_from_buffer(&entry, &data) {
Ok(_) => panic!("glyph decode should reject short buffer"),
Err(err) => err,
};
assert_eq!(err, "Buffer too small for GLYPHBANK");
}
#[test]
fn test_decode_glyph_bank_requires_palette_count_64() {
let mut entry = test_glyph_asset_entry("glyphs", 16, 16);
entry.metadata["palette_count"] = serde_json::json!(32);
let err =
match AssetManager::decode_glyph_bank_from_buffer(&entry, &test_glyph_asset_data()) {
Ok(_) => panic!("glyph decode should reject invalid palette_count"),
Err(err) => err,
};
assert_eq!(err, "Invalid palette_count: 32");
}
#[test]
fn test_op_mode_for_glyphs_none_stages_in_memory() {
let entry = test_glyph_asset_entry("glyphs", 16, 16);
assert_eq!(AssetManager::op_mode_for(&entry), Ok(AssetOpMode::StageInMemory));
}
#[test]
fn test_op_mode_for_glyphs_none_uses_typed_codec() {
let entry = test_glyph_asset_entry("glyphs", 16, 16);
assert_eq!(AssetManager::op_mode_for(&entry), Ok(AssetOpMode::StageInMemory));
}
#[test]
fn test_op_mode_for_sounds_none_reads_direct_from_slice() {
let entry = AssetEntry {
asset_id: 1,
asset_name: "sound".to_string(),
bank_type: BankType::SOUNDS,
offset: 0,
size: 8,
decoded_size: 8,
codec: AssetCodec::None,
metadata: serde_json::json!({
"sample_rate": 44100,
"channels": 1
}),
};
assert_eq!(AssetManager::op_mode_for(&entry), Ok(AssetOpMode::DirectFromSlice));
}
#[test]
fn test_op_mode_for_scene_none_reads_direct_from_slice() {
let scene = test_scene();
let data = encode_scene_payload(&scene);
let entry = test_scene_asset_entry("scene", &data, &scene);
assert_eq!(AssetManager::op_mode_for(&entry), Ok(AssetOpMode::DirectFromSlice));
}
#[test]
fn test_decode_scene_bank_from_binary_payload() {
let scene = test_scene();
let data = encode_scene_payload(&scene);
let entry = test_scene_asset_entry("scene", &data, &scene);
let decoded = AssetManager::decode_scene_bank_from_buffer(&entry, &data).expect("scene");
assert_eq!(decoded.layers[1].glyph_asset_id, 101);
assert_eq!(decoded.layers[1].parallax_factor.x, 0.5);
assert_eq!(decoded.layers[2].tile_size, TileSize::Size32);
assert!(decoded.layers[0].tilemap.tiles[1].flip_x);
assert!(decoded.layers[2].tilemap.tiles[2].flip_y);
assert!(!decoded.layers[3].active);
}
#[test]
fn test_decode_scene_bank_rejects_invalid_version() {
let scene = test_scene();
let mut data = encode_scene_payload(&scene);
data[4..6].copy_from_slice(&2_u16.to_le_bytes());
let entry = test_scene_asset_entry("scene", &data, &scene);
let err = AssetManager::decode_scene_bank_from_buffer(&entry, &data).unwrap_err();
assert_eq!(err, "Unsupported SCENE version: 2");
}
#[test]
fn test_decode_scene_bank_rejects_invalid_tile_size() {
let scene = test_scene();
let mut data = encode_scene_payload(&scene);
data[17] = 12;
let entry = test_scene_asset_entry("scene", &data, &scene);
let err = AssetManager::decode_scene_bank_from_buffer(&entry, &data).unwrap_err();
assert_eq!(err, "Invalid SCENE tile size: 12");
}
#[test]
fn test_decode_scene_bank_rejects_layer_count_mismatch() {
let scene = test_scene();
let mut data = encode_scene_payload(&scene);
data[6..8].copy_from_slice(&3_u16.to_le_bytes());
let entry = test_scene_asset_entry("scene", &data, &scene);
let err = AssetManager::decode_scene_bank_from_buffer(&entry, &data).unwrap_err();
assert_eq!(err, "Invalid SCENE layer count: 3");
}
#[test]
fn test_decode_scene_bank_rejects_tile_count_mismatch() {
let scene = test_scene();
let mut data = encode_scene_payload(&scene);
data[36..40].copy_from_slice(&5_u32.to_le_bytes());
let entry = test_scene_asset_entry("scene", &data, &scene);
let err = AssetManager::decode_scene_bank_from_buffer(&entry, &data).unwrap_err();
assert_eq!(err, "Invalid SCENE tile count for layer: expected 4, got 5");
}
#[test]
fn test_asset_loading_flow() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = test_glyph_asset_data();
let asset_entry = test_glyph_asset_entry("test_glyphs", 16, 16);
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let handle = am.load(0, 0).expect("Should start loading");
let mut status = am.status(handle);
let start = Instant::now();
while status != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
status = am.status(handle);
}
assert_eq!(status, LoadStatus::READY);
{
let staging = am.gfx_policy.staging.read().unwrap();
assert!(staging.contains_key(&handle));
}
am.commit(handle);
am.apply_commits();
assert_eq!(am.status(handle), LoadStatus::COMMITTED);
assert!(banks.glyph_bank_slot(0).is_some());
}
#[test]
fn test_asset_dedup() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = test_glyph_asset_data();
let asset_entry = test_glyph_asset_entry("test_glyphs", 16, 16);
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let handle1 = am.load(0, 0).unwrap();
let start = Instant::now();
while am.status(handle1) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
let handle2 = am.load(0, 1).unwrap();
assert_eq!(am.status(handle2), LoadStatus::READY);
let staging = am.gfx_policy.staging.read().unwrap();
let bank1 = &staging.get(&handle1).unwrap().0;
let bank2 = &staging.get(&handle2).unwrap().0;
assert!(Arc::ptr_eq(bank1, bank2));
}
#[test]
fn test_asset_load_uses_stable_handle_for_target_slot() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = test_glyph_asset_data();
let asset_entry = test_glyph_asset_entry("test_glyphs", 16, 16);
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let first_handle = am.load(0, 0).unwrap();
let second_handle = am.load(0, 0).unwrap();
assert_eq!(first_handle, second_handle);
let start = Instant::now();
while am.status(second_handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
assert_eq!(am.status(second_handle), LoadStatus::READY);
}
#[test]
fn test_asset_load_supersedes_previous_request_for_target_slot() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let first_payload = test_glyph_asset_data();
let second_payload = test_glyph_asset_data();
let second_offset = first_payload.len();
let mut data = first_payload;
data.extend_from_slice(&second_payload);
let first_asset = test_glyph_asset_entry("first_glyphs", 16, 16);
let mut second_asset = test_glyph_asset_entry("second_glyphs", 16, 16);
second_asset.asset_id = 1;
second_asset.offset = second_offset as u64;
let am = AssetManager::new(
vec![first_asset, second_asset],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let first_handle = am.load(0, 0).unwrap();
let second_handle = am.load(1, 0).unwrap();
assert_eq!(first_handle, second_handle);
let start = Instant::now();
while am.status(second_handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
assert_eq!(am.status(second_handle), LoadStatus::READY);
assert_eq!(am.commit(second_handle), AssetOpStatus::Ok);
am.apply_commits();
let slot = am.slot_info(SlotRef::gfx(0));
assert_eq!(slot.asset_id, Some(1));
}
#[test]
fn test_asset_backlog_queue_reorders_pending_requests() {
let queue = AssetLoadQueue::default();
let entry = test_glyph_asset_entry("queued_glyphs", 16, 16);
let job = |handle_id: HandleId, slot: usize| AssetLoadJob {
handle_id,
asset_id: handle_id as AssetId,
slot: SlotRef::gfx(slot),
request_generation: handle_id as u64,
entry: entry.clone(),
};
queue.submit(job(1, 0));
queue.submit(job(2, 1));
queue.submit(job(3, 2));
assert_eq!(queue.pending_position(3), Some(3));
assert!(queue.move_pending(3, 1));
assert_eq!(queue.pending_position(3), Some(1));
assert_eq!(queue.pending_position(1), Some(2));
assert!(queue.move_pending(3, queue.pending_count()));
assert_eq!(queue.pending_position(3), Some(3));
}
#[test]
fn test_asset_backlog_queue_keeps_only_latest_request_per_target() {
let queue = AssetLoadQueue::default();
let entry = test_glyph_asset_entry("queued_glyphs", 16, 16);
let job = |handle_id: HandleId, request_generation: u64| AssetLoadJob {
handle_id,
asset_id: handle_id as AssetId,
slot: SlotRef::gfx(0),
request_generation,
entry: entry.clone(),
};
assert_eq!(queue.submit(job(1, 1)), 0);
assert_eq!(queue.submit(job(2, 2)), 1);
assert_eq!(queue.pending_count(), 1);
assert_eq!(queue.pending_position(1), None);
assert_eq!(queue.pending_position(2), Some(1));
}
#[test]
fn test_asset_target_status_exposes_empty_and_ready_slot_handle() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = test_glyph_asset_data();
let asset_entry = test_glyph_asset_entry("test_glyphs", 16, 16);
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let empty = am.target_status(BankType::GLYPH, 0);
assert_eq!(empty.status, AssetOpStatus::Ok);
assert_eq!(empty.asset_id, None);
assert_eq!(empty.state, LoadStatus::EMPTY);
assert_ne!(empty.handle, 0);
assert_eq!(am.status(empty.handle), LoadStatus::EMPTY);
let handle = am.load(0, 0).unwrap();
assert_eq!(handle, empty.handle);
let start = Instant::now();
while am.status(handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
let ready = am.target_status(BankType::GLYPH, 0);
assert_eq!(ready.status, AssetOpStatus::Ok);
assert_eq!(ready.asset_id, Some(0));
assert_eq!(ready.handle, handle);
assert_eq!(ready.state, LoadStatus::READY);
}
#[test]
fn test_asset_pipeline_progress_and_telemetry_snapshot_for_success() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = test_glyph_asset_data();
let asset_entry = test_glyph_asset_entry("test_glyphs", 16, 16);
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let empty = am.target_status(BankType::GLYPH, 0);
assert_eq!(empty.progress, ASSET_PROGRESS_IDLE);
let handle = am.load(0, 0).unwrap();
let start = Instant::now();
while am.status(handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
let position = am.backlog_position(handle);
assert_eq!(position.state, LoadStatus::READY);
assert_eq!(position.progress, ASSET_PROGRESS_DONE);
let telemetry = am.pipeline_telemetry();
assert_eq!(telemetry.submitted_requests, 1);
assert_eq!(telemetry.started_jobs, 1);
assert_eq!(telemetry.completed_jobs, 1);
assert_eq!(telemetry.failed_jobs, 0);
assert_eq!(telemetry.active_progress, ASSET_PROGRESS_IDLE);
assert_eq!(telemetry.last_asset_id, Some(0));
}
#[test]
fn test_asset_pipeline_progress_and_telemetry_snapshot_for_failure() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let asset_entry = test_glyph_asset_entry("broken_glyphs", 16, 16);
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(Vec::new()),
gfx_installer,
sound_installer,
scene_installer,
);
let handle = am.load(0, 0).unwrap();
let start = Instant::now();
while am.status(handle) != LoadStatus::ERROR && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
let position = am.backlog_position(handle);
assert_eq!(position.state, LoadStatus::ERROR);
assert_eq!(position.progress, ASSET_PROGRESS_DONE);
let telemetry = am.pipeline_telemetry();
assert_eq!(telemetry.submitted_requests, 1);
assert_eq!(telemetry.started_jobs, 1);
assert_eq!(telemetry.completed_jobs, 0);
assert_eq!(telemetry.failed_jobs, 1);
assert_eq!(telemetry.last_asset_id, Some(0));
}
#[test]
fn test_asset_backlog_reorder_rejects_non_queued_request() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = test_glyph_asset_data();
let asset_entry = test_glyph_asset_entry("test_glyphs", 16, 16);
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let handle = am.load(0, 0).unwrap();
let start = Instant::now();
while am.status(handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
assert_eq!(am.backlog_promote(handle), AssetOpStatus::InvalidState);
assert_eq!(am.backlog_demote(handle), AssetOpStatus::InvalidState);
assert_eq!(am.backlog_move(999, 1), AssetOpStatus::UnknownHandle);
}
#[test]
fn test_sound_asset_loading() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
// 100 samples of 16-bit PCM (zeros)
let data = vec![0u8; 200];
let asset_entry = AssetEntry {
asset_id: 1,
asset_name: "test_sound".to_string(),
bank_type: BankType::SOUNDS,
offset: 0,
size: data.len() as u64,
decoded_size: data.len() as u64,
codec: AssetCodec::None,
metadata: serde_json::json!({
"sample_rate": 44100,
"channels": 1
}),
};
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let handle = am.load(1, 0).expect("Should start loading");
let start = Instant::now();
while am.status(handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
assert_eq!(am.status(handle), LoadStatus::READY);
am.commit(handle);
am.apply_commits();
assert_eq!(am.status(handle), LoadStatus::COMMITTED);
assert!(banks.sound_bank_slot(0).is_some());
}
#[test]
fn test_preload_on_init() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = vec![0u8; 200];
let asset_entry = AssetEntry {
asset_id: 2,
asset_name: "preload_sound".to_string(),
bank_type: BankType::SOUNDS,
offset: 0,
size: data.len() as u64,
decoded_size: data.len() as u64,
codec: AssetCodec::None,
metadata: serde_json::json!({
"sample_rate": 44100,
"channels": 1
}),
};
let preload = vec![PreloadEntry { asset_id: 2, slot: 5 }];
let am = AssetManager::new(
vec![],
AssetsPayloadSource::empty(),
gfx_installer,
sound_installer,
scene_installer,
);
// Before init, slot 5 is empty
assert!(banks.sound_bank_slot(5).is_none());
am.initialize_for_cartridge(
vec![asset_entry],
preload,
AssetsPayloadSource::from_bytes(data),
);
// After init, slot 5 should be occupied because of preload
assert!(banks.sound_bank_slot(5).is_some());
assert_eq!(am.slot_info(SlotRef::audio(5)).asset_id, Some(2));
}
#[test]
fn test_scene_asset_loading() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let scene = test_scene();
let data = encode_scene_payload(&scene);
let asset_entry = test_scene_asset_entry("test_scene", &data, &scene);
let am = AssetManager::new(
vec![asset_entry],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let handle = am.load(2, 0).expect("Should start loading scene");
let start = Instant::now();
while am.status(handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
assert_eq!(am.status(handle), LoadStatus::READY);
am.commit(handle);
am.apply_commits();
assert_eq!(am.status(handle), LoadStatus::COMMITTED);
assert!(banks.scene_bank_slot(0).is_some());
}
#[test]
fn test_scene_preload_on_init() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let scene = test_scene();
let data = encode_scene_payload(&scene);
let asset_entry = test_scene_asset_entry("preload_scene", &data, &scene);
let preload = vec![PreloadEntry { asset_id: 2, slot: 4 }];
let am = AssetManager::new(
vec![],
AssetsPayloadSource::empty(),
gfx_installer,
sound_installer,
scene_installer,
);
assert!(banks.scene_bank_slot(4).is_none());
am.initialize_for_cartridge(
vec![asset_entry],
preload,
AssetsPayloadSource::from_bytes(data),
);
assert!(banks.scene_bank_slot(4).is_some());
assert_eq!(am.slot_info(SlotRef::scene(4)).asset_id, Some(2));
}
#[test]
fn shutdown_clears_physical_memory_bank_slots() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let glyph_data = test_glyph_asset_data();
let scene = test_scene();
let scene_data = encode_scene_payload(&scene);
let glyph_entry = test_glyph_asset_entry("preload_glyphs", 16, 16);
let scene_entry = test_scene_asset_entry("preload_scene", &scene_data, &scene);
let mut payload = Vec::new();
payload.extend_from_slice(&glyph_data);
payload.extend_from_slice(&scene_data);
let mut shifted_scene_entry = scene_entry;
shifted_scene_entry.offset = glyph_data.len() as u64;
let am = AssetManager::new(
vec![],
AssetsPayloadSource::empty(),
gfx_installer,
sound_installer,
scene_installer,
);
am.initialize_for_cartridge(
vec![glyph_entry, shifted_scene_entry],
vec![PreloadEntry { asset_id: 0, slot: 0 }, PreloadEntry { asset_id: 2, slot: 0 }],
AssetsPayloadSource::from_bytes(payload),
);
assert!(banks.glyph_bank_slot(0).is_some());
assert!(banks.scene_bank_slot(0).is_some());
assert_eq!(am.slot_info(SlotRef::gfx(0)).asset_id, Some(0));
assert_eq!(am.slot_info(SlotRef::scene(0)).asset_id, Some(2));
am.shutdown();
assert!(banks.glyph_bank_slot(0).is_none());
assert!(banks.scene_bank_slot(0).is_none());
assert_eq!(am.slot_info(SlotRef::gfx(0)).asset_id, None);
assert_eq!(am.slot_info(SlotRef::scene(0)).asset_id, None);
}
#[test]
fn test_load_returns_asset_not_found() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let am = AssetManager::new(
vec![],
AssetsPayloadSource::empty(),
gfx_installer,
sound_installer,
scene_installer,
);
let result = am.load(999, 0);
assert_eq!(result, Err(AssetLoadError::AssetNotFound));
}
#[test]
fn test_load_returns_slot_index_invalid() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = test_glyph_asset_data();
let am = AssetManager::new(
vec![test_glyph_asset_entry("test_glyphs", 16, 16)],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
let result = am.load(0, 16);
assert_eq!(result, Err(AssetLoadError::SlotIndexInvalid));
}
#[test]
fn test_status_returns_unknown_handle() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let am = AssetManager::new(
vec![],
AssetsPayloadSource::empty(),
gfx_installer,
sound_installer,
scene_installer,
);
assert_eq!(am.status(999), LoadStatus::UnknownHandle);
}
#[test]
fn test_commit_and_cancel_return_explicit_statuses() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let data = test_glyph_asset_data();
let am = AssetManager::new(
vec![test_glyph_asset_entry("test_glyphs", 16, 16)],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
assert_eq!(am.commit(999), AssetOpStatus::UnknownHandle);
assert_eq!(am.cancel(999), AssetOpStatus::UnknownHandle);
let handle = am.load(0, 0).expect("load must allocate handle");
let start = Instant::now();
while am.status(handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
assert_eq!(am.cancel(handle), AssetOpStatus::Ok);
assert_eq!(am.status(handle), LoadStatus::CANCELED);
assert_eq!(am.commit(handle), AssetOpStatus::InvalidState);
}
#[test]
fn test_asset_telemetry_incremental() {
let banks = Arc::new(MemoryBanks::new());
let gfx_installer = Arc::clone(&banks) as Arc<dyn GlyphBankPoolInstaller>;
let sound_installer = Arc::clone(&banks) as Arc<dyn SoundBankPoolInstaller>;
let scene_installer = Arc::clone(&banks) as Arc<dyn SceneBankPoolInstaller>;
let width = 16;
let height = 16;
let data = test_glyph_asset_data();
let am = AssetManager::new(
vec![test_glyph_asset_entry("test_glyphs", width, height)],
AssetsPayloadSource::from_bytes(data),
gfx_installer,
sound_installer,
scene_installer,
);
// Initially zero
let info = am.bank_telemetry();
assert_eq!(info[0].bank_type, BankType::GLYPH);
assert_eq!(info[0].used_slots, 0);
assert_eq!(info[0].total_slots, 16);
// Loading
let handle = am.load(0, 0).expect("load must allocate handle");
// While LOADING or READY, it should be in inflight_bytes
let start = Instant::now();
while am.status(handle) != LoadStatus::READY && start.elapsed().as_secs() < 5 {
thread::sleep(std::time::Duration::from_millis(10));
}
let info = am.bank_telemetry();
assert_eq!(info[0].used_slots, 0);
// Commit
am.commit(handle);
am.apply_commits();
let info = am.bank_telemetry();
assert_eq!(info[0].used_slots, 1);
assert_eq!(info[1].bank_type, BankType::SOUNDS);
assert_eq!(info[1].used_slots, 0);
// Shutdown resets
am.shutdown();
let info = am.bank_telemetry();
assert_eq!(info[0].used_slots, 0);
}
}