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