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