use super::*;
use kcode_k1_audio_classification_event_format::{
AudioClassificationEventV3 as Event, DiscardedV2, FailedV2, FragmentStageV1, ProgressUpdateV1,
ProgressV1, 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;
#[derive(Default)]
struct Engine {
started: Mutex<Vec<FragmentId>>,
aborted: AtomicUsize,
release: AtomicBool,
}
impl Engine {
fn starts(&self, id: FragmentId) -> usize {
lock(&self.started)
.iter()
.filter(|entry| **entry == id)
.count()
}
}
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);
while active() && !self.release.load(AtomicOrdering::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 submit(peering: &K1Peering, event: Event) {
let payload = encode_event(&event).expect("event");
peering
.submit_txn(subsystem(), &payload)
.expect("submit event");
}
fn queue(peering: &K1Peering, fragment_id: FragmentId) {
submit(
peering,
Event::Queue(QueueV2 {
audio_object_id: fragment_id,
}),
);
}
fn object(objects: &K1Objects, seed: u8) -> FragmentId {
objects
.save("", "audio/ogg", "", &[seed])
.expect("audio object")
}
fn state(coordinator: &AudioClassificationCoordinator, id: FragmentId) -> OverallState {
coordinator
.status(id)
.expect("status")
.expect("fragment")
.state
}
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 startup_replay_is_effect_free_and_schedules_the_complete_work_list() {
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);
submit(
&dependencies.peering,
Event::Progress(ProgressV1 {
fragment_id: running,
update: ProgressUpdateV1::LlmJobStarted {
sequence: 1,
stage: FragmentStageV1::Transcript,
name: "transcript".to_owned(),
},
}),
);
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::default());
let coordinator = open(root.path(), &dependencies, engine.clone());
wait_for(|| engine.starts(running) == 1 && engine.starts(queued) == 1);
assert_eq!(dependencies.ordering.tip(), old_tip);
assert_eq!(
AudioClassificationCoordinator::startup_replay_count(&root.path().join("projection")),
Ok(3)
);
assert_eq!(state(&coordinator, running), OverallState::Running);
assert_eq!(state(&coordinator, queued), OverallState::Queued);
let plan = coordinator
.resume_plan(running)
.expect("artifact read")
.expect("fragment");
assert_eq!(plan.fragment_id, running);
assert!(plan.transcript.is_none());
engine.release.store(true, AtomicOrdering::Release);
}
#[test]
fn live_lifecycle_health_and_independent_lanes_remain_coordinator_owned() {
let root = TempDir::new().expect("root");
let live_dependencies = dependencies(root.path());
let engine = Arc::new(Engine::default());
let coordinator = open(root.path(), &live_dependencies, engine.clone());
let fragments = [
object(&live_dependencies.objects, 3),
object(&live_dependencies.objects, 4),
object(&live_dependencies.objects, 5),
];
for fragment in fragments {
queue(&live_dependencies.peering, fragment);
}
wait_for(|| {
fragments
.iter()
.all(|fragment| engine.starts(*fragment) == 1)
});
coordinator
.driver
.start(fragments[0])
.expect("already active");
assert_eq!(engine.starts(fragments[0]), 1);
let isolated_root = TempDir::new().expect("isolated root");
let isolated_dependencies = dependencies(isolated_root.path());
let isolated_engine = Arc::new(Engine::default());
let isolated = open(
isolated_root.path(),
&isolated_dependencies,
isolated_engine.clone(),
);
let isolated_fragment = object(&isolated_dependencies.objects, 6);
queue(&isolated_dependencies.peering, isolated_fragment);
wait_for(|| isolated_engine.starts(isolated_fragment) == 1);
submit(
&live_dependencies.peering,
Event::Failed(FailedV2 {
fragment_id: fragments[0],
stage: FragmentStageV1::Queue,
llm_job_sequence: None,
error: "failed".to_owned(),
}),
);
wait_for(|| engine.aborted.load(AtomicOrdering::SeqCst) == 1);
submit(
&live_dependencies.peering,
Event::Discarded(DiscardedV2 {
fragment_id: fragments[1],
}),
);
wait_for(|| engine.aborted.load(AtomicOrdering::SeqCst) == 2);
coordinator
.callback
.react(fragments[2], ProjectionEffect::LabelsCommitted)
.expect("labels committed");
wait_for(|| engine.aborted.load(AtomicOrdering::SeqCst) == 3);
assert_eq!(state(&coordinator, fragments[0]), OverallState::Failed);
assert_eq!(state(&coordinator, fragments[1]), OverallState::Discarded);
assert_eq!(state(&coordinator, fragments[2]), OverallState::Queued);
coordinator
.inject_errors(fragments[2], vec!["retained".to_owned()])
.expect("inject error");
assert_eq!(
coordinator
.status(fragments[2])
.expect("status")
.expect("fragment")
.errors,
["retained"]
);
assert!(coordinator.validate_labels(fragments[2], &[]).is_err());
fn typed_identity(label: &SpeakerLabelV1) {
let _: &Option<kcode_k1_person_types::PersonId> = &label.person_id;
}
let _: fn(&SpeakerLabelV1) = typed_identity;
assert!(
live_dependencies
.peering
.submit_txn(subsystem(), b"invalid")
.is_err()
);
let callback_fault = coordinator.ensure_healthy().expect_err("callback fault");
coordinator.shutdown();
coordinator.shutdown();
assert_eq!(coordinator.ensure_healthy(), Err(callback_fault));
Subsystem::reorg(isolated.callback.as_ref()).expect("clear projection");
assert_eq!(
isolated
.projection
.status(isolated_fragment)
.expect("projection"),
None
);
let reorg_fault = isolated.ensure_healthy().expect_err("reorg fault");
isolated.shutdown();
isolated.shutdown();
assert_eq!(isolated.ensure_healthy(), Err(reorg_fault));
}