Skip to main content

kcode_k1_daemon_lib/
lib.rs

1#![doc = include_str!("../Documentation.md")]
2
3use kcode_gemini_3_1_pro::Gemini31Pro;
4use kcode_k1_access::K1Access;
5use kcode_k1_access_full_audio::K1AccessFullAudio;
6use kcode_k1_access_persons::K1AccessPersons;
7use kcode_k1_access_profiles::K1AccessProfiles;
8use kcode_k1_accounting::Accounting;
9use kcode_k1_accounts::K1Accounts;
10use kcode_k1_audio_classification::AudioClassification;
11use kcode_k1_chat_service::K1ChatService;
12use kcode_k1_codex_adapter::{Adapter as CodexAdapter, Error as CodexAdapterError};
13use kcode_k1_daemon_files::DaemonFiles;
14use kcode_k1_daemon_http_boundary::{
15    Boundary, PUBLIC_ORIGIN, api_not_found, warn_if_slow, write_readiness,
16};
17use kcode_k1_daemon_provider_config::{
18    CODEX_EXECUTABLE_ENV, audio_models, chat_model, codex_configs, codex_executable, people_models,
19    resolve_ffmpeg,
20};
21use kcode_k1_daemon_vault_unlock::VaultUnlock;
22use kcode_k1_full_audio::K1FullAudio;
23use kcode_k1_groups::K1Groups;
24use kcode_k1_http::{Config as HttpConfig, K1Http};
25use kcode_k1_http_accounts::K1HttpAccounts;
26use kcode_k1_http_people::K1HttpPeople;
27use kcode_k1_http_replay::{ReplayConfig, ReplayWindow};
28use kcode_k1_invites::K1Invites;
29use kcode_k1_objects::K1Objects;
30use kcode_k1_peering::K1Peering;
31use kcode_k1_persons::K1Persons;
32use kcode_k1_txn_ordering::K1TxnOrdering;
33use kcode_k1_users::K1Users;
34use kcode_k1_vault::{ExposeSecret, K1Vault};
35use kcode_speaker_v3_analysis::Analyzer;
36use std::fmt;
37use std::path::{Path, PathBuf};
38use std::process::ExitCode;
39use std::sync::Arc;
40use std::time::Instant;
41
42const INVITE_LINK_URL: &str = "http://localhost:4321/lib/kcode-k1-ui/*/account.html";
43const GEMINI_API_KEY: &str = "gemini-api-key";
44
45struct Prepared {
46    boundary: Boundary,
47    unused_invites: usize,
48    vault: Arc<K1Vault>,
49}
50
51enum StartupError {
52    Stage(&'static str),
53    CodexAdapter(CodexAdapterError),
54    Chat(String),
55}
56
57impl fmt::Display for StartupError {
58    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
59        match self {
60            Self::Stage(stage) => {
61                write!(formatter, "kcode-k1-daemon: startup failed at {stage}")
62            }
63            Self::CodexAdapter(error) => {
64                write!(
65                    formatter,
66                    "kcode-k1-daemon: startup failed at Codex adapter: {error}"
67                )
68            }
69            Self::Chat(child) => {
70                write!(
71                    formatter,
72                    "kcode-k1-daemon: startup failed at chat service: {child}"
73                )
74            }
75        }
76    }
77}
78
79fn stage<E>(name: &'static str) -> impl FnOnce(E) -> StartupError {
80    move |_| StartupError::Stage(name)
81}
82
83pub fn run(k1_root: PathBuf) -> ExitCode {
84    let runtime = match tokio::runtime::Builder::new_multi_thread()
85        .enable_all()
86        .build()
87    {
88        Ok(runtime) => runtime,
89        Err(_) => {
90            eprintln!("kcode-k1-daemon: startup failed");
91            return ExitCode::from(1);
92        }
93    };
94    let unlock = match VaultUnlock::prompt() {
95        Ok(unlock) => unlock,
96        Err(_) => {
97            eprintln!("kcode-k1-daemon: startup failed");
98            return ExitCode::from(1);
99        }
100    };
101    runtime.block_on(run_async(k1_root, unlock))
102}
103
104async fn run_async(k1_root: PathBuf, unlock: VaultUnlock) -> ExitCode {
105    let started = Instant::now();
106    let prepared = match startup(k1_root, unlock).await {
107        Ok(prepared) => prepared,
108        Err(error) => {
109            warn_if_slow(started.elapsed(), "error");
110            eprintln!("{error}");
111            return ExitCode::from(1);
112        }
113    };
114    let elapsed = started.elapsed();
115    if write_readiness(prepared.unused_invites).is_err() {
116        warn_if_slow(elapsed, "error");
117        eprintln!("kcode-k1-daemon: startup failed");
118        return ExitCode::from(1);
119    }
120    warn_if_slow(elapsed, "ready");
121    let Prepared {
122        boundary, vault, ..
123    } = prepared;
124    let result = boundary.serve().await;
125    drop(vault);
126    match result {
127        Ok(()) => ExitCode::SUCCESS,
128        Err(()) => {
129            eprintln!("kcode-k1-daemon: listener failed");
130            ExitCode::from(1)
131        }
132    }
133}
134
135async fn startup(k1_root: PathBuf, unlock: VaultUnlock) -> Result<Prepared, StartupError> {
136    let state_root = state_root(&k1_root);
137    let files = DaemonFiles::open(&state_root).map_err(stage("daemon files"))?;
138    let ordering = Arc::new(
139        K1TxnOrdering::open(&state_root.join("ordering")).map_err(stage("transaction ordering"))?,
140    );
141    let peering = Arc::new(
142        K1Peering::open(&state_root.join("peering"), Arc::clone(&ordering))
143            .map_err(stage("peering"))?,
144    );
145    let vault = unlock
146        .open(&state_root, Arc::clone(&ordering), Arc::clone(&peering))
147        .map_err(stage("Vault"))?;
148    let persons = Arc::new(
149        K1Persons::open(
150            &state_root.join("persons"),
151            Arc::clone(&ordering),
152            Arc::clone(&peering),
153        )
154        .map_err(stage("Persons"))?,
155    );
156    let invites = Arc::new(
157        K1Invites::open(
158            &state_root.join("invites"),
159            Arc::clone(&ordering),
160            Arc::clone(&peering),
161        )
162        .map_err(stage("Invites"))?,
163    );
164    let accounts = Arc::new(K1Accounts::open(Arc::clone(&invites)).map_err(stage("Accounts"))?);
165    let users = Arc::new(K1Users::new(Arc::clone(&accounts), Arc::clone(&persons)));
166    let groups = Arc::new(
167        K1Groups::open(
168            &state_root.join("groups"),
169            Arc::clone(&ordering),
170            Arc::clone(&peering),
171        )
172        .map_err(stage("Groups"))?,
173    );
174    let profiles = Arc::new(
175        K1AccessProfiles::open(
176            &state_root.join("access-profiles"),
177            Arc::clone(&ordering),
178            Arc::clone(&peering),
179        )
180        .map_err(stage("Access Profiles"))?,
181    );
182    let gemini_key = vault
183        .secret(GEMINI_API_KEY)
184        .map_err(stage("Gemini API key"))?
185        .ok_or(StartupError::Stage("Gemini API key"))?;
186    let gemini = Gemini31Pro::new(
187        gemini_key.expose_secret().to_owned(),
188        Accounting::new(),
189        std::time::Duration::from_secs(30 * 60),
190    )
191    .map_err(stage("Gemini client"))?;
192    let executable = codex_executable(std::env::var_os(CODEX_EXECUTABLE_ENV));
193    let working_directory = std::env::current_dir()
194        .map_err(stage("working directory"))?
195        .to_string_lossy()
196        .into_owned();
197    let (audio_config, chat_config) = codex_configs(executable, working_directory);
198    let audio_codex_adapter = CodexAdapter::open(audio_config)
199        .await
200        .map_err(StartupError::CodexAdapter)?;
201    let chat_codex_adapter = audio_codex_adapter
202        .with_config(chat_config)
203        .map_err(StartupError::CodexAdapter)?;
204    let analyzer = Analyzer::from_codex_adapter(gemini, audio_codex_adapter);
205    let objects = Arc::new(
206        K1Objects::open(Arc::clone(&ordering), Arc::clone(&peering)).map_err(stage("Objects"))?,
207    );
208    let classification = Arc::new(
209        AudioClassification::open(
210            &state_root.join("audio-classification"),
211            Arc::clone(&ordering),
212            Arc::clone(&peering),
213            Arc::clone(&objects),
214            analyzer,
215        )
216        .map_err(stage("Audio Classification"))?,
217    );
218    let ffmpeg = resolve_ffmpeg().map_err(stage("FFmpeg"))?;
219    let full_audio = Arc::new(
220        K1FullAudio::open(ffmpeg, Arc::clone(&objects), Arc::clone(&classification))
221            .map_err(stage("Full Audio"))?,
222    );
223    let access = Arc::new(
224        K1Access::open(
225            &state_root.join("access"),
226            Arc::clone(&ordering),
227            Arc::clone(&peering),
228            Arc::clone(&groups),
229        )
230        .map_err(stage("Access"))?,
231    );
232    let chat = K1ChatService::open(
233        &state_root.join("chat"),
234        Arc::clone(&ordering),
235        Arc::clone(&peering),
236        Arc::clone(&access),
237        Arc::clone(&profiles),
238        chat_codex_adapter,
239    )
240    .map_err(StartupError::Chat)?;
241    let access_persons = Arc::new(
242        K1AccessPersons::open(
243            Arc::clone(&access),
244            Arc::clone(&profiles),
245            Arc::clone(&persons),
246        )
247        .map_err(stage("Access Persons"))?,
248    );
249    let models = audio_models();
250    let audio = Arc::new(
251        K1AccessFullAudio::open_for_models(
252            Arc::clone(&access),
253            Arc::clone(&profiles),
254            full_audio,
255            classification,
256            Arc::clone(&groups),
257            models.to_vec(),
258        )
259        .map_err(stage("Access Full Audio"))?,
260    );
261    let replay = ReplayWindow::open(ReplayConfig {
262        epoch_file: files.replay_epoch_path().to_owned(),
263        max_nonces_per_epoch: usize::MAX,
264    })
265    .await
266    .map_err(stage("HTTP replay"))?;
267    let unused_invites = kcode_k1_daemon_invite_stock::reconcile(
268        &invites,
269        files.invite_links_path(),
270        INVITE_LINK_URL,
271    )
272    .map_err(stage("invite stock"))?;
273    if unused_invites < 100 {
274        return Err(StartupError::Stage("minimum invite stock"));
275    }
276    let adapter = K1HttpAccounts::new(
277        Arc::clone(&accounts),
278        Arc::clone(&invites),
279        Arc::clone(&users),
280    );
281    let people_models: Arc<[kcode_k1_http_people::LocalModel]> = Arc::from(people_models());
282    let people = K1HttpPeople::new_with_models(accounts, users, groups, profiles, people_models)
283        .map_err(stage("People HTTP"))?;
284    let http = K1Http::new(
285        HttpConfig {
286            server_id: files.server_id().to_owned(),
287            public_origin: PUBLIC_ORIGIN.to_owned(),
288            max_body_bytes: usize::MAX,
289        },
290        replay,
291        adapter.identity_provider(),
292    )
293    .map_err(stage("K1 HTTP"))?;
294    let person_routes =
295        kcode_k1_http_persons::authenticated_routes(access_persons, access, models[0])
296            .map_err(stage("Persons HTTP"))?;
297    let authenticated = adapter
298        .authenticated_routes()
299        .merge(people.authenticated_routes())
300        .merge(kcode_k1_http_audio::authenticated_routes(Arc::clone(
301            &audio,
302        )))
303        .merge(kcode_k1_http_audio_artifacts::authenticated_routes(audio))
304        .merge(person_routes)
305        .merge(kcode_k1_http_chat::router(chat, chat_model()))
306        .fallback(api_not_found);
307    let api = http.router(
308        adapter.registration_endpoint(),
309        kcode_k1_terms::endpoint(),
310        authenticated,
311    );
312    let boundary = Boundary::bind(api, files.server_id().to_owned())
313        .await
314        .map_err(stage("listener bind"))?;
315    Ok(Prepared {
316        boundary,
317        unused_invites,
318        vault,
319    })
320}
321
322fn state_root(k1_root: &Path) -> PathBuf {
323    k1_root.join("state")
324}
325
326#[cfg(test)]
327mod tests {
328    use super::*;
329
330    #[test]
331    fn public_operation_and_state_root_are_fixed() {
332        let _: fn(PathBuf) -> ExitCode = run;
333        assert_eq!(
334            state_root(Path::new("/trusted/k1")),
335            PathBuf::from("/trusted/k1/state")
336        );
337    }
338
339    #[test]
340    fn fixed_provider_key_and_origins_remain_exact_and_distinct() {
341        assert_eq!(GEMINI_API_KEY, "gemini-api-key");
342        assert_eq!(
343            INVITE_LINK_URL,
344            "http://localhost:4321/lib/kcode-k1-ui/*/account.html"
345        );
346        assert_eq!(PUBLIC_ORIGIN, "http://localhost:4450");
347        assert_ne!(INVITE_LINK_URL, PUBLIC_ORIGIN);
348    }
349
350    #[test]
351    fn startup_error_rendering_is_stage_specific_and_secret_safe() {
352        assert_eq!(
353            StartupError::Stage("Vault").to_string(),
354            "kcode-k1-daemon: startup failed at Vault"
355        );
356        let error = CodexAdapterError {
357            kind: kcode_k1_codex_adapter::ErrorKind::Unavailable,
358            message: "safe adapter display".to_owned(),
359            diagnostics: b"RAW_SECRET_DIAGNOSTIC".to_vec(),
360        };
361        let rendered = StartupError::CodexAdapter(error).to_string();
362        assert_eq!(
363            rendered,
364            "kcode-k1-daemon: startup failed at Codex adapter: safe adapter display"
365        );
366        assert!(!rendered.contains("RAW_SECRET_DIAGNOSTIC"));
367        let rendered =
368            StartupError::Chat("open chat service: safe child failure".to_owned()).to_string();
369        assert_eq!(
370            rendered,
371            "kcode-k1-daemon: startup failed at chat service: open chat service: safe child failure"
372        );
373    }
374
375    #[test]
376    fn selected_composition_dependencies_and_constructor_are_compatible() {
377        const MANIFEST: &str = include_str!("../Cargo.toml");
378        for selected in [
379            "kcode-k1-access-full-audio = \"0.7.4\"",
380            "kcode-k1-audio-classification = \"0.5.6\"",
381            "kcode-k1-chat-service = \"0.2.0\"",
382            "kcode-k1-codex-adapter = \"0.5.0\"",
383            "kcode-k1-daemon-http-boundary = \"0.1.0\"",
384            "kcode-k1-daemon-provider-config = \"0.1.2\"",
385            "kcode-k1-daemon-vault-unlock = \"0.1.0\"",
386            "kcode-k1-full-audio = \"0.3.6\"",
387            "kcode-k1-groups = \"0.3.1\"",
388            "kcode-k1-http-audio = \"0.1.4\"",
389            "kcode-k1-http-audio-artifacts = \"0.1.0\"",
390            "kcode-k1-http-chat = \"0.1.0\"",
391            "kcode-speaker-v3-analysis = { version = \"0.3.4\", default-features = false, features = [\"adapter-providers\"] }",
392        ] {
393            assert!(MANIFEST.contains(selected));
394        }
395        assert!(!MANIFEST.lines().any(|line| {
396            line.trim_start()
397                .starts_with("kcode-speaker-v3-terra-analysis ")
398        }));
399        fn require_constructor(_: fn(Gemini31Pro, CodexAdapter) -> Analyzer) {}
400        require_constructor(Analyzer::from_codex_adapter);
401        let models = people_models();
402        assert_eq!(
403            models
404                .iter()
405                .map(kcode_k1_http_people::LocalModel::name)
406                .collect::<Vec<_>>(),
407            [
408                "All models — special; includes current and future models",
409                "GPT-5.6 Terra",
410                "GPT-5.6 Sol",
411                "GPT-5.6 Luna",
412                "Gemini 3.1 Pro",
413            ]
414        );
415    }
416}