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