kcode-k1-audio-classification-coordinator 0.2.3

Lifecycle coordinator for K1 audio classification
Documentation
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 speaker_identity_uses_typed_optional_person_id() {
    fn assert_type(label: &SpeakerLabelV1) {
        let _: &Option<kcode_k1_audio_classification_format::PersonId> = &label.person_id;
    }
    let _: fn(&SpeakerLabelV1) = assert_type;
}

#[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());
}