use super::*;
use kcode_k1_audio_classification_format::{
AudioClassificationEventV3, ProgressUpdateV1, QueueV2, encode_event,
};
use kcode_k1_txn_ordering::REGISTER_AT_TIP;
use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
use std::thread;
use std::time::{Duration, Instant};
use tempfile::TempDir;
struct Engine {
ordering: Arc<K1TxnOrdering>,
started: Mutex<Vec<(FragmentId, Option<TxId>)>>,
aborted: AtomicUsize,
release: AtomicBool,
}
impl Engine {
fn new(ordering: Arc<K1TxnOrdering>) -> Self {
Self {
ordering,
started: Mutex::new(Vec::new()),
aborted: AtomicUsize::new(0),
release: AtomicBool::new(false),
}
}
fn start_tip(&self, id: FragmentId) -> Option<Option<TxId>> {
lock(&self.started)
.iter()
.find(|entry| entry.0 == id)
.map(|entry| entry.1)
}
}
impl FragmentEngine for Engine {
fn run<'a>(
&'a self,
fragment_id: FragmentId,
_ogg_bytes: &'a [u8],
active: &'a (dyn Fn() -> bool + Send + Sync),
) -> EngineFuture<'a> {
Box::pin(async move {
lock(&self.started).push((fragment_id, self.ordering.tip()));
while active() && !self.release.load(Ordering::Acquire) {
thread::yield_now();
}
if !active() {
self.aborted.fetch_add(1, AtomicOrdering::SeqCst);
}
Ok(())
})
}
}
struct Dependencies {
ordering: Arc<K1TxnOrdering>,
peering: Arc<K1Peering>,
objects: Arc<K1Objects>,
}
fn dependencies(root: &Path) -> Dependencies {
let ordering = Arc::new(K1TxnOrdering::open(&root.join("ordering")).expect("ordering"));
let peering =
Arc::new(K1Peering::open(&root.join("peering"), ordering.clone()).expect("peering"));
let objects = Arc::new(K1Objects::open(ordering.clone(), peering.clone()).expect("objects"));
Dependencies {
ordering,
peering,
objects,
}
}
fn open(
root: &Path,
dependencies: &Dependencies,
engine: Arc<Engine>,
) -> AudioClassificationCoordinator {
AudioClassificationCoordinator::with_engine(
&root.join("projection"),
dependencies.ordering.clone(),
dependencies.peering.clone(),
dependencies.objects.clone(),
engine,
)
.expect("coordinator")
}
fn subsystem() -> SubsystemId {
subsystem_id().expect("subsystem")
}
fn queue(peering: &K1Peering, fragment_id: FragmentId) -> TxId {
let payload = encode_event(&AudioClassificationEventV3::Queue(QueueV2 {
audio_object_id: fragment_id,
}))
.expect("queue event");
peering.submit_txn(subsystem(), &payload).expect("queue")
}
fn object(objects: &K1Objects, seed: u8) -> FragmentId {
objects
.save("", "audio/ogg", "", &[seed])
.expect("audio object")
}
fn wait_for(condition: impl Fn() -> bool) {
let deadline = Instant::now() + Duration::from_secs(2);
while !condition() && Instant::now() < deadline {
thread::sleep(Duration::from_millis(2));
}
assert!(condition());
}
struct Noop;
impl Subsystem for Noop {
fn submit_txn(&self, _id: TxId, _payload: &[u8]) -> Result<(), String> {
Ok(())
}
fn reorg(&self) -> Result<(), String> {
Ok(())
}
}
#[test]
fn replay_suppresses_starts_and_commits_restart_failure_before_queued_start() {
let root = TempDir::new().expect("root");
let (running, queued, old_tip) = {
let dependencies = dependencies(root.path());
dependencies
.ordering
.register_subsystem(subsystem(), Some(REGISTER_AT_TIP), Arc::new(Noop))
.expect("temporary registration");
let running = object(&dependencies.objects, 1);
queue(&dependencies.peering, running);
transactions::submit_progress(
&dependencies.peering,
running,
ProgressUpdateV1::LlmJobStarted {
sequence: 1,
stage: kcode_k1_audio_classification_projection::FragmentStageV1::Transcript,
name: "transcript".to_owned(),
},
)
.expect("running progress");
let queued = object(&dependencies.objects, 2);
queue(&dependencies.peering, queued);
(running, queued, dependencies.ordering.tip())
};
let dependencies = dependencies(root.path());
let engine = Arc::new(Engine::new(dependencies.ordering.clone()));
let coordinator = open(root.path(), &dependencies, engine.clone());
assert_eq!(
coordinator
.status(running)
.expect("status")
.expect("running")
.state,
OverallState::Failed
);
wait_for(|| engine.start_tip(queued).is_some());
assert_eq!(engine.start_tip(running), None);
assert_ne!(engine.start_tip(queued).expect("queued start"), old_tip);
engine.release.store(true, Ordering::Release);
}
#[test]
fn live_start_terminal_none_invalid_callback_reorg_and_independent_lanes() {
let root = TempDir::new().expect("root");
let live_dependencies = dependencies(root.path());
let engine = Arc::new(Engine::new(live_dependencies.ordering.clone()));
let _coordinator = open(root.path(), &live_dependencies, engine.clone());
let fragment = object(&live_dependencies.objects, 3);
queue(&live_dependencies.peering, fragment);
wait_for(|| engine.start_tip(fragment).is_some());
let isolated_root = TempDir::new().expect("isolated root");
let isolated_dependencies = dependencies(isolated_root.path());
let isolated_engine = Arc::new(Engine::new(isolated_dependencies.ordering.clone()));
let isolated = open(
isolated_root.path(),
&isolated_dependencies,
isolated_engine.clone(),
);
let isolated_fragment = object(&isolated_dependencies.objects, 4);
queue(&isolated_dependencies.peering, isolated_fragment);
wait_for(|| isolated_engine.start_tip(isolated_fragment).is_some());
Subsystem::reorg(isolated.callback.as_ref()).expect("clear projection");
assert!(isolated.ensure_healthy().is_err());
transactions::submit_failure(
&live_dependencies.peering,
fragment,
kcode_k1_audio_classification_projection::FragmentStageV1::Queue,
None,
"failed".to_owned(),
)
.expect("failure");
wait_for(|| engine.aborted.load(AtomicOrdering::SeqCst) == 1);
let invalid_root = TempDir::new().expect("invalid root");
let invalid_dependencies = dependencies(invalid_root.path());
let invalid = open(
invalid_root.path(),
&invalid_dependencies,
Arc::new(Engine::new(invalid_dependencies.ordering.clone())),
);
assert!(
invalid_dependencies
.peering
.submit_txn(subsystem(), b"invalid")
.is_err()
);
assert!(invalid.ensure_healthy().is_err());
}