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