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