kcode-k1-audio-classification-coordinator 0.2.5

Lifecycle coordinator for K1 audio classification
Documentation
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);
    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));
}