use std::sync::Arc;
use std::sync::mpsc::{Receiver, Sender};
use super::{StreamPlanner, StreamState};
use crate::bake::texture::TextureImage;
pub(crate) struct DecodedTexture {
pub image: TextureImage,
}
pub(crate) trait PayloadSource: Send + Sync {
fn fetch(&self, id: usize) -> Result<DecodedTexture, String>;
}
pub(crate) struct MemPayloadSource {
payloads: Vec<Vec<u8>>,
}
impl MemPayloadSource {
pub(crate) fn new(payloads: Vec<Vec<u8>>) -> Self {
Self { payloads }
}
}
impl PayloadSource for MemPayloadSource {
fn fetch(&self, id: usize) -> Result<DecodedTexture, String> {
let bytes = self
.payloads
.get(id)
.ok_or_else(|| format!("no payload for streamed texture {}", id))?;
let image = crate::bake::texture::deserialise(bytes)?;
Ok(DecodedTexture { image })
}
}
#[derive(Clone)]
pub(crate) struct DiskTextureLocator {
pub path: String,
pub(crate) file_offset: u64,
pub len: u64,
}
pub(crate) struct DiskPayloadSource {
locators: Vec<DiskTextureLocator>,
}
impl DiskPayloadSource {
pub(crate) fn new(locators: Vec<DiskTextureLocator>) -> Self {
Self { locators }
}
}
impl PayloadSource for DiskPayloadSource {
fn fetch(&self, id: usize) -> Result<DecodedTexture, String> {
let loc = self
.locators
.get(id)
.ok_or_else(|| format!("no disk locator for streamed texture {}", id))?;
let bytes = super::file_range::read_at(&loc.path, loc.file_offset, loc.len)?;
let image = crate::bake::texture::deserialise(&bytes)?;
Ok(DecodedTexture { image })
}
}
struct LoadResult {
id: usize,
decoded: Result<DecodedTexture, String>,
}
pub(crate) struct TextureStreamer {
planner: StreamPlanner,
centers: Vec<Vec<[f32; 3]>>,
worker: super::worker::Worker<usize>,
result_rx: Receiver<LoadResult>,
}
impl TextureStreamer {
pub(crate) fn new(
source: Arc<dyn PayloadSource>,
centers: Vec<Vec<[f32; 3]>>,
load_budget: usize,
resident_cap: usize,
) -> Self {
let planner = StreamPlanner::new(centers.len(), load_budget, resident_cap);
let (request_tx, request_rx) = std::sync::mpsc::channel::<usize>();
let (result_tx, result_rx) = std::sync::mpsc::channel::<LoadResult>();
let worker =
super::worker::Worker::spawn("cn-texture-stream", request_rx, request_tx, move |rx| {
worker_loop(source, rx, result_tx)
});
Self {
planner,
centers,
result_rx,
worker,
}
}
pub(crate) fn len(&self) -> usize {
self.planner.len()
}
pub(crate) fn set_byte_budget(&mut self, budget: Option<u64>) {
self.planner.set_byte_budget(budget);
}
pub(crate) fn resident_bytes(&self) -> u64 {
self.planner.resident_bytes()
}
pub(crate) fn set_blocked(&mut self, slot: usize, blocked: bool) {
self.planner.set_blocked(slot, blocked);
}
pub(crate) fn byte_budget(&self) -> Option<u64> {
self.planner.byte_budget()
}
pub(crate) fn update_scores(&mut self, camera: [f32; 3], frame: u64) {
for id in 0..self.planner.len() {
self.planner
.set_score(id, nearest_sq_distance(&self.centers[id], camera));
if self.planner.state(id) == Some(StreamState::Resident) {
self.planner.touch(id, frame);
}
}
}
pub(crate) fn plan_and_dispatch(&mut self) -> Vec<usize> {
let plan = self.planner.plan();
for &id in &plan.to_load {
let sent = self.worker.send(id);
if !sent {
self.planner.mark_unloaded(id);
}
}
plan.to_evict
}
pub(crate) fn drain_completed(
&mut self,
frame: u64,
mut upload: impl FnMut(usize, TextureImage),
) -> usize {
let mut applied = 0;
while let Ok(result) = self.result_rx.try_recv() {
match result.decoded {
Ok(tex) => {
let bytes = tex.image.byte_len() as u64;
upload(result.id, tex.image);
self.planner.mark_resident(result.id, frame, bytes);
applied += 1;
}
Err(e) => {
tracing::warn!("texture stream: load of slot {} failed: {}", result.id, e);
self.planner.mark_resident(result.id, frame, 0);
}
}
}
applied
}
pub(crate) fn stats(&self) -> (usize, usize, usize) {
self.planner.counts()
}
}
fn worker_loop(
source: Arc<dyn PayloadSource>,
requests: Receiver<usize>,
results: Sender<LoadResult>,
) {
while let Ok(id) = requests.recv() {
let decoded = source.fetch(id);
if results.send(LoadResult { id, decoded }).is_err() {
break;
}
}
}
fn nearest_sq_distance(centers: &[[f32; 3]], camera: [f32; 3]) -> f32 {
let mut nearest = f32::MAX;
for c in centers {
let dx = c[0] - camera[0];
let dy = c[1] - camera[1];
let dz = c[2] - camera[2];
let d = dx * dx + dy * dy + dz * dz;
if d < nearest {
nearest = d;
}
}
if centers.is_empty() { 0.0 } else { nearest }
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn nearest_sq_distance_picks_the_closest_center() {
let centers = [[10.0, 0.0, 0.0], [3.0, 0.0, 0.0], [7.0, 0.0, 0.0]];
assert_eq!(nearest_sq_distance(¢ers, [0.0, 0.0, 0.0]), 9.0);
}
#[test]
fn nearest_sq_distance_of_no_centers_is_zero() {
assert_eq!(nearest_sq_distance(&[], [5.0, 5.0, 5.0]), 0.0);
}
fn make_payload(w: u32, h: u32, fill: u8) -> Vec<u8> {
let pixels = std::iter::repeat_n(fill, (w * h * 4) as usize).collect();
crate::bake::texture::serialise(&TextureImage::rgba8(w, h, pixels))
}
#[test]
fn mem_payload_source_decodes_a_payload() {
let source = MemPayloadSource::new(vec![make_payload(2, 1, 0xAB)]);
let tex = source.fetch(0).expect("fetch ok");
assert_eq!((tex.image.width(), tex.image.height()), (2, 1));
assert_eq!(tex.image.mips[0].data.len(), 2 * 4);
assert!(tex.image.mips[0].data.iter().all(|&b| b == 0xAB));
}
#[test]
fn mem_payload_source_errors_on_unknown_id() {
let source = MemPayloadSource::new(vec![make_payload(1, 1, 0)]);
assert!(source.fetch(9).is_err());
}
#[test]
fn disk_payload_source_reads_a_payload_at_offset() {
let tree = concinnity_testing::TempTree::new();
let payload = make_payload(2, 1, 0xCD);
let prefix = vec![0u8; 37];
let mut bytes = prefix.clone();
bytes.extend_from_slice(&payload);
let path = tree.write("payload.bin", &bytes);
let source = DiskPayloadSource::new(vec![DiskTextureLocator {
path: path.to_string_lossy().into_owned(),
file_offset: prefix.len() as u64,
len: payload.len() as u64,
}]);
let tex = source.fetch(0).expect("fetch ok");
assert_eq!((tex.image.width(), tex.image.height()), (2, 1));
assert_eq!(tex.image.mips[0].data.len(), 2 * 4);
assert!(tex.image.mips[0].data.iter().all(|&b| b == 0xCD));
let _ = std::fs::remove_file(&path);
}
#[test]
fn disk_payload_source_errors_on_unknown_id() {
let source = DiskPayloadSource::new(vec![]);
assert!(source.fetch(0).is_err());
}
#[test]
fn disk_payload_source_errors_on_missing_file() {
let source = DiskPayloadSource::new(vec![DiskTextureLocator {
path: "/nonexistent/cn_disk_payload_missing.bin".to_string(),
file_offset: 0,
len: 4,
}]);
assert!(source.fetch(0).is_err());
}
struct ConstSource;
impl PayloadSource for ConstSource {
fn fetch(&self, _id: usize) -> Result<DecodedTexture, String> {
Ok(DecodedTexture {
image: TextureImage::rgba8(1, 1, vec![1, 2, 3, 4]),
})
}
}
fn drain_until(streamer: &mut TextureStreamer, frame: u64, want: usize) -> usize {
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
let mut uploads = 0;
while std::time::Instant::now() < deadline {
uploads += streamer.drain_completed(frame, |_, _| {});
if streamer.stats().0 >= want {
break;
}
std::thread::sleep(std::time::Duration::from_millis(1));
}
uploads
}
#[test]
fn streamer_loads_nearest_slots_within_budget() {
let centers = vec![
vec![[100.0, 0.0, 0.0]], vec![[2.0, 0.0, 0.0]], vec![[50.0, 0.0, 0.0]], ];
let mut streamer = TextureStreamer::new(Arc::new(ConstSource), centers, 1, 8);
assert_eq!(streamer.len(), 3);
streamer.update_scores([0.0, 0.0, 0.0], 1);
let evict = streamer.plan_and_dispatch();
assert!(evict.is_empty());
drain_until(&mut streamer, 1, 1);
assert_eq!(streamer.stats().0, 1);
streamer.update_scores([0.0, 0.0, 0.0], 2);
streamer.plan_and_dispatch();
drain_until(&mut streamer, 2, 2);
assert_eq!(streamer.stats().0, 2);
streamer.update_scores([0.0, 0.0, 0.0], 3);
streamer.plan_and_dispatch();
drain_until(&mut streamer, 3, 3);
assert_eq!(streamer.stats(), (3, 0, 0));
}
struct FailingSource;
impl PayloadSource for FailingSource {
fn fetch(&self, _id: usize) -> Result<DecodedTexture, String> {
Err("undecodable payload".to_string())
}
}
#[test]
fn a_failed_fetch_is_not_retried() {
let centers = vec![vec![[1.0, 0.0, 0.0]]];
let mut streamer = TextureStreamer::new(Arc::new(FailingSource), centers, 4, 8);
streamer.update_scores([0.0, 0.0, 0.0], 1);
streamer.plan_and_dispatch();
let uploads = drain_until(&mut streamer, 1, 1);
assert_eq!(uploads, 0, "a failed load uploads no pixels");
assert_eq!(streamer.stats(), (1, 0, 0));
assert_eq!(streamer.resident_bytes(), 0);
}
#[test]
fn upload_callback_receives_decoded_pixels() {
let centers = vec![vec![[1.0, 0.0, 0.0]]];
let mut streamer = TextureStreamer::new(Arc::new(ConstSource), centers, 4, 8);
streamer.update_scores([0.0, 0.0, 0.0], 1);
streamer.plan_and_dispatch();
let mut seen: Option<(usize, u32, u32, Vec<u8>)> = None;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
while std::time::Instant::now() < deadline && seen.is_none() {
streamer.drain_completed(1, |id, image| {
seen = Some((
id,
image.width(),
image.height(),
image.mips[0].data.clone(),
));
});
std::thread::sleep(std::time::Duration::from_millis(1));
}
assert_eq!(seen, Some((0, 1, 1, vec![1, 2, 3, 4])));
}
}