#![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, text_inference_config,
};
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_text_inference::{CodexTextInference, TextInference};
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 text_config = text_inference_config(executable.clone(), working_directory.clone());
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 text_inference = Arc::new(
CodexTextInference::new(audio_codex_adapter.clone(), text_config)
.map_err(with_cause("Text inference"))?,
);
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),
text_inference: text_inference as Arc<dyn TextInference>,
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"));
}
}