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