kcode-k1-daemon-lib 0.14.6

Library-only K1 loopback daemon and authority Web composition root
Documentation
#![doc = include_str!("../Documentation.md")]

use kcode_gemini_3_1_pro::Gemini31Pro;
use kcode_k1_access::K1Access;
use kcode_k1_access_audio_classifier::K1AccessAudioClassifiers;
use kcode_k1_access_full_audio::K1AccessFullAudio;
use kcode_k1_access_launch_nodes::K1AccessLaunchNodes;
use kcode_k1_access_persons::K1AccessPersons;
use kcode_k1_access_profiles::K1AccessProfiles;
use kcode_k1_accounting::Accounting;
use kcode_k1_accounts::K1Accounts;
use kcode_k1_audio_classification::AudioClassification;
use kcode_k1_audio_classifier::K1AudioClassifiers;
use kcode_k1_authority_filters::K1AuthorityFilters;
use kcode_k1_bootstrap_state::K1BootstrapState;
use kcode_k1_chat_service::K1ChatService;
use kcode_k1_codex_adapter::Adapter as CodexAdapter;
use kcode_k1_codex_websearch::Runner as WebSearchRunner;
use kcode_k1_daemon_audio_lifetime::ClassificationLifetime;
use kcode_k1_daemon_code_services::CodeServices;
use kcode_k1_daemon_files::DaemonFiles;
use kcode_k1_daemon_http_boundary::{Boundary, warn_if_slow, write_readiness_for};
use kcode_k1_daemon_provider_config::{
    CODEX_EXECUTABLE_ENV, audio_access_model, chat_access_model, codex_configs, codex_executable,
    people_models, persons_access_model, resolve_ffmpeg,
};
use kcode_k1_daemon_startup_error::{StartupError, redacted_stage, with_cause};
use kcode_k1_daemon_vault_unlock::VaultUnlock;
use kcode_k1_daemon_web_startup::{open, select_public_origin};
use kcode_k1_ese::K1Ese;
use kcode_k1_full_audio::K1FullAudio;
use kcode_k1_groups::K1Groups;
use kcode_k1_invites::K1Invites;
use kcode_k1_ktool_set_launch_node::SetLaunchNodeKtool;
use kcode_k1_ktool_social::SocialKtools;
use kcode_k1_launch_nodes::LaunchNodes;
use kcode_k1_loom_bootstrap::{BootstrapServices, ensure_with_topology};
use kcode_k1_objects::K1Objects;
use kcode_k1_peering::K1Peering;
use kcode_k1_persons::K1Persons;
use kcode_k1_txn_ordering::K1TxnOrdering;
use kcode_k1_users::K1Users;
use kcode_k1_vault::{ExposeSecret, K1Vault};
use kcode_speaker_v3_analysis::Analyzer;
use std::path::PathBuf;
use std::process::ExitCode;
use std::sync::Arc;
use std::time::Instant;

const INVITE_LINK_URL: &str = "http://localhost:4321/lib/kcode-k1-ui/*/account.html";
const GEMINI_API_KEY: &str = "gemini-api-key";

struct Prepared {
    boundary: Boundary,
    classification: ClassificationLifetime,
    public_origin: String,
    unused_invites: usize,
    vault: Arc<K1Vault>,
}

pub fn run(k1_root: PathBuf) -> ExitCode {
    let runtime = match tokio::runtime::Builder::new_multi_thread()
        .enable_all()
        .build()
    {
        Ok(runtime) => runtime,
        Err(error) => {
            eprintln!("{}", with_cause("runtime")(error));
            return ExitCode::from(1);
        }
    };
    let unlock = match VaultUnlock::prompt() {
        Ok(unlock) => unlock,
        Err(error) => {
            eprintln!("{}", redacted_stage("Vault unlock")(error));
            return ExitCode::from(1);
        }
    };
    runtime.block_on(run_async(k1_root, unlock))
}

async fn run_async(k1_root: PathBuf, unlock: VaultUnlock) -> ExitCode {
    let started = Instant::now();
    let prepared = match startup(k1_root, unlock).await {
        Ok(prepared) => prepared,
        Err(error) => {
            warn_if_slow(started.elapsed(), "error");
            eprintln!("{error}");
            return ExitCode::from(1);
        }
    };
    let elapsed = started.elapsed();
    if write_readiness_for(&prepared.public_origin, prepared.unused_invites).is_err() {
        warn_if_slow(elapsed, "error");
        eprintln!(
            "{}",
            with_cause("readiness output")("write or flush failed")
        );
        return ExitCode::from(1);
    }
    warn_if_slow(elapsed, "ready");
    let Prepared {
        boundary,
        classification,
        vault,
        ..
    } = prepared;
    let result = boundary.serve().await;
    classification.shutdown();
    drop(vault);
    match result {
        Ok(()) => ExitCode::SUCCESS,
        Err(()) => {
            eprintln!("kcode-k1-daemon: listener failed");
            ExitCode::from(1)
        }
    }
}

async fn startup(k1_root: PathBuf, unlock: VaultUnlock) -> Result<Prepared, StartupError> {
    let public_origin = select_public_origin(std::env::var_os("K1_PUBLIC_ORIGIN"))
        .map_err(|error| with_cause("public origin")(error.to_string_lossy()))?;
    let state_root = k1_root.join("state");
    let files = DaemonFiles::open(&state_root).map_err(with_cause("daemon files"))?;
    let ordering = Arc::new(
        K1TxnOrdering::open(&state_root.join("ordering"))
            .map_err(with_cause("transaction ordering"))?,
    );
    let peering = Arc::new(
        K1Peering::open(&state_root.join("peering"), Arc::clone(&ordering))
            .map_err(with_cause("peering"))?,
    );
    macro_rules! open_kto_subsystem {
        ($component:ty, $directory:literal, $name:literal) => {
            Arc::new(
                <$component>::open(
                    &state_root.join($directory),
                    Arc::clone(&ordering),
                    Arc::clone(&peering),
                )
                .map_err(with_cause($name))?,
            )
        };
    }
    let web = open(&state_root, Arc::clone(&ordering), Arc::clone(&peering))
        .map_err(with_cause("Web HTTP"))?;
    let vault = unlock
        .open(&state_root, Arc::clone(&ordering), Arc::clone(&peering))
        .map_err(redacted_stage("Vault"))?;
    let persons = open_kto_subsystem!(K1Persons, "persons", "Persons");
    let invites = open_kto_subsystem!(K1Invites, "invites", "Invites");
    let accounts =
        Arc::new(K1Accounts::open(Arc::clone(&invites)).map_err(with_cause("Accounts"))?);
    let users = Arc::new(K1Users::new(Arc::clone(&accounts), Arc::clone(&persons)));
    let groups = open_kto_subsystem!(K1Groups, "groups", "Groups");
    let ese = Arc::new(
        K1Ese::open(
            &state_root.join("ese"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
            Arc::clone(&groups),
        )
        .map_err(with_cause("ESE"))?,
    );
    let launch_nodes = open_kto_subsystem!(LaunchNodes, "launch-nodes", "Launch Nodes");
    let profiles = open_kto_subsystem!(K1AccessProfiles, "access-profiles", "Access Profiles");
    let filters = open_kto_subsystem!(K1AuthorityFilters, "authority-filters", "Authority Filters");
    let bootstrap_state = Arc::new(
        K1BootstrapState::open(Arc::clone(&ordering), Arc::clone(&peering))
            .map_err(with_cause("Loom bootstrap state"))?,
    );
    let gemini_key = vault
        .secret(GEMINI_API_KEY)
        .map_err(redacted_stage("Gemini API key"))?
        .ok_or_else(|| redacted_stage("Gemini API key")(()))?;
    let gemini = Gemini31Pro::new(
        gemini_key.expose_secret().to_owned(),
        Accounting::new(),
        std::time::Duration::from_secs(30 * 60),
    )
    .map_err(redacted_stage("Gemini client"))?;
    let executable = codex_executable(std::env::var_os(CODEX_EXECUTABLE_ENV));
    let working_directory = std::env::current_dir()
        .map_err(with_cause("working directory"))?
        .to_string_lossy()
        .into_owned();
    let (audio_config, chat_config) = codex_configs(executable, working_directory);
    let audio_codex_adapter = CodexAdapter::open(audio_config)
        .await
        .map_err(StartupError::CodexAdapter)?;
    let chat_codex_adapter = audio_codex_adapter
        .with_config(chat_config)
        .map_err(StartupError::CodexAdapter)?;
    let web_search = WebSearchRunner::new(chat_codex_adapter.clone());
    let analyzer = Analyzer::from_codex_adapter(gemini, audio_codex_adapter);
    let objects = Arc::new(
        K1Objects::open(Arc::clone(&ordering), Arc::clone(&peering))
            .map_err(with_cause("Objects"))?,
    );
    let code_services = CodeServices::open(
        &state_root,
        Arc::clone(&ordering),
        Arc::clone(&peering),
        Arc::clone(&groups),
        Arc::clone(&objects),
        web.projection(),
    )
    .map_err(with_cause("Code services"))?;
    let web_bootstrap_importer = code_services.web_bootstrap_importer();
    let rust_projection = code_services.rust_projection();
    let (rust_code, web_code) = code_services.into_parts();
    let classification = AudioClassification::open(
        Arc::clone(&ordering),
        Arc::clone(&peering),
        Arc::clone(&objects),
        analyzer,
    )
    .map_err(StartupError::AudioClassification)?;
    let classification = ClassificationLifetime::new(Arc::new(classification));
    let ffmpeg =
        resolve_ffmpeg().map_err(|_| with_cause("FFmpeg")("executable not found on PATH"))?;
    let full_audio = Arc::new(
        K1FullAudio::open(ffmpeg, Arc::clone(&objects), classification.clone_value())
            .map_err(with_cause("Full Audio"))?,
    );
    let access = Arc::new(
        K1Access::open(
            &state_root.join("access"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
            Arc::clone(&groups),
        )
        .map_err(with_cause("Access"))?,
    );
    let classifiers = Arc::new(
        K1AudioClassifiers::open(
            state_root.join("audio-classifiers"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
        )
        .map_err(with_cause("Audio Classifiers"))?,
    );
    let access_classifiers = Arc::new(
        K1AccessAudioClassifiers::open(Arc::clone(&access), Arc::clone(&classifiers))
            .map_err(with_cause("Access Audio Classifiers"))?,
    );
    let access_launch_nodes = Arc::new(
        K1AccessLaunchNodes::open(
            Arc::clone(&access),
            Arc::clone(&groups),
            Arc::clone(&launch_nodes),
        )
        .map_err(with_cause("Access Launch Nodes"))?,
    );
    let people_models: Arc<[_]> = Arc::from(people_models());
    let model_names = people_models
        .iter()
        .map(|model| (model.id(), Some(model.name().to_owned())))
        .collect();
    let social = SocialKtools::new(
        Arc::clone(&users),
        Arc::clone(&groups),
        Arc::clone(&access_launch_nodes),
        model_names,
    );
    let chat = K1ChatService::open_with_social_and_set_launch_node_and_rust_code_and_web_code(
        &state_root.join("chat-v2"),
        &state_root.join("kmap"),
        Arc::clone(&ordering),
        Arc::clone(&peering),
        Arc::clone(&access),
        Arc::clone(&profiles),
        chat_codex_adapter,
        social,
        SetLaunchNodeKtool::new(Arc::clone(&access_launch_nodes)),
        rust_code,
        web_code,
        web_search,
    )
    .map_err(StartupError::Chat)?;
    let access_kmap = chat.access_kmap();
    let ensure_loom_bootstrap =
        |services| ensure_with_topology(&k1_root, &access, &web_bootstrap_importer, services);
    let bootstrap = ensure_loom_bootstrap(BootstrapServices {
        state: &bootstrap_state,
        invites: &invites,
        users: &users,
        groups: &groups,
        profiles: &profiles,
        access_kmap: &access_kmap,
        access_launch_nodes: &access_launch_nodes,
        launch_nodes: &launch_nodes,
        rust_projection: &rust_projection,
        model: chat_access_model(),
    })
    .map_err(with_cause("Loom bootstrap"))?;
    let access_persons = Arc::new(
        K1AccessPersons::open(Arc::clone(&access), Arc::clone(&persons))
            .map_err(with_cause("Access Persons"))?,
    );
    let audio = Arc::new(
        K1AccessFullAudio::open_with_classifier(
            Arc::clone(&access),
            full_audio,
            classification.clone_value(),
            access_classifiers,
            classifiers,
        )
        .map_err(with_cause("Access Full Audio"))?,
    );
    let prepared_http =
        kcode_k1_daemon_http_composition::prepare(kcode_k1_daemon_http_composition::Inputs {
            replay_epoch: files.replay_epoch_path().to_owned(),
            server_id: files.server_id().to_owned(),
            public_origin: public_origin.clone(),
            public: web.router(),
            launch_nodes: bootstrap.launch_nodes(),
            accounts: Arc::clone(&accounts),
            invites: Arc::clone(&invites),
            users: Arc::clone(&users),
            groups: Arc::clone(&groups),
            profiles,
            filters,
            access,
            access_persons,
            access_launch_nodes,
            access_kmap,
            audio,
            chat,
            ese: Arc::clone(&ese),
            people_models,
            chat_model: chat_access_model(),
            persons_model: persons_access_model(),
            audio_model: audio_access_model(),
        })
        .await?;
    let unused_invites = kcode_k1_daemon_invite_stock::reconcile(
        &invites,
        files.invite_links_path(),
        INVITE_LINK_URL,
    )
    .map_err(with_cause("invite stock"))?;
    if unused_invites < 100 {
        return Err(with_cause("minimum invite stock")(
            "fewer than 100 unused invites",
        ));
    }
    let boundary = prepared_http.bind().await?;
    Ok(Prepared {
        boundary,
        classification,
        public_origin,
        unused_invites,
        vault,
    })
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn public_operation_and_state_roots_are_fixed() {
        let _: fn(PathBuf) -> ExitCode = run;
        let root = PathBuf::from("/trusted/k1/state");
        assert_eq!(
            root.join("authority-filters"),
            PathBuf::from("/trusted/k1/state/authority-filters")
        );
        assert_eq!(
            root.join("launch-nodes"),
            PathBuf::from("/trusted/k1/state/launch-nodes")
        );
        let chat = root.join("chat-v2");
        assert_eq!(chat, PathBuf::from("/trusted/k1/state/chat-v2"));
        assert_ne!(chat, PathBuf::from("/trusted/k1/state/chat"));
    }

    #[test]
    fn fixed_provider_key_and_invite_url_remain_exact() {
        assert_eq!(GEMINI_API_KEY, "gemini-api-key");
        assert_eq!(
            INVITE_LINK_URL,
            "http://localhost:4321/lib/kcode-k1-ui/*/account.html"
        );
    }

    #[test]
    fn classification_shutdown_is_synchronous() {
        let _: fn(&AudioClassification) = AudioClassification::shutdown;
    }

    #[test]
    fn selected_composition_dependencies_are_current() {
        kcode_k1_daemon_lib_testkit::verify_manifest(include_str!("../Cargo.toml"));
        kcode_k1_daemon_lib_testkit::verify_source(include_str!("lib.rs"));
    }
}