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