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::{
20    Boundary, PUBLIC_ORIGIN, api_not_found, warn_if_slow, write_readiness,
21};
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, stage};
27use kcode_k1_daemon_vault_unlock::VaultUnlock;
28use kcode_k1_full_audio::K1FullAudio;
29use kcode_k1_groups::K1Groups;
30use kcode_k1_http::{Config as HttpConfig, K1Http};
31use kcode_k1_http_access_context::K1HttpAccessContext;
32use kcode_k1_http_accounts::K1HttpAccounts;
33use kcode_k1_http_people::K1HttpPeople;
34use kcode_k1_http_replay::{ReplayConfig, ReplayWindow};
35use kcode_k1_invites::K1Invites;
36use kcode_k1_launch_nodes::LaunchNodes;
37use kcode_k1_objects::K1Objects;
38use kcode_k1_peering::K1Peering;
39use kcode_k1_persons::K1Persons;
40use kcode_k1_txn_ordering::K1TxnOrdering;
41use kcode_k1_users::K1Users;
42use kcode_k1_vault::{ExposeSecret, K1Vault};
43use kcode_speaker_v3_analysis::Analyzer;
44use std::path::{Path, PathBuf};
45use std::process::ExitCode;
46use std::sync::Arc;
47use std::time::Instant;
48
49const INVITE_LINK_URL: &str = "http://localhost:4321/lib/kcode-k1-ui/*/account.html";
50const GEMINI_API_KEY: &str = "gemini-api-key";
51
52struct Prepared {
53    boundary: Boundary,
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(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 state_root = state_root(&k1_root);
112    let files = DaemonFiles::open(&state_root).map_err(stage("daemon files"))?;
113    let ordering = Arc::new(
114        K1TxnOrdering::open(&state_root.join("ordering")).map_err(stage("transaction ordering"))?,
115    );
116    let peering = Arc::new(
117        K1Peering::open(&state_root.join("peering"), Arc::clone(&ordering))
118            .map_err(stage("peering"))?,
119    );
120    let vault = unlock
121        .open(&state_root, Arc::clone(&ordering), Arc::clone(&peering))
122        .map_err(stage("Vault"))?;
123    let persons = Arc::new(
124        K1Persons::open(
125            &state_root.join("persons"),
126            Arc::clone(&ordering),
127            Arc::clone(&peering),
128        )
129        .map_err(stage("Persons"))?,
130    );
131    let invites = Arc::new(
132        K1Invites::open(
133            &state_root.join("invites"),
134            Arc::clone(&ordering),
135            Arc::clone(&peering),
136        )
137        .map_err(stage("Invites"))?,
138    );
139    let accounts = Arc::new(K1Accounts::open(Arc::clone(&invites)).map_err(stage("Accounts"))?);
140    let users = Arc::new(K1Users::new(Arc::clone(&accounts), Arc::clone(&persons)));
141    let groups = Arc::new(
142        K1Groups::open(
143            &state_root.join("groups"),
144            Arc::clone(&ordering),
145            Arc::clone(&peering),
146        )
147        .map_err(stage("Groups"))?,
148    );
149    let launch_nodes = Arc::new(
150        LaunchNodes::open(
151            &state_root.join("launch-nodes"),
152            Arc::clone(&ordering),
153            Arc::clone(&peering),
154        )
155        .map_err(stage("Launch Nodes"))?,
156    );
157    let profiles = Arc::new(
158        K1AccessProfiles::open(
159            &state_root.join("access-profiles"),
160            Arc::clone(&ordering),
161            Arc::clone(&peering),
162        )
163        .map_err(stage("Access Profiles"))?,
164    );
165    let filters = Arc::new(
166        K1AuthorityFilters::open(
167            &state_root.join("authority-filters"),
168            Arc::clone(&ordering),
169            Arc::clone(&peering),
170        )
171        .map_err(stage("Authority Filters"))?,
172    );
173    let gemini_key = vault
174        .secret(GEMINI_API_KEY)
175        .map_err(stage("Gemini API key"))?
176        .ok_or(StartupError::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(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(stage("working directory"))?
186        .to_string_lossy()
187        .into_owned();
188    let (audio_config, chat_config) = codex_configs(executable, working_directory);
189    let audio_codex_adapter = CodexAdapter::open(audio_config)
190        .await
191        .map_err(StartupError::CodexAdapter)?;
192    let chat_codex_adapter = audio_codex_adapter
193        .with_config(chat_config)
194        .map_err(StartupError::CodexAdapter)?;
195    let web_search = WebSearchRunner::new(chat_codex_adapter.clone());
196    let analyzer = Analyzer::from_codex_adapter(gemini, audio_codex_adapter);
197    let objects = Arc::new(
198        K1Objects::open(Arc::clone(&ordering), Arc::clone(&peering)).map_err(stage("Objects"))?,
199    );
200    let classification = Arc::new(
201        AudioClassification::open(
202            &state_root.join("audio-classification-v2"),
203            Arc::clone(&ordering),
204            Arc::clone(&peering),
205            Arc::clone(&objects),
206            analyzer,
207        )
208        .map_err(StartupError::AudioClassification)?,
209    );
210    let ffmpeg = resolve_ffmpeg().map_err(stage("FFmpeg"))?;
211    let full_audio = Arc::new(
212        K1FullAudio::open(ffmpeg, Arc::clone(&objects), Arc::clone(&classification))
213            .map_err(stage("Full Audio"))?,
214    );
215    let access = Arc::new(
216        K1Access::open(
217            &state_root.join("access"),
218            Arc::clone(&ordering),
219            Arc::clone(&peering),
220            Arc::clone(&groups),
221        )
222        .map_err(stage("Access"))?,
223    );
224    let classifiers = Arc::new(
225        K1AudioClassifiers::open(
226            state_root.join("audio-classifiers"),
227            Arc::clone(&ordering),
228            Arc::clone(&peering),
229        )
230        .map_err(stage("Audio Classifiers"))?,
231    );
232    let access_classifiers = Arc::new(
233        K1AccessAudioClassifiers::open(Arc::clone(&access), Arc::clone(&classifiers))
234            .map_err(stage("Access Audio Classifiers"))?,
235    );
236    let access_launch_nodes = Arc::new(
237        K1AccessLaunchNodes::open(Arc::clone(&access), Arc::clone(&groups), launch_nodes)
238            .map_err(stage("Access Launch Nodes"))?,
239    );
240    let chat = K1ChatService::open(
241        &chat_root(&state_root),
242        &state_root.join("kmap"),
243        Arc::clone(&ordering),
244        Arc::clone(&peering),
245        Arc::clone(&access),
246        Arc::clone(&profiles),
247        chat_codex_adapter,
248        web_search,
249    )
250    .map_err(StartupError::Chat)?;
251    let access_kmap = chat.access_kmap();
252    let access_persons = Arc::new(
253        K1AccessPersons::open(Arc::clone(&access), Arc::clone(&persons))
254            .map_err(stage("Access Persons"))?,
255    );
256    let audio = Arc::new(
257        K1AccessFullAudio::open_with_classifier(
258            Arc::clone(&access),
259            full_audio,
260            classification,
261            access_classifiers,
262            classifiers,
263        )
264        .map_err(stage("Access Full Audio"))?,
265    );
266    let chat_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
267    let presentation_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
268    let launch_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
269    let kmap_context = K1HttpAccessContext::new(Arc::clone(&filters), chat_access_model());
270    let persons_context = K1HttpAccessContext::new(Arc::clone(&filters), persons_access_model());
271    let audio_context = K1HttpAccessContext::new(Arc::clone(&filters), audio_access_model());
272    let replay = ReplayWindow::open(ReplayConfig {
273        epoch_file: files.replay_epoch_path().to_owned(),
274        max_nonces_per_epoch: usize::MAX,
275    })
276    .await
277    .map_err(stage("HTTP replay"))?;
278    let unused_invites = kcode_k1_daemon_invite_stock::reconcile(
279        &invites,
280        files.invite_links_path(),
281        INVITE_LINK_URL,
282    )
283    .map_err(stage("invite stock"))?;
284    if unused_invites < 100 {
285        return Err(StartupError::Stage("minimum invite stock"));
286    }
287    let adapter = K1HttpAccounts::new(
288        Arc::clone(&accounts),
289        Arc::clone(&invites),
290        Arc::clone(&users),
291    );
292    let people_models: Arc<[kcode_k1_http_people::LocalModel]> = Arc::from(people_models());
293    let people = K1HttpPeople::new_with_models(
294        accounts,
295        users,
296        groups,
297        Arc::clone(&profiles),
298        Arc::clone(&filters),
299        people_models,
300    )
301    .map_err(stage("People HTTP"))?;
302    let http = K1Http::new(
303        HttpConfig {
304            server_id: files.server_id().to_owned(),
305            public_origin: PUBLIC_ORIGIN.to_owned(),
306            max_body_bytes: usize::MAX,
307        },
308        replay,
309        adapter.identity_provider(),
310    )
311    .map_err(stage("K1 HTTP"))?;
312    let presentation_routes = kcode_k1_http_access_profile_presentation::authenticated_routes(
313        Arc::clone(&access),
314        Arc::clone(&profiles),
315        presentation_context,
316    );
317    let person_routes = kcode_k1_http_persons::authenticated_routes(
318        access_persons,
319        access,
320        Arc::clone(&profiles),
321        persons_context,
322    )
323    .map_err(stage("Persons HTTP"))?;
324    let launch_routes = kcode_k1_http_launch_nodes::router(
325        access_launch_nodes,
326        Arc::clone(&access_kmap),
327        Arc::clone(&profiles),
328        launch_context,
329    );
330    let kmap_routes =
331        kcode_k1_http_kmap::authenticated_routes(access_kmap, Arc::clone(&profiles), kmap_context);
332    let audio_routes = kcode_k1_http_audio::authenticated_routes(
333        Arc::clone(&audio),
334        profiles,
335        audio_context.clone(),
336    );
337    let authenticated = adapter
338        .authenticated_routes()
339        .merge(people.authenticated_routes())
340        .merge(audio_routes)
341        .merge(kcode_k1_http_audio_artifacts::authenticated_routes(
342            audio,
343            audio_context,
344        ))
345        .merge(person_routes)
346        .merge(launch_routes)
347        .merge(kmap_routes)
348        .merge(kcode_k1_http_chat::router(chat, chat_context))
349        .merge(presentation_routes)
350        .fallback(api_not_found);
351    let api = http.router(
352        adapter.registration_endpoint(),
353        kcode_k1_terms::endpoint(),
354        authenticated,
355    );
356    let boundary = Boundary::bind(api, files.server_id().to_owned())
357        .await
358        .map_err(stage("listener bind"))?;
359    Ok(Prepared {
360        boundary,
361        unused_invites,
362        vault,
363    })
364}
365
366fn state_root(k1_root: &Path) -> PathBuf {
367    k1_root.join("state")
368}
369
370fn chat_root(state_root: &Path) -> PathBuf {
371    state_root.join("chat-v2")
372}
373
374#[cfg(test)]
375mod tests {
376    use super::*;
377
378    #[test]
379    fn public_operation_and_state_roots_are_fixed() {
380        let _: fn(PathBuf) -> ExitCode = run;
381        let root = state_root(Path::new("/trusted/k1"));
382        assert_eq!(root, PathBuf::from("/trusted/k1/state"));
383        assert_eq!(
384            root.join("authority-filters"),
385            PathBuf::from("/trusted/k1/state/authority-filters")
386        );
387        assert_eq!(
388            root.join("launch-nodes"),
389            PathBuf::from("/trusted/k1/state/launch-nodes")
390        );
391        let chat = chat_root(&root);
392        assert_eq!(chat, PathBuf::from("/trusted/k1/state/chat-v2"));
393        assert_ne!(chat, PathBuf::from("/trusted/k1/state/chat"));
394    }
395
396    #[test]
397    fn fixed_provider_key_and_origins_remain_exact_and_distinct() {
398        assert_eq!(GEMINI_API_KEY, "gemini-api-key");
399        assert_eq!(
400            INVITE_LINK_URL,
401            "http://localhost:4321/lib/kcode-k1-ui/*/account.html"
402        );
403        assert_eq!(PUBLIC_ORIGIN, "http://localhost:4450");
404        assert_ne!(INVITE_LINK_URL, PUBLIC_ORIGIN);
405    }
406
407    #[test]
408    fn selected_composition_dependencies_are_current_and_compatible() {
409        kcode_k1_daemon_lib_testkit::verify_manifest(include_str!("../Cargo.toml"));
410    }
411}