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_audio_classifier::K1AccessAudioClassifiers;
6use kcode_k1_access_full_audio::K1AccessFullAudio;
7use kcode_k1_access_launch_nodes::K1AccessLaunchNodes;
8use kcode_k1_access_persons::K1AccessPersons;
9use kcode_k1_access_profiles::K1AccessProfiles;
10use kcode_k1_accounting::Accounting;
11use kcode_k1_accounts::K1Accounts;
12use kcode_k1_audio_classification::AudioClassification;
13use kcode_k1_audio_classifier::K1AudioClassifiers;
14use kcode_k1_authority_filters::K1AuthorityFilters;
15use kcode_k1_chat_service::K1ChatService;
16use kcode_k1_codex_adapter::Adapter as CodexAdapter;
17use kcode_k1_codex_websearch::Runner as WebSearchRunner;
18use kcode_k1_daemon_files::DaemonFiles;
19use kcode_k1_daemon_http_boundary::{Boundary, api_not_found, warn_if_slow, write_readiness_for};
20use kcode_k1_daemon_provider_config::{
21    CODEX_EXECUTABLE_ENV, audio_access_model, chat_access_model, codex_configs, codex_executable,
22    people_models, persons_access_model, resolve_ffmpeg,
23};
24use kcode_k1_daemon_startup_error::{StartupError, stage};
25use kcode_k1_daemon_vault_unlock::VaultUnlock;
26use kcode_k1_daemon_web_startup::{open_router, select_public_origin};
27use kcode_k1_full_audio::K1FullAudio;
28use kcode_k1_groups::K1Groups;
29use kcode_k1_http::{Config as HttpConfig, K1Http};
30use kcode_k1_http_access_context::K1HttpAccessContext;
31use kcode_k1_http_accounts::K1HttpAccounts;
32use kcode_k1_http_people::K1HttpPeople;
33use kcode_k1_http_replay::{ReplayConfig, ReplayWindow};
34use kcode_k1_invites::K1Invites;
35use kcode_k1_launch_nodes::LaunchNodes;
36use kcode_k1_objects::K1Objects;
37use kcode_k1_peering::K1Peering;
38use kcode_k1_persons::K1Persons;
39use kcode_k1_txn_ordering::K1TxnOrdering;
40use kcode_k1_users::K1Users;
41use kcode_k1_vault::{ExposeSecret, K1Vault};
42use kcode_speaker_v3_analysis::Analyzer;
43use std::path::{Path, PathBuf};
44use std::process::ExitCode;
45use std::sync::Arc;
46use std::time::Instant;
47
48const INVITE_LINK_URL: &str = "http://localhost:4321/lib/kcode-k1-ui/*/account.html";
49const GEMINI_API_KEY: &str = "gemini-api-key";
50
51struct ClassificationLifetime {
52    value: Arc<AudioClassification>,
53    stopped: bool,
54}
55
56impl ClassificationLifetime {
57    fn new(value: Arc<AudioClassification>) -> Self {
58        Self {
59            value,
60            stopped: false,
61        }
62    }
63
64    fn clone_value(&self) -> Arc<AudioClassification> {
65        Arc::clone(&self.value)
66    }
67
68    fn shutdown(&mut self) {
69        if !self.stopped {
70            self.value.shutdown();
71            self.stopped = true;
72        }
73    }
74}
75
76impl Drop for ClassificationLifetime {
77    fn drop(&mut self) {
78        self.shutdown();
79    }
80}
81
82struct Prepared {
83    boundary: Boundary,
84    classification: ClassificationLifetime,
85    public_origin: String,
86    unused_invites: usize,
87    vault: Arc<K1Vault>,
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_for(&prepared.public_origin, 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,
130        mut classification,
131        vault,
132        ..
133    } = prepared;
134    let result = boundary.serve().await;
135    classification.shutdown();
136    drop(vault);
137    match result {
138        Ok(()) => ExitCode::SUCCESS,
139        Err(()) => {
140            eprintln!("kcode-k1-daemon: listener failed");
141            ExitCode::from(1)
142        }
143    }
144}
145
146async fn startup(k1_root: PathBuf, unlock: VaultUnlock) -> Result<Prepared, StartupError> {
147    let public_origin = select_public_origin(std::env::var_os("K1_PUBLIC_ORIGIN"))
148        .map_err(|_| StartupError::Stage("public origin"))?;
149    let state_root = state_root(&k1_root);
150    let files = DaemonFiles::open(&state_root).map_err(stage("daemon files"))?;
151    let ordering = Arc::new(
152        K1TxnOrdering::open(&state_root.join("ordering")).map_err(stage("transaction ordering"))?,
153    );
154    let peering = Arc::new(
155        K1Peering::open(&state_root.join("peering"), Arc::clone(&ordering))
156            .map_err(stage("peering"))?,
157    );
158    macro_rules! open_kto_subsystem {
159        ($component:ty, $directory:literal, $name:literal) => {
160            Arc::new(
161                <$component>::open(
162                    &state_root.join($directory),
163                    Arc::clone(&ordering),
164                    Arc::clone(&peering),
165                )
166                .map_err(stage($name))?,
167            )
168        };
169    }
170    let web = open_router(&state_root, Arc::clone(&ordering), Arc::clone(&peering))
171        .map_err(stage("Web HTTP"))?;
172    let vault = unlock
173        .open(&state_root, Arc::clone(&ordering), Arc::clone(&peering))
174        .map_err(stage("Vault"))?;
175    let persons = open_kto_subsystem!(K1Persons, "persons", "Persons");
176    let invites = open_kto_subsystem!(K1Invites, "invites", "Invites");
177    let accounts = Arc::new(K1Accounts::open(Arc::clone(&invites)).map_err(stage("Accounts"))?);
178    let users = Arc::new(K1Users::new(Arc::clone(&accounts), Arc::clone(&persons)));
179    let groups = open_kto_subsystem!(K1Groups, "groups", "Groups");
180    let launch_nodes = open_kto_subsystem!(LaunchNodes, "launch-nodes", "Launch Nodes");
181    let profiles = open_kto_subsystem!(K1AccessProfiles, "access-profiles", "Access Profiles");
182    let filters = open_kto_subsystem!(K1AuthorityFilters, "authority-filters", "Authority Filters");
183    let gemini_key = vault
184        .secret(GEMINI_API_KEY)
185        .map_err(stage("Gemini API key"))?
186        .ok_or(StartupError::Stage("Gemini API key"))?;
187    let gemini = Gemini31Pro::new(
188        gemini_key.expose_secret().to_owned(),
189        Accounting::new(),
190        std::time::Duration::from_secs(30 * 60),
191    )
192    .map_err(stage("Gemini client"))?;
193    let executable = codex_executable(std::env::var_os(CODEX_EXECUTABLE_ENV));
194    let working_directory = std::env::current_dir()
195        .map_err(stage("working directory"))?
196        .to_string_lossy()
197        .into_owned();
198    let (audio_config, chat_config) = codex_configs(executable, working_directory);
199    let audio_codex_adapter = CodexAdapter::open(audio_config)
200        .await
201        .map_err(StartupError::CodexAdapter)?;
202    let chat_codex_adapter = audio_codex_adapter
203        .with_config(chat_config)
204        .map_err(StartupError::CodexAdapter)?;
205    let web_search = WebSearchRunner::new(chat_codex_adapter.clone());
206    let analyzer = Analyzer::from_codex_adapter(gemini, audio_codex_adapter);
207    let objects = Arc::new(
208        K1Objects::open(Arc::clone(&ordering), Arc::clone(&peering)).map_err(stage("Objects"))?,
209    );
210    let classification = AudioClassification::open(
211        Arc::clone(&ordering),
212        Arc::clone(&peering),
213        Arc::clone(&objects),
214        analyzer,
215    )
216    .map_err(StartupError::AudioClassification)?;
217    let classification = ClassificationLifetime::new(Arc::new(classification));
218    let ffmpeg = resolve_ffmpeg().map_err(stage("FFmpeg"))?;
219    let full_audio = Arc::new(
220        K1FullAudio::open(ffmpeg, Arc::clone(&objects), classification.clone_value())
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 classifiers = Arc::new(
233        K1AudioClassifiers::open(
234            state_root.join("audio-classifiers"),
235            Arc::clone(&ordering),
236            Arc::clone(&peering),
237        )
238        .map_err(stage("Audio Classifiers"))?,
239    );
240    let access_classifiers = Arc::new(
241        K1AccessAudioClassifiers::open(Arc::clone(&access), Arc::clone(&classifiers))
242            .map_err(stage("Access Audio Classifiers"))?,
243    );
244    let access_launch_nodes = Arc::new(
245        K1AccessLaunchNodes::open(Arc::clone(&access), Arc::clone(&groups), launch_nodes)
246            .map_err(stage("Access Launch Nodes"))?,
247    );
248    let chat = K1ChatService::open(
249        &chat_root(&state_root),
250        &state_root.join("kmap"),
251        Arc::clone(&ordering),
252        Arc::clone(&(peering)),
253        Arc::clone(&access),
254        Arc::clone(&profiles),
255        chat_codex_adapter,
256        web_search,
257    )
258    .map_err(StartupError::Chat)?;
259    let access_kmap = chat.access_kmap();
260    let access_persons = Arc::new(
261        K1AccessPersons::open(Arc::clone(&access), Arc::clone(&persons))
262            .map_err(stage("Access Persons"))?,
263    );
264    let audio = Arc::new(
265        K1AccessFullAudio::open_with_classifier(
266            Arc::clone(&access),
267            full_audio,
268            classification.clone_value(),
269            access_classifiers,
270            classifiers,
271        )
272        .map_err(stage("Access Full Audio"))?,
273    );
274    let chat_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
275    let presentation_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
276    let launch_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
277    let kmap_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
278    let persons_context = K1HttpAccessContext::new(Arc::clone(&filters), persons_access_model());
279    let audio_context = K1HttpAccessContext::new(Arc::clone(&filters), audio_access_model());
280    let replay = ReplayWindow::open(ReplayConfig {
281        epoch_file: files.replay_epoch_path().to_owned(),
282        max_nonces_per_epoch: usize::MAX,
283    })
284    .await
285    .map_err(stage("HTTP replay"))?;
286    let unused_invites = kcode_k1_daemon_invite_stock::reconcile(
287        &invites,
288        files.invite_links_path(),
289        INVITE_LINK_URL,
290    )
291    .map_err(stage("invite stock"))?;
292    if unused_invites < 100 {
293        return Err(StartupError::Stage("minimum invite stock"));
294    }
295    let adapter = K1HttpAccounts::new(
296        Arc::clone(&accounts),
297        Arc::clone(&invites),
298        Arc::clone(&users),
299    );
300    let people_models: Arc<[kcode_k1_http_people::LocalModel]> = Arc::from(people_models());
301    let people = K1HttpPeople::new_with_models(
302        accounts,
303        users,
304        groups,
305        Arc::clone(&profiles),
306        Arc::clone(&filters),
307        people_models,
308    )
309    .map_err(stage("People HTTP"))?;
310    let http = K1Http::new(
311        HttpConfig {
312            server_id: files.server_id().to_owned(),
313            public_origin: public_origin.clone(),
314            max_body_bytes: usize::MAX,
315        },
316        replay,
317        adapter.identity_provider(),
318    )
319    .map_err(stage("K1 HTTP"))?;
320    let presentation_routes = kcode_k1_http_access_profile_presentation::authenticated_routes(
321        Arc::clone(&access),
322        Arc::clone(&profiles),
323        presentation_context,
324    );
325    let person_routes = kcode_k1_http_persons::authenticated_routes(
326        access_persons,
327        access,
328        Arc::clone(&profiles),
329        persons_context,
330    )
331    .map_err(stage("Persons HTTP"))?;
332    let launch_routes = kcode_k1_http_launch_nodes::router(
333        access_launch_nodes,
334        Arc::clone(&access_kmap),
335        Arc::clone(&profiles),
336        launch_context,
337    );
338    let kmap_routes =
339        kcode_k1_http_kmap::authenticated_routes(access_kmap, Arc::clone(&profiles), kmap_context);
340    let audio_routes = kcode_k1_http_audio::authenticated_routes(audio, profiles, audio_context);
341    let authenticated = adapter
342        .authenticated_routes()
343        .merge(people.authenticated_routes())
344        .merge(audio_routes)
345        .merge(person_routes)
346        .merge(launch_routes)
347        .merge(kmap_routes)
348        .merge(kcode_k1_http_chat::router(chat, chat_context))
349        .merge(presentation_routes)
350        .fallback(api_not_found);
351    let api = http.router(
352        adapter.registration_endpoint(),
353        kcode_k1_terms::endpoint(),
354        authenticated,
355    );
356    let boundary = Boundary::bind_with_public(
357        api,
358        web,
359        files.server_id().to_owned(),
360        public_origin.clone(),
361    )
362    .await
363    .map_err(stage("listener bind"))?;
364    Ok(Prepared {
365        boundary,
366        classification,
367        public_origin,
368        unused_invites,
369        vault,
370    })
371}
372
373fn state_root(k1_root: &Path) -> PathBuf {
374    k1_root.join("state")
375}
376
377fn chat_root(state_root: &Path) -> PathBuf {
378    state_root.join("chat-v2")
379}
380
381#[cfg(test)]
382mod tests {
383    use super::*;
384
385    #[test]
386    fn public_operation_and_state_roots_are_fixed() {
387        let _: fn(PathBuf) -> ExitCode = run;
388        let root = state_root(Path::new("/trusted/k1"));
389        assert_eq!(root, PathBuf::from("/trusted/k1/state"));
390        assert_eq!(
391            root.join("authority-filters"),
392            PathBuf::from("/trusted/k1/state/authority-filters")
393        );
394        assert_eq!(
395            root.join("launch-nodes"),
396            PathBuf::from("/trusted/k1/state/launch-nodes")
397        );
398        let chat = chat_root(&root);
399        assert_eq!(chat, PathBuf::from("/trusted/k1/state/chat-v2"));
400        assert_ne!(chat, PathBuf::from("/trusted/k1/state/chat"));
401    }
402
403    #[test]
404    fn fixed_provider_key_and_invite_url_remain_exact() {
405        assert_eq!(GEMINI_API_KEY, "gemini-api-key");
406        assert_eq!(
407            INVITE_LINK_URL,
408            "http://localhost:4321/lib/kcode-k1-ui/*/account.html"
409        );
410    }
411
412    #[test]
413    fn classification_shutdown_is_synchronous() {
414        let _: fn(&AudioClassification) = AudioClassification::shutdown;
415    }
416
417    #[test]
418    fn selected_composition_dependencies_are_current() {
419        kcode_k1_daemon_lib_testkit::verify_manifest(include_str!("../Cargo.toml"));
420        kcode_k1_daemon_lib_testkit::verify_source(include_str!("lib.rs"));
421    }
422}