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