From 9b21f8861c8f9756551c28b88fc49799b2fb5988 Mon Sep 17 00:00:00 2001 From: TheM14 Date: Tue, 6 Oct 2026 15:03:49 +0800 Subject: [PATCH] Fix the objects16 input layout and bound its render-ahead --- CMakeLists.txt | 4 +- include/joc_stream.h | 8 +-- src/simd/cpu_probe.cpp | 4 +- src/stream/stream.cpp | 125 +++++++++++++++++++++++++++++------------ src/stream/stream.h | 21 ++++--- tests/test_core.cpp | 80 ++++++++++++++++++++++++++ 6 files changed, 190 insertions(+), 52 deletions(-) diff --git a/CMakeLists.txt b/CMakeLists.txt index b389ebf..c58ae30 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -260,7 +260,9 @@ endif() if(JOC_BUILD_TESTS) enable_testing() add_executable(joc_tests tests/test_core.cpp) - target_link_libraries(joc_tests PRIVATE joc_core_impl) + # The stream tests reach the core the shared library exports, so the static + # library alone is not enough to link them. + target_link_libraries(joc_tests PRIVATE joc_core_impl joc_core) target_include_directories(joc_tests PRIVATE "${CMAKE_CURRENT_SOURCE_DIR}/src") add_test(NAME core COMMAND joc_tests) # The public headers must stay valid C and C++: these targets exist to prove it. diff --git a/include/joc_stream.h b/include/joc_stream.h index 4e04962..280b37a 100644 --- a/include/joc_stream.h +++ b/include/joc_stream.h @@ -39,14 +39,14 @@ typedef struct joc_stream joc_stream; typedef enum joc_stream_input { JOC_STREAM_IN_EAC3 = 0, /* bare E-AC-3 syncframes (the metadata stream) */ - JOC_STREAM_IN_PCM_OBJECTS16 = 1, /* 16-channel objects16, decoded by the host */ + JOC_STREAM_IN_PCM_OBJECTS16 = 1, /* objects16 frames, planar [16][1536] each */ JOC_STREAM_IN_CORE_PCM = 3 /* the 5.1 core PCM of the pushed E-AC-3 frames */ } joc_stream_input; typedef enum joc_stream_output { - JOC_STREAM_OUT_PCM_OBJECTS16 = 0, /* planar [16][samples] float32 */ - JOC_STREAM_OUT_SPEAKER = 1, /* interleaved [samples][channels] f32 */ - JOC_STREAM_OUT_BINAURAL = 2 /* interleaved [samples][2] f32 */ + JOC_STREAM_OUT_PCM_OBJECTS16 = 0, /* objects16 frames, planar [16][1536] each */ + JOC_STREAM_OUT_SPEAKER = 1, /* interleaved [samples][channels] f32 */ + JOC_STREAM_OUT_BINAURAL = 2 /* interleaved [samples][2] f32 */ } joc_stream_output; typedef struct joc_stream_config { diff --git a/src/simd/cpu_probe.cpp b/src/simd/cpu_probe.cpp index 206d973..fe03e99 100644 --- a/src/simd/cpu_probe.cpp +++ b/src/simd/cpu_probe.cpp @@ -22,9 +22,7 @@ namespace joc::simd { namespace { // ---------------------------------------------------------------------- x86 -- -// Both pointer sizes are probed: AVX2 is not an x86-64-only ISA, and gating this -// on _M_X64 / __x86_64__ left every 32-bit x86 build reporting "no features", -// which pinned the dispatcher to the scalar kernels. +// AVX2 is not an x86-64-only ISA, so both pointer sizes are probed. #if defined(_M_X64) || defined(_M_IX86) || defined(__x86_64__) || defined(__i386__) #if defined(_M_X64) || defined(_M_IX86) diff --git a/src/stream/stream.cpp b/src/stream/stream.cpp index 399d5ea..f2b2d5e 100644 --- a/src/stream/stream.cpp +++ b/src/stream/stream.cpp @@ -38,7 +38,10 @@ void Stream::reset_state() { reader_ = eac3::FrameReader(); metadata_.clear(); bed_pending_.clear(); + bed_read_offset_ = 0; objects16_.clear(); + objects_pending_.clear(); + objects_read_offset_ = 0; output_.clear(); read_offset_ = 0; info_ = Info(); @@ -245,54 +248,82 @@ Status Stream::push_objects16(const float* planar16, std::size_t samples, std::s if (planar16 == nullptr || samples == 0u) { return Status::success(); } - // Rendered immediately: the host has already done the JOC rebuild. + // Queued as whole frames, each planar [16][kFrameSamples]: the shape the + // objects16 output writes, so a host can feed a batch back unchanged. + constexpr std::size_t kFrameValues = + static_cast(JOC_OUTPUT_CHANNELS) * kFrameSamples; for (std::size_t offset = 0; offset < samples; offset += kFrameSamples) { const std::size_t count = std::min(kFrameSamples, samples - offset); - std::vector frame(static_cast(JOC_OUTPUT_CHANNELS) * kFrameSamples, - 0.0f); + const std::size_t base = objects_pending_.size(); + const float* frame = planar16 + (offset / kFrameSamples) * kFrameValues; + objects_pending_.resize(base + kFrameValues, 0.0f); for (std::size_t channel = 0; channel < JOC_OUTPUT_CHANNELS; ++channel) { - std::memcpy(frame.data() + channel * kFrameSamples, - planar16 + channel * samples + offset, count * sizeof(float)); + std::memcpy(objects_pending_.data() + base + channel * kFrameSamples, + frame + channel * kFrameSamples, count * sizeof(float)); } - ++info_.frames_in; - info_.samples_in += count; - const Status rendered = render_objects16(frame); + } + return process_objects16_frames(false); +} + +Status Stream::process_objects16_frames(bool drain_all) { + constexpr std::size_t kFrameValues = + static_cast(JOC_OUTPUT_CHANNELS) * kFrameSamples; + while (objects_pending_.size() - objects_read_offset_ >= kFrameValues) { + // Leave the rest queued, in order, for a later push or for flush(). + if (!drain_all && buffered_samples() >= kMaxRenderAheadSamples) { + break; + } + const auto first = objects_pending_.begin() + + static_cast(objects_read_offset_); + objects_frame_.assign(first, first + static_cast(kFrameValues)); + objects_read_offset_ += kFrameValues; + const Status rendered = render_objects16(objects_frame_); if (!rendered.ok()) { return rendered; } - if (count != kFrameSamples) { - break; // a partial frame is dropped; the host should push whole frames - } + ++info_.frames_in; + info_.samples_in += kFrameSamples; + } + if (objects_read_offset_ == objects_pending_.size()) { + objects_pending_.clear(); + objects_read_offset_ = 0; + } else if (objects_read_offset_ >= (1u << 20)) { + // Erasing from the front moves the remainder, so it is only worth doing + // once the consumed prefix is large enough to pay for the move. + objects_pending_.erase( + objects_pending_.begin(), + objects_pending_.begin() + static_cast(objects_read_offset_)); + objects_read_offset_ = 0; } return Status::success(); } Status Stream::process_ready_frames(bool drain_all) { - while (bed_pending_.size() / kBedChannels >= kFrameSamples && !metadata_.empty()) { - // Stop before rendering what the caller is not about to take: the frames - // stay queued, in order, and are rendered by a later push or by flush(). + while (bed_pending_.size() - bed_read_offset_ >= kFrameSamples * kBedChannels && + !metadata_.empty()) { + // Leave the rest queued, in order, for a later push or for flush(). if (!drain_all && buffered_samples() >= kMaxRenderAheadSamples) { break; } const FrameMetadata entry = metadata_.front(); metadata_.pop_front(); - std::vector bed5(static_cast(JOC_CORE_CHANNELS) * kFrameSamples, 0.0f); - std::vector lfe(kFrameSamples, 0.0f); + const float* bed = bed_pending_.data() + bed_read_offset_; + bed5_.resize(static_cast(JOC_CORE_CHANNELS) * kFrameSamples); + lfe_.resize(kFrameSamples); for (std::size_t sample = 0; sample < kFrameSamples; ++sample) { for (std::size_t channel = 0; channel < JOC_CORE_CHANNELS; ++channel) { - bed5[channel * kFrameSamples + sample] = - bed_pending_[sample * kBedChannels + kCoreChannels[channel]]; + bed5_[channel * kFrameSamples + sample] = + bed[sample * kBedChannels + kCoreChannels[channel]]; } - lfe[sample] = bed_pending_[sample * kBedChannels + kLfeChannel]; + lfe_[sample] = bed[sample * kBedChannels + kLfeChannel]; } - bed_pending_.erase(bed_pending_.begin(), - bed_pending_.begin() + static_cast(kFrameSamples * - kBedChannels)); + bed_read_offset_ += kFrameSamples * kBedChannels; + compact_bed_pending(); std::string error; - const Status rebuilt = joc::rebuild_objects16(rebuilder_, entry.params, bed5.data(), - lfe.data(), gain_, &objects16_, &error); + const Status rebuilt = joc::rebuild_objects16(rebuilder_, entry.params, bed5_.data(), + lfe_.data(), gain_, &objects16_, &error); if (!rebuilt.ok()) { return Status::fail(rebuilt.code(), stage::kDsp, error); } @@ -307,6 +338,22 @@ Status Stream::process_ready_frames(bool drain_all) { return Status::success(); } +void Stream::compact_bed_pending() { + if (bed_read_offset_ == 0) { + return; + } + if (bed_read_offset_ == bed_pending_.size()) { + bed_pending_.clear(); + bed_read_offset_ = 0; + } else if (bed_read_offset_ >= (1u << 20)) { + // Erasing from the front moves the remainder, so it is only worth doing + // once the consumed prefix is large enough to pay for the move. + bed_pending_.erase(bed_pending_.begin(), + bed_pending_.begin() + static_cast(bed_read_offset_)); + bed_read_offset_ = 0; + } +} + Status Stream::render_objects16(const std::vector& objects16) { if (config_.output == JOC_STREAM_OUT_PCM_OBJECTS16) { output_.insert(output_.end(), objects16.begin(), objects16.end()); @@ -324,6 +371,7 @@ Status Stream::render_objects16(const std::vector& objects16) { if (!stepped.ok()) { return Status::fail(stepped.code(), stage::kRender, error); } + output_.reserve(output_.size() + speaker_.output.size()); for (const double value : speaker_.output) { output_.push_back(static_cast(value)); } @@ -344,13 +392,13 @@ Status Stream::render_objects16(const std::vector& objects16) { if (!submitted.ok()) { return submitted; } - std::vector produced; - binaural_.take_output(&produced); - for (const double value : produced) { + binaural_.take_output(&produced_); + output_.reserve(output_.size() + produced_.size()); + for (const double value : produced_) { output_.push_back(static_cast(value)); } info_.frames_out++; - info_.samples_out += produced.size() / 2u; + info_.samples_out += produced_.size() / 2u; return Status::success(); } @@ -371,16 +419,15 @@ Status Stream::render_rosella_objects16(const std::vector& objects16) { if (!submitted.ok()) { return submitted; } - std::vector produced; - rosella_.take_output(&produced); - if (!produced.empty()) { - rosella_pending_.insert(rosella_pending_.end(), produced.begin(), produced.end()); + rosella_.take_output(&produced_); + if (!produced_.empty()) { + rosella_pending_.insert(rosella_pending_.end(), produced_.begin(), produced_.end()); } release_rosella_output(kFrameSamples); info_.frames_out++; // Counted as the runtime produces it, which is also how the SOFA path counts: // the totals are identical, only the frame they appear on differs. - info_.samples_out += produced.size() / 2u; + info_.samples_out += produced_.size() / 2u; return Status::success(); } @@ -393,6 +440,7 @@ void Stream::release_rosella_output(std::size_t limit) { return; } const std::size_t values = count * 2u; + output_.reserve(output_.size() + values); for (std::size_t index = 0; index < values; ++index) { output_.push_back(static_cast(rosella_pending_[rosella_read_offset_ + index])); } @@ -434,13 +482,16 @@ Status Stream::pull(float* destination, std::size_t capacity_samples, std::size_ } Status Stream::flush() { - // Input has ended, so the render-ahead bound has nothing left to wait for: - // every frame still queued has to reach the renderer before its tail is - // drained, or the end of the file would be dropped. + // Input has ended, so drain what the cap held back: nothing else will + // trigger rendering. const Status remaining = process_ready_frames(true); if (!remaining.ok()) { return remaining; } + const Status objects = process_objects16_frames(true); + if (!objects.ok()) { + return objects; + } if (binaural_ready_) { std::vector tail; const Status drained = @@ -448,6 +499,7 @@ Status Stream::flush() { if (!drained.ok()) { return drained; } + output_.reserve(output_.size() + tail.size()); for (const double value : tail) { output_.push_back(static_cast(value)); } @@ -464,6 +516,7 @@ Status Stream::flush() { // tail only sounds after it. The program samples were already counted by // render_rosella_objects16, so only the tail is added here. release_rosella_output(rosella_pending_samples()); + output_.reserve(output_.size() + tail.size()); for (const double value : tail) { output_.push_back(static_cast(value)); } diff --git a/src/stream/stream.h b/src/stream/stream.h index e0ad56a..7f686af 100644 --- a/src/stream/stream.h +++ b/src/stream/stream.h @@ -85,13 +85,8 @@ public: : 0u; } - // A push renders every frame it makes ready, and the caller decides how far - // its demuxer runs ahead of playback. Without a bound, a demuxer that runs - // far ahead turns its whole read-ahead burst into latency on whichever pull() - // happens to follow it: the samples are not wasted, but they are rendered at - // the worst possible moment. Rendering therefore stops once this many - // samples are rendered and unpulled; flush() lifts the bound so the frames - // still waiting when the input ends are drained rather than dropped. + // Cap on rendered samples that have not been pulled. A push renders what it + // makes ready, so a caller that feeds faster than it pulls renders ahead. static constexpr std::size_t kMaxRenderAheadSamples = 16384; private: @@ -106,6 +101,10 @@ private: return rosella_ready_ ? (rosella_pending_.size() - rosella_read_offset_) / 2u : 0u; } void reset_state(); + // Drops the bed samples that have already been rendered, keeping the rest. + void compact_bed_pending(); + // Renders the queued objects16 frames, bounded by kMaxRenderAheadSamples. + Status process_objects16_frames(bool drain_all); Config config_; Info info_; @@ -113,7 +112,12 @@ private: std::deque metadata_; FrameMetadata pending_metadata_; std::vector bed_pending_; - std::vector frame_copy_; + std::size_t bed_read_offset_ = 0; + std::vector bed5_; + std::vector lfe_; + std::vector objects_pending_; + std::size_t objects_read_offset_ = 0; + std::vector objects_frame_; std::vector objects16_; std::vector output_; std::size_t read_offset_ = 0; @@ -125,6 +129,7 @@ private: hrtf::RosellaRuntime rosella_; std::vector rosella_pending_; std::size_t rosella_read_offset_ = 0; + std::vector produced_; bool speaker_enabled_ = false; bool binaural_enabled_ = false; bool binaural_ready_ = false; diff --git a/tests/test_core.cpp b/tests/test_core.cpp index e5458a3..918291f 100644 --- a/tests/test_core.cpp +++ b/tests/test_core.cpp @@ -1,5 +1,6 @@ // Unit tests for the engine's self-contained parts: no test data files, no // reference implementation, no external framework. Run with `ctest` or directly. +#include #include #include #include @@ -30,6 +31,7 @@ #include "io/wav_writer.h" #include "io/zip_reader.h" #include "oamd/oamd_parser.h" +#include "stream/stream.h" namespace { @@ -707,6 +709,83 @@ void test_builtin_kernel_tables() { "f5beb3220e4530fcf28e7f4da7f07e821074265d118c911d61a590e00753a573"); } +// Drives one objects16 stream over `block` and returns what it renders. With +// `frame_at_a_time` each push carries a single frame; otherwise the whole batch +// goes in at once. `peak` reports the largest backlog a push left behind. +std::vector drive_objects16_stream(const std::vector& block, std::size_t frames, + bool frame_at_a_time, std::size_t* peak) { + constexpr std::size_t kFrameValues = JOC_OUTPUT_CHANNELS * JOC_FRAME_SAMPLES; + joc::stream::Stream stream; + joc::stream::Config config; + config.input = JOC_STREAM_IN_PCM_OBJECTS16; + config.output = JOC_STREAM_OUT_SPEAKER; + config.layout = "5.1"; + const joc::Status created = stream.create(config); + CHECK(created.ok()); + if (!created.ok()) { + return {}; + } + const std::size_t channels = stream.info().output_channels; + CHECK(channels != 0u); + std::vector pulled(4096u * channels); + std::vector rendered; + const auto drain = [&] { + for (;;) { + std::size_t produced = 0; + const joc::Status status = stream.pull(pulled.data(), 4096u, &produced); + CHECK(status.ok()); + if (!status.ok() || produced == 0u) { + return; + } + rendered.insert(rendered.end(), pulled.begin(), + pulled.begin() + static_cast(produced * channels)); + } + }; + + std::size_t consumed = 0; + while (consumed < frames) { + const std::size_t count = frame_at_a_time ? 1u : frames - consumed; + std::size_t taken = 0; + const joc::Status pushed = stream.push_objects16( + block.data() + consumed * kFrameValues, count * JOC_FRAME_SAMPLES, &taken); + CHECK(pushed.ok()); + if (!pushed.ok()) { + return rendered; + } + // A push takes everything it is given; only the rendering is bounded. + CHECK(taken == count * JOC_FRAME_SAMPLES); + *peak = std::max(*peak, stream.buffered_samples()); + consumed += taken / JOC_FRAME_SAMPLES; + drain(); + } + CHECK(stream.flush().ok()); + drain(); + return rendered; +} + +// A push renders every frame it makes ready, so the stream bounds how far ahead of +// the caller it renders. A batch pushed at once must therefore leave a bounded +// backlog, and it must render the same samples as pushing one frame at a time. +void test_stream_render_ahead() { + constexpr std::size_t kFrames = 32; + constexpr std::size_t kFrameValues = JOC_OUTPUT_CHANNELS * JOC_FRAME_SAMPLES; + std::vector block(kFrames * kFrameValues); + for (std::size_t index = 0; index < block.size(); ++index) { + block[index] = static_cast(std::sin(static_cast(index) * 0.0007) * 0.25); + } + + std::size_t frame_peak = 0; + const std::vector reference = + drive_objects16_stream(block, kFrames, true, &frame_peak); + std::size_t batch_peak = 0; + const std::vector batched = drive_objects16_stream(block, kFrames, false, &batch_peak); + + CHECK(!reference.empty()); + CHECK(batched == reference); + // The batch is larger than the bound, so it cannot all be rendered at once. + CHECK(batch_peak < kFrames * JOC_FRAME_SAMPLES / 2u); +} + } // namespace int main() { @@ -728,6 +807,7 @@ int main() { test_sofa_reader_errors(); test_hrtf_cache_policy(); test_rosella_model_errors(); + test_stream_render_ahead(); std::printf("%d checks, %d failure(s)\n", g_checks, g_failures); return g_failures == 0 ? 0 : 1; }