mirror of
https://github.com/ruvnet/RuView
synced 2026-08-06 19:51:43 +00:00
38676aa2bd
Iter 34. Closes the gap where BfldPipelineHandle had no path for an
operator-supplied SoulMatchOracle to reach the worker thread. The
emit_with_oracle surface added in iter 14 was unreachable through the
handle API — Soul Signature deployments (ADR-118 §1.4) had to either
drop down to BfldEmitter directly or accept Recalibrate gate-drops on
known-enrolled matches.
Added (in src/pipeline.rs):
- BfldPipeline::process_with_oracle<O: SoulMatchOracle>(
inputs, embedding, oracle,
) -> Option<BfldEvent>
Wraps emitter.emit_with_oracle then applies the same privacy_mode
post-processing as process(). Privacy_mode and oracle are independent
— class-3 demote still happens AFTER any oracle Recalibrate exemption.
Added (in src/pipeline_handle.rs):
- BfldPipelineHandle::spawn_with_oracle<P, O>(pipeline, publisher, oracle) -> Self
where O: SoulMatchOracle + Send + Sync + 'static
The worker thread owns the oracle and consults it on every recv().
Worker loop now calls pipeline.process_with_oracle(...) instead of
pipeline.process(...).
tests/handle_soul_oracle.rs (3 named tests, all green):
spawn_with_oracle_null_is_equivalent_to_spawn
Parity: 3 identical low-risk inputs through spawn() and
spawn_with_oracle(NullOracle) produce the same publish count
and the same motion-topic count.
spawn_with_always_match_oracle_lets_events_publish_under_high_risk
*** Headline test ***
3 high-risk inputs spaced > DEBOUNCE_NS apart. With AlwaysMatch
oracle, all 3 produce motion topics — the gate never reaches
Recalibrate because the oracle reports an enrolled-person match.
spawn_with_null_oracle_drops_events_under_sustained_recalibrate_score
Negative control for the above: same 3 inputs through NullOracle,
only 1 motion topic survives (the first input lands at Accept;
the second and third hit Recalibrate after debounce and are
dropped per ADR-121 §2.4).
ADR-124 status (iter step 0 sibling check):
- docs/adr/ADR-124-rvagent-mcp-ruvector-npm-integration.md unchanged
at 431 lines. SENSE-BRIDGE scope remains orthogonal to BFLD core;
no overlap with this iter.
ACs progressed:
- ADR-118 §1.4 Soul Signature companion contract end-to-end through
the public handle API. Operators wiring Soul Signature into a
RuView deployment now use:
BfldPipelineHandle::spawn_with_oracle(pipeline, publisher, my_oracle)
…and the rest of the per-frame flow stays identical to spawn().
- ADR-121 §2.6 Recalibrate exemption proven over the worker-thread
boundary, not just at the unit level (iter 12 covered the gate-only
case).
Test config:
- cargo test --no-default-features → 72 passed
- cargo test → 227 passed (224 + 3)
Out of scope (next iter target):
- GitHub Actions workflow with mosquitto Docker (lifts iters 24+29
live-broker e2e from skip-mode). Remaining unmet ACs require
either external resources (KIT BFId, Pi5/Nexmon) or CI infra.
Co-Authored-By: claude-flow <ruv@ruv.net>
135 lines
5.1 KiB
Rust
135 lines
5.1 KiB
Rust
//! `BfldPipelineHandle` — worker-thread wrapper around [`BfldPipeline`] and a
|
|
//! [`Publish`]er. ADR-118 §2.1 single-call operator surface.
|
|
//!
|
|
//! `spawn()` returns a handle owning the inbound channel sender. The worker
|
|
//! thread loops on `recv()`, drives one `pipeline.process()` per input, and
|
|
//! forwards any emitted `BfldEvent` through `publish_event()`. `shutdown()`
|
|
//! closes the channel and joins the thread.
|
|
|
|
#![cfg(feature = "std")]
|
|
|
|
use std::sync::mpsc::{channel, RecvError, SendError, Sender};
|
|
use std::thread::{self, JoinHandle};
|
|
|
|
use crate::coherence_gate::SoulMatchOracle;
|
|
use crate::mqtt_topics::{publish_event, Publish};
|
|
use crate::pipeline::BfldPipeline;
|
|
use crate::{IdentityEmbedding, SensingInputs};
|
|
|
|
/// Frame-level input to the spawned worker. The pipeline state — gate,
|
|
/// embedding ring, hasher — lives behind the worker thread; callers only
|
|
/// send the per-frame sensing data.
|
|
pub struct PipelineInput {
|
|
/// Sensing fields fed to `pipeline.process`.
|
|
pub inputs: SensingInputs,
|
|
/// Optional embedding for the iter-15 hasher input + iter-8 ring.
|
|
pub embedding: Option<IdentityEmbedding>,
|
|
}
|
|
|
|
/// Handle to the spawned worker. Drop or `shutdown()` to stop. `send()`
|
|
/// returns an error after shutdown.
|
|
pub struct BfldPipelineHandle {
|
|
sender: Sender<PipelineInput>,
|
|
worker: Option<JoinHandle<()>>,
|
|
}
|
|
|
|
impl BfldPipelineHandle {
|
|
/// Spawn a worker that owns `pipeline` and `publisher`. Returns a handle
|
|
/// whose `send()` enqueues sensing inputs into the worker thread.
|
|
///
|
|
/// Publish errors are logged to stderr and the worker continues — single
|
|
/// frame failures should not kill the long-running pipeline.
|
|
#[must_use]
|
|
pub fn spawn<P>(mut pipeline: BfldPipeline, mut publisher: P) -> Self
|
|
where
|
|
P: Publish + Send + 'static,
|
|
P::Error: core::fmt::Debug,
|
|
{
|
|
let (sender, receiver) = channel::<PipelineInput>();
|
|
let worker = thread::spawn(move || loop {
|
|
match receiver.recv() {
|
|
Ok(PipelineInput { inputs, embedding }) => {
|
|
if let Some(event) = pipeline.process(inputs, embedding) {
|
|
if let Err(e) = publish_event(&mut publisher, &event) {
|
|
eprintln!("BFLD publish error: {e:?}");
|
|
}
|
|
}
|
|
}
|
|
Err(RecvError) => break, // channel closed by shutdown / drop
|
|
}
|
|
});
|
|
Self {
|
|
sender,
|
|
worker: Some(worker),
|
|
}
|
|
}
|
|
|
|
/// Variant of [`Self::spawn`] that installs a long-lived
|
|
/// [`SoulMatchOracle`] used on every per-frame `process` call. The oracle
|
|
/// must be `Send + Sync + 'static` because the worker thread consults it
|
|
/// on every recv. Pairs with ADR-121 §2.6: when the oracle reports a
|
|
/// `Match`, a would-be Recalibrate gate transition is downgraded to
|
|
/// `PredictOnly` (high score is the *intended* outcome of a known-enrolled
|
|
/// person match, not an attacker-grade sniffer arrival).
|
|
#[must_use]
|
|
pub fn spawn_with_oracle<P, O>(
|
|
mut pipeline: BfldPipeline,
|
|
mut publisher: P,
|
|
oracle: O,
|
|
) -> Self
|
|
where
|
|
P: Publish + Send + 'static,
|
|
P::Error: core::fmt::Debug,
|
|
O: SoulMatchOracle + Send + Sync + 'static,
|
|
{
|
|
let (sender, receiver) = channel::<PipelineInput>();
|
|
let worker = thread::spawn(move || loop {
|
|
match receiver.recv() {
|
|
Ok(PipelineInput { inputs, embedding }) => {
|
|
if let Some(event) =
|
|
pipeline.process_with_oracle(inputs, embedding, &oracle)
|
|
{
|
|
if let Err(e) = publish_event(&mut publisher, &event) {
|
|
eprintln!("BFLD publish error: {e:?}");
|
|
}
|
|
}
|
|
}
|
|
Err(RecvError) => break,
|
|
}
|
|
});
|
|
Self {
|
|
sender,
|
|
worker: Some(worker),
|
|
}
|
|
}
|
|
|
|
/// Enqueue an input. Returns `SendError<PipelineInput>` (carrying the
|
|
/// rejected input) if the worker has already shut down.
|
|
pub fn send(&self, input: PipelineInput) -> Result<(), SendError<PipelineInput>> {
|
|
self.sender.send(input)
|
|
}
|
|
|
|
/// Close the input channel and join the worker. Panics from the worker
|
|
/// thread propagate here; otherwise returns cleanly.
|
|
pub fn shutdown(mut self) {
|
|
if let Some(worker) = self.worker.take() {
|
|
drop(std::mem::replace(&mut self.sender, channel().0));
|
|
worker
|
|
.join()
|
|
.expect("BFLD pipeline worker panicked during shutdown");
|
|
}
|
|
}
|
|
}
|
|
|
|
impl Drop for BfldPipelineHandle {
|
|
/// Best-effort cleanup if `shutdown()` was not called explicitly.
|
|
fn drop(&mut self) {
|
|
if let Some(worker) = self.worker.take() {
|
|
// Replace the sender with a fresh disconnected one so the worker
|
|
// recv() returns Err(RecvError) and the loop exits.
|
|
drop(std::mem::replace(&mut self.sender, channel().0));
|
|
let _ = worker.join();
|
|
}
|
|
}
|
|
}
|