kcode-k1-daemon-lib 0.12.2

Library-only private K1 loopback daemon 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_chat_service::K1ChatService;
use kcode_k1_codex_adapter::Adapter as CodexAdapter;
use kcode_k1_codex_websearch::Runner as WebSearchRunner;
use kcode_k1_daemon_files::DaemonFiles;
use kcode_k1_daemon_http_boundary::{
    Boundary, PUBLIC_ORIGIN, api_not_found, warn_if_slow, write_readiness,
};
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, stage};
use kcode_k1_daemon_vault_unlock::VaultUnlock;
use kcode_k1_full_audio::K1FullAudio;
use kcode_k1_groups::K1Groups;
use kcode_k1_http::{Config as HttpConfig, K1Http};
use kcode_k1_http_access_context::K1HttpAccessContext;
use kcode_k1_http_accounts::K1HttpAccounts;
use kcode_k1_http_people::K1HttpPeople;
use kcode_k1_http_replay::{ReplayConfig, ReplayWindow};
use kcode_k1_invites::K1Invites;
use kcode_k1_launch_nodes::LaunchNodes;
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::{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,
    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(_) => {
            eprintln!("kcode-k1-daemon: startup failed");
            return ExitCode::from(1);
        }
    };
    let unlock = match VaultUnlock::prompt() {
        Ok(unlock) => unlock,
        Err(_) => {
            eprintln!("kcode-k1-daemon: startup failed");
            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(prepared.unused_invites).is_err() {
        warn_if_slow(elapsed, "error");
        eprintln!("kcode-k1-daemon: startup failed");
        return ExitCode::from(1);
    }
    warn_if_slow(elapsed, "ready");
    let Prepared {
        boundary, vault, ..
    } = prepared;
    let result = boundary.serve().await;
    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 state_root = state_root(&k1_root);
    let files = DaemonFiles::open(&state_root).map_err(stage("daemon files"))?;
    let ordering = Arc::new(
        K1TxnOrdering::open(&state_root.join("ordering")).map_err(stage("transaction ordering"))?,
    );
    let peering = Arc::new(
        K1Peering::open(&state_root.join("peering"), Arc::clone(&ordering))
            .map_err(stage("peering"))?,
    );
    let vault = unlock
        .open(&state_root, Arc::clone(&ordering), Arc::clone(&peering))
        .map_err(stage("Vault"))?;
    let persons = Arc::new(
        K1Persons::open(
            &state_root.join("persons"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
        )
        .map_err(stage("Persons"))?,
    );
    let invites = Arc::new(
        K1Invites::open(
            &state_root.join("invites"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
        )
        .map_err(stage("Invites"))?,
    );
    let accounts = Arc::new(K1Accounts::open(Arc::clone(&invites)).map_err(stage("Accounts"))?);
    let users = Arc::new(K1Users::new(Arc::clone(&accounts), Arc::clone(&persons)));
    let groups = Arc::new(
        K1Groups::open(
            &state_root.join("groups"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
        )
        .map_err(stage("Groups"))?,
    );
    let launch_nodes = Arc::new(
        LaunchNodes::open(&state_root.join("launch-nodes")).map_err(stage("Launch Nodes"))?,
    );
    let profiles = Arc::new(
        K1AccessProfiles::open(
            &state_root.join("access-profiles"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
        )
        .map_err(stage("Access Profiles"))?,
    );
    let filters = Arc::new(
        K1AuthorityFilters::open(
            &state_root.join("authority-filters"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
        )
        .map_err(stage("Authority Filters"))?,
    );
    let gemini_key = vault
        .secret(GEMINI_API_KEY)
        .map_err(stage("Gemini API key"))?
        .ok_or(StartupError::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(stage("Gemini client"))?;
    let executable = codex_executable(std::env::var_os(CODEX_EXECUTABLE_ENV));
    let web_search = WebSearchRunner::new(executable.clone());
    let working_directory = std::env::current_dir()
        .map_err(stage("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 analyzer = Analyzer::from_codex_adapter(gemini, audio_codex_adapter);
    let objects = Arc::new(
        K1Objects::open(Arc::clone(&ordering), Arc::clone(&peering)).map_err(stage("Objects"))?,
    );
    let classification = Arc::new(
        AudioClassification::open(
            &state_root.join("audio-classification-v2"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
            Arc::clone(&objects),
            analyzer,
        )
        .map_err(StartupError::AudioClassification)?,
    );
    let ffmpeg = resolve_ffmpeg().map_err(stage("FFmpeg"))?;
    let full_audio = Arc::new(
        K1FullAudio::open(ffmpeg, Arc::clone(&objects), Arc::clone(&classification))
            .map_err(stage("Full Audio"))?,
    );
    let access = Arc::new(
        K1Access::open(
            &state_root.join("access"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
            Arc::clone(&groups),
        )
        .map_err(stage("Access"))?,
    );
    let classifiers = Arc::new(
        K1AudioClassifiers::open(
            state_root.join("audio-classifiers"),
            Arc::clone(&ordering),
            Arc::clone(&peering),
        )
        .map_err(stage("Audio Classifiers"))?,
    );
    let access_classifiers = Arc::new(
        K1AccessAudioClassifiers::open(Arc::clone(&access), Arc::clone(&classifiers))
            .map_err(stage("Access Audio Classifiers"))?,
    );
    let access_launch_nodes = Arc::new(
        K1AccessLaunchNodes::open(Arc::clone(&access), Arc::clone(&groups), launch_nodes)
            .map_err(stage("Access Launch Nodes"))?,
    );
    let chat = K1ChatService::open(
        &chat_root(&state_root),
        &state_root.join("kmap"),
        Arc::clone(&ordering),
        Arc::clone(&peering),
        Arc::clone(&access),
        Arc::clone(&profiles),
        chat_codex_adapter,
        web_search,
    )
    .map_err(StartupError::Chat)?;
    let access_kmap = chat.access_kmap();
    let access_persons = Arc::new(
        K1AccessPersons::open(Arc::clone(&access), Arc::clone(&persons))
            .map_err(stage("Access Persons"))?,
    );
    let audio = Arc::new(
        K1AccessFullAudio::open_with_classifier(
            Arc::clone(&access),
            full_audio,
            classification,
            access_classifiers,
            classifiers,
        )
        .map_err(stage("Access Full Audio"))?,
    );
    let chat_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
    let presentation_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
    let launch_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
    let persons_context = K1HttpAccessContext::new(Arc::clone(&filters), persons_access_model());
    let audio_context = K1HttpAccessContext::new(Arc::clone(&filters), audio_access_model());
    let replay = ReplayWindow::open(ReplayConfig {
        epoch_file: files.replay_epoch_path().to_owned(),
        max_nonces_per_epoch: usize::MAX,
    })
    .await
    .map_err(stage("HTTP replay"))?;
    let unused_invites = kcode_k1_daemon_invite_stock::reconcile(
        &invites,
        files.invite_links_path(),
        INVITE_LINK_URL,
    )
    .map_err(stage("invite stock"))?;
    if unused_invites < 100 {
        return Err(StartupError::Stage("minimum invite stock"));
    }
    let adapter = K1HttpAccounts::new(
        Arc::clone(&accounts),
        Arc::clone(&invites),
        Arc::clone(&users),
    );
    let people_models: Arc<[kcode_k1_http_people::LocalModel]> = Arc::from(people_models());
    let people = K1HttpPeople::new_with_models(
        accounts,
        users,
        groups,
        Arc::clone(&profiles),
        Arc::clone(&filters),
        people_models,
    )
    .map_err(stage("People HTTP"))?;
    let http = K1Http::new(
        HttpConfig {
            server_id: files.server_id().to_owned(),
            public_origin: PUBLIC_ORIGIN.to_owned(),
            max_body_bytes: usize::MAX,
        },
        replay,
        adapter.identity_provider(),
    )
    .map_err(stage("K1 HTTP"))?;
    let presentation_routes = kcode_k1_http_access_profile_presentation::authenticated_routes(
        Arc::clone(&access),
        Arc::clone(&profiles),
        presentation_context,
    );
    let person_routes = kcode_k1_http_persons::authenticated_routes(
        access_persons,
        access,
        Arc::clone(&profiles),
        persons_context,
    )
    .map_err(stage("Persons HTTP"))?;
    let launch_routes = kcode_k1_http_launch_nodes::router(
        access_launch_nodes,
        access_kmap,
        Arc::clone(&profiles),
        launch_context,
    );
    let audio_routes = kcode_k1_http_audio::authenticated_routes(
        Arc::clone(&audio),
        profiles,
        audio_context.clone(),
    );
    let authenticated = adapter
        .authenticated_routes()
        .merge(people.authenticated_routes())
        .merge(audio_routes)
        .merge(kcode_k1_http_audio_artifacts::authenticated_routes(
            audio,
            audio_context,
        ))
        .merge(person_routes)
        .merge(launch_routes)
        .merge(kcode_k1_http_chat::router(chat, chat_context))
        .merge(presentation_routes)
        .fallback(api_not_found);
    let api = http.router(
        adapter.registration_endpoint(),
        kcode_k1_terms::endpoint(),
        authenticated,
    );
    let boundary = Boundary::bind(api, files.server_id().to_owned())
        .await
        .map_err(stage("listener bind"))?;
    Ok(Prepared {
        boundary,
        unused_invites,
        vault,
    })
}

fn state_root(k1_root: &Path) -> PathBuf {
    k1_root.join("state")
}

fn chat_root(state_root: &Path) -> PathBuf {
    state_root.join("chat-v2")
}

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

    #[test]
    fn public_operation_and_state_roots_are_fixed() {
        let _: fn(PathBuf) -> ExitCode = run;
        let root = state_root(Path::new("/trusted/k1"));
        assert_eq!(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 = chat_root(&root);
        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_origins_remain_exact_and_distinct() {
        assert_eq!(GEMINI_API_KEY, "gemini-api-key");
        assert_eq!(
            INVITE_LINK_URL,
            "http://localhost:4321/lib/kcode-k1-ui/*/account.html"
        );
        assert_eq!(PUBLIC_ORIGIN, "http://localhost:4450");
        assert_ne!(INVITE_LINK_URL, PUBLIC_ORIGIN);
    }

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