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