Skip to main content

kcode_k1_daemon_lib/
lib.rs

1#![doc = include_str!("../Documentation.md")]
2
3use axum::body::{Body, to_bytes};
4use axum::extract::Request;
5use axum::http::header::{CACHE_CONTROL, CONTENT_LENGTH, CONTENT_TYPE, HOST};
6use axum::http::{HeaderValue, StatusCode};
7use axum::middleware::{self, Next};
8use axum::response::Response;
9use axum::routing::get;
10use axum::{Json, Router};
11use kcode_codex_terra::CodexTerra;
12use kcode_gemini_3_1_pro::Gemini31Pro;
13use kcode_k1_access::K1Access;
14use kcode_k1_access_full_audio::K1AccessFullAudio;
15use kcode_k1_access_persons::K1AccessPersons;
16use kcode_k1_access_profiles::K1AccessProfiles;
17use kcode_k1_accounting::Accounting;
18use kcode_k1_accounts::K1Accounts;
19use kcode_k1_audio_classification::AudioClassification;
20use kcode_k1_daemon_files::DaemonFiles;
21use kcode_k1_full_audio::K1FullAudio;
22use kcode_k1_groups::{K1Groups, ModelId};
23use kcode_k1_http::{Config, K1Http};
24use kcode_k1_http_accounts::K1HttpAccounts;
25use kcode_k1_http_people::K1HttpPeople;
26use kcode_k1_http_replay::{ReplayConfig, ReplayWindow};
27use kcode_k1_invites::K1Invites;
28use kcode_k1_objects::K1Objects;
29use kcode_k1_peering::K1Peering;
30use kcode_k1_persons::K1Persons;
31use kcode_k1_txn_ordering::K1TxnOrdering;
32use kcode_k1_users::K1Users;
33use kcode_k1_vault::{ExposeSecret, K1Vault, SecretString};
34use kcode_speaker_v3_analysis::Analyzer;
35use serde::Serialize;
36use serde_json::Value;
37use std::io::Write as _;
38use std::path::{Path, PathBuf};
39use std::process::ExitCode;
40use std::sync::Arc;
41use std::time::{Duration, Instant};
42use tokio::net::TcpListener;
43use tokio::signal::unix::{Signal, SignalKind, signal};
44
45const LISTEN_ADDRESS: &str = "127.0.0.1:4450";
46const PUBLIC_ORIGIN: &str = "http://localhost:4450";
47const INVITE_LINK_URL: &str = "http://localhost:4321/lib/kcode-k1-ui/*/account.html";
48const AUTHORITY: &str = "localhost:4450";
49const STARTUP_BOUND: Duration = Duration::from_millis(100);
50const PROVIDER_OPERATION_TIMEOUT: Duration = Duration::from_secs(30 * 60);
51const GEMINI_API_KEY: &str = "gemini-api-key";
52const GEMINI_MODEL_BYTES: [u8; 32] = *b"gemini-3.1-pro-preview..........";
53const TERRA_MODEL_BYTES: [u8; 32] = *b"gpt-5.6-terra...................";
54const API_OPERATION: &str = "serve API request";
55
56#[derive(Clone, Serialize)]
57struct PublicConfig {
58    protocol: &'static str,
59    server_id: String,
60    public_origin: &'static str,
61}
62
63#[derive(Serialize)]
64struct Ready {
65    event: &'static str,
66    public_origin: &'static str,
67    unused_invites: usize,
68}
69
70struct Prepared {
71    app: Router,
72    listener: TcpListener,
73    signals: Signals,
74    unused_invites: usize,
75    vault: Arc<K1Vault>,
76}
77
78struct Signals {
79    interrupt: Signal,
80    terminate: Signal,
81}
82
83pub fn run(k1_root: PathBuf) -> ExitCode {
84    let runtime = match tokio::runtime::Builder::new_multi_thread()
85        .enable_all()
86        .build()
87    {
88        Ok(runtime) => runtime,
89        Err(_) => {
90            eprintln!("kcode-k1-daemon: startup failed");
91            return ExitCode::from(1);
92        }
93    };
94    let passphrase = match rpassword::prompt_password("Unlock K1 vault: ") {
95        Ok(passphrase) => match protect_passphrase(passphrase) {
96            Ok(passphrase) => passphrase,
97            Err(()) => {
98                eprintln!("kcode-k1-daemon: startup failed");
99                return ExitCode::from(1);
100            }
101        },
102        Err(_) => {
103            eprintln!("kcode-k1-daemon: startup failed");
104            return ExitCode::from(1);
105        }
106    };
107    runtime.block_on(run_async(k1_root, passphrase))
108}
109
110fn protect_passphrase(passphrase: String) -> Result<SecretString, ()> {
111    (!passphrase.is_empty())
112        .then(|| SecretString::from(passphrase))
113        .ok_or(())
114}
115
116async fn run_async(k1_root: PathBuf, passphrase: SecretString) -> ExitCode {
117    let started = Instant::now();
118    let prepared = match startup(k1_root, passphrase).await {
119        Ok(prepared) => prepared,
120        Err(()) => {
121            warn_if_slow(started.elapsed(), "error");
122            eprintln!("kcode-k1-daemon: startup failed");
123            return ExitCode::from(1);
124        }
125    };
126    let elapsed = started.elapsed();
127    if write_readiness(prepared.unused_invites).is_err() {
128        warn_if_slow(elapsed, "error");
129        eprintln!("kcode-k1-daemon: startup failed");
130        return ExitCode::from(1);
131    }
132    warn_if_slow(elapsed, "ready");
133    let Prepared {
134        app,
135        listener,
136        signals,
137        vault,
138        ..
139    } = prepared;
140    let result = axum::serve(listener, app)
141        .with_graceful_shutdown(signals.wait())
142        .await;
143    drop(vault);
144    match result {
145        Ok(()) => ExitCode::SUCCESS,
146        Err(_) => {
147            eprintln!("kcode-k1-daemon: listener failed");
148            ExitCode::from(1)
149        }
150    }
151}
152
153async fn startup(k1_root: PathBuf, passphrase: SecretString) -> Result<Prepared, ()> {
154    let state_root = state_root(&k1_root);
155    let files = DaemonFiles::open(&state_root).map_err(|_| ())?;
156    let ordering = Arc::new(K1TxnOrdering::open(&state_root.join("ordering")).map_err(|_| ())?);
157    let peering = Arc::new(
158        K1Peering::open(&state_root.join("peering"), Arc::clone(&ordering)).map_err(|_| ())?,
159    );
160    let vault = open_vault(
161        &state_root,
162        passphrase,
163        Arc::clone(&ordering),
164        Arc::clone(&peering),
165    )?;
166    let persons = Arc::new(
167        K1Persons::open(
168            &state_root.join("persons"),
169            Arc::clone(&ordering),
170            Arc::clone(&peering),
171        )
172        .map_err(|_| ())?,
173    );
174    let invites = Arc::new(
175        K1Invites::open(
176            &state_root.join("invites"),
177            Arc::clone(&ordering),
178            Arc::clone(&peering),
179        )
180        .map_err(|_| ())?,
181    );
182    let accounts = Arc::new(K1Accounts::open(Arc::clone(&invites)).map_err(|_| ())?);
183    let users = Arc::new(K1Users::new(Arc::clone(&accounts), Arc::clone(&persons)));
184    let groups = Arc::new(
185        K1Groups::open(
186            &state_root.join("groups"),
187            Arc::clone(&ordering),
188            Arc::clone(&peering),
189        )
190        .map_err(|_| ())?,
191    );
192    let profiles = Arc::new(
193        K1AccessProfiles::open(
194            &state_root.join("access-profiles"),
195            Arc::clone(&ordering),
196            Arc::clone(&peering),
197        )
198        .map_err(|_| ())?,
199    );
200    let gemini_key = vault.secret(GEMINI_API_KEY).map_err(|_| ())?.ok_or(())?;
201    let accounting = Accounting::new();
202    let gemini = Gemini31Pro::new(
203        gemini_key.expose_secret().to_owned(),
204        accounting.clone(),
205        PROVIDER_OPERATION_TIMEOUT,
206    )
207    .map_err(|_| ())?;
208    let terra = CodexTerra::new(
209        accounting,
210        "codex",
211        std::env::current_dir().map_err(|_| ())?,
212        PROVIDER_OPERATION_TIMEOUT,
213    )
214    .map_err(|_| ())?;
215    let analyzer = Analyzer::new(gemini, terra);
216    let objects =
217        Arc::new(K1Objects::open(Arc::clone(&ordering), Arc::clone(&peering)).map_err(|_| ())?);
218    let classification = Arc::new(
219        AudioClassification::open(
220            &state_root.join("audio-classification"),
221            Arc::clone(&ordering),
222            Arc::clone(&peering),
223            Arc::clone(&objects),
224            analyzer,
225        )
226        .map_err(|_| ())?,
227    );
228    let full_audio = Arc::new(
229        K1FullAudio::open(
230            resolve_ffmpeg()?,
231            Arc::clone(&objects),
232            Arc::clone(&classification),
233        )
234        .map_err(|_| ())?,
235    );
236    let access = Arc::new(
237        K1Access::open(
238            &state_root.join("access"),
239            Arc::clone(&ordering),
240            Arc::clone(&peering),
241            Arc::clone(&groups),
242        )
243        .map_err(|_| ())?,
244    );
245    let access_persons = Arc::new(
246        K1AccessPersons::open(
247            Arc::clone(&access),
248            Arc::clone(&profiles),
249            Arc::clone(&persons),
250        )
251        .map_err(|_| ())?,
252    );
253    let audio = Arc::new(
254        K1AccessFullAudio::open_for_models(
255            Arc::clone(&access),
256            Arc::clone(&profiles),
257            full_audio,
258            classification,
259            Arc::clone(&groups),
260            audio_models().to_vec(),
261        )
262        .map_err(|_| ())?,
263    );
264    let replay = ReplayWindow::open(ReplayConfig {
265        epoch_file: files.replay_epoch_path().to_owned(),
266        max_nonces_per_epoch: usize::MAX,
267    })
268    .await
269    .map_err(|_| ())?;
270    let unused_invites = kcode_k1_daemon_invite_stock::reconcile(
271        &invites,
272        files.invite_links_path(),
273        INVITE_LINK_URL,
274    )
275    .map_err(|_| ())?;
276    if unused_invites < 100 {
277        return Err(());
278    }
279    let adapter = K1HttpAccounts::new(
280        Arc::clone(&accounts),
281        Arc::clone(&invites),
282        Arc::clone(&users),
283    );
284    let people = K1HttpPeople::new(accounts, users, groups, profiles);
285    let http = K1Http::new(
286        Config {
287            server_id: files.server_id().to_owned(),
288            public_origin: PUBLIC_ORIGIN.to_owned(),
289            max_body_bytes: usize::MAX,
290        },
291        replay,
292        adapter.identity_provider(),
293    )
294    .map_err(|_| ())?;
295    let person_routes =
296        kcode_k1_http_persons::authenticated_routes(access_persons, access, audio_models()[0])
297            .map_err(|_| ())?;
298    let authenticated = adapter
299        .authenticated_routes()
300        .merge(people.authenticated_routes())
301        .merge(kcode_k1_http_audio::authenticated_routes(audio))
302        .merge(person_routes)
303        .fallback(api_not_found);
304    let api = http
305        .router(
306            adapter.registration_endpoint(),
307            kcode_k1_terms::endpoint(),
308            authenticated,
309        )
310        .layer(middleware::from_fn(contextualize_api_error));
311    let config = PublicConfig {
312        protocol: "K1-HTTP-1",
313        server_id: files.server_id().to_owned(),
314        public_origin: PUBLIC_ORIGIN,
315    };
316    let config_route = get(move || {
317        let config = config.clone();
318        async move { ([(CACHE_CONTROL, "no-store")], Json(config)) }
319    });
320    let app = Router::new()
321        .route("/config.json", config_route)
322        .merge(api)
323        .layer(middleware::from_fn(require_authority));
324    Ok(Prepared {
325        app,
326        listener: TcpListener::bind(LISTEN_ADDRESS).await.map_err(|_| ())?,
327        signals: Signals::install()?,
328        unused_invites,
329        vault,
330    })
331}
332
333fn audio_models() -> [ModelId; 2] {
334    [
335        ModelId::from_bytes(GEMINI_MODEL_BYTES),
336        ModelId::from_bytes(TERRA_MODEL_BYTES),
337    ]
338}
339
340fn resolve_ffmpeg() -> Result<PathBuf, ()> {
341    let path = std::env::var_os("PATH").ok_or(())?;
342    resolve_executable("ffmpeg", std::env::split_paths(&path))
343}
344
345fn resolve_executable(name: &str, paths: impl IntoIterator<Item = PathBuf>) -> Result<PathBuf, ()> {
346    paths
347        .into_iter()
348        .find_map(|directory| {
349            let candidate = directory.join(name);
350            executable(&candidate)
351                .then(|| std::fs::canonicalize(candidate).ok())
352                .flatten()
353                .filter(|path| path.is_absolute())
354        })
355        .ok_or(())
356}
357
358#[cfg(unix)]
359fn executable(path: &Path) -> bool {
360    use std::os::unix::fs::PermissionsExt as _;
361    std::fs::metadata(path)
362        .is_ok_and(|metadata| metadata.is_file() && metadata.permissions().mode() & 0o111 != 0)
363}
364
365#[cfg(not(unix))]
366fn executable(path: &Path) -> bool {
367    std::fs::metadata(path).is_ok_and(|metadata| metadata.is_file())
368}
369
370fn open_vault(
371    state_root: &Path,
372    passphrase: SecretString,
373    ordering: Arc<K1TxnOrdering>,
374    peering: Arc<K1Peering>,
375) -> Result<Arc<K1Vault>, ()> {
376    K1Vault::open(&state_root.join("vault"), passphrase, ordering, peering)
377        .map(Arc::new)
378        .map_err(|_| ())
379}
380
381fn state_root(k1_root: &Path) -> PathBuf {
382    k1_root.join("state")
383}
384
385async fn api_not_found() -> Response {
386    json_error(
387        StatusCode::NOT_FOUND,
388        "not_found",
389        "authenticated API route not found",
390    )
391}
392
393async fn contextualize_api_error(request: Request, next: Next) -> Response {
394    let response = next.run(request).await;
395    if !(response.status().is_client_error() || response.status().is_server_error()) {
396        return response;
397    }
398    let (mut parts, body) = response.into_parts();
399    let bytes = match to_bytes(body, usize::MAX).await {
400        Ok(bytes) => bytes,
401        Err(_) => return Response::from_parts(parts, Body::empty()),
402    };
403    let Some(contextualized) = contextualize_error_body(&bytes) else {
404        return Response::from_parts(parts, Body::from(bytes));
405    };
406    parts.headers.remove(CONTENT_LENGTH);
407    Response::from_parts(parts, Body::from(contextualized))
408}
409
410fn contextualize_error_body(bytes: &[u8]) -> Option<Vec<u8>> {
411    let mut payload: Value = serde_json::from_slice(bytes).ok()?;
412    let object = payload.as_object_mut()?;
413    let code = object.get("error")?.as_str()?.to_owned();
414    let source = object
415        .get("message")
416        .and_then(Value::as_str)
417        .map(str::to_owned)
418        .unwrap_or_else(|| format!("error code {code}"));
419    object.insert(
420        "message".to_owned(),
421        Value::String(format!("{API_OPERATION}: {source}")),
422    );
423    Some(payload.to_string().into_bytes())
424}
425
426async fn require_authority(request: Request, next: Next) -> Response {
427    let mut values = request.headers().get_all(HOST).iter();
428    if values
429        .next()
430        .is_some_and(|value| value.as_bytes() == AUTHORITY.as_bytes())
431        && values.next().is_none()
432    {
433        next.run(request).await
434    } else {
435        json_error(
436            StatusCode::MISDIRECTED_REQUEST,
437            "invalid_request_authority",
438            "validate request authority: request authority is invalid",
439        )
440    }
441}
442
443fn json_error(status: StatusCode, code: &'static str, message: &'static str) -> Response {
444    let mut response = Response::new(Body::from(
445        serde_json::json!({"error": code, "message": message}).to_string(),
446    ));
447    *response.status_mut() = status;
448    response
449        .headers_mut()
450        .insert(CONTENT_TYPE, HeaderValue::from_static("application/json"));
451    response
452        .headers_mut()
453        .insert(CACHE_CONTROL, HeaderValue::from_static("no-store"));
454    response
455}
456
457fn write_readiness(unused_invites: usize) -> Result<(), ()> {
458    let stdout = std::io::stdout();
459    let mut output = stdout.lock();
460    serde_json::to_writer(
461        &mut output,
462        &Ready {
463            event: "ready",
464            public_origin: PUBLIC_ORIGIN,
465            unused_invites,
466        },
467    )
468    .map_err(|_| ())?;
469    output.write_all(b"\n").map_err(|_| ())?;
470    output.flush().map_err(|_| ())
471}
472
473fn warn_if_slow(elapsed: Duration, outcome: &'static str) {
474    if elapsed > STARTUP_BOUND {
475        eprintln!(
476            "{{\"module\":\"kcode-k1-daemon\",\"operation\":\"startup\",\"elapsed_us\":{},\"outcome\":\"{outcome}\"}}",
477            elapsed.as_micros()
478        );
479    }
480}
481
482impl Signals {
483    fn install() -> Result<Self, ()> {
484        Ok(Self {
485            interrupt: signal(SignalKind::interrupt()).map_err(|_| ())?,
486            terminate: signal(SignalKind::terminate()).map_err(|_| ())?,
487        })
488    }
489
490    async fn wait(mut self) {
491        tokio::select! {
492            _ = self.interrupt.recv() => {}
493            _ = self.terminate.recv() => {}
494        }
495    }
496}
497
498#[cfg(test)]
499mod tests {
500    use super::*;
501
502    #[test]
503    fn public_operation_accepts_only_the_state_root() {
504        let _: fn(PathBuf) -> ExitCode = run;
505    }
506
507    #[test]
508    fn accepted_passphrase_boundary_is_strict_and_protected() {
509        assert!(protect_passphrase(String::new()).is_err());
510        let text = "conspicuous-fake-passphrase-never-real";
511        let protected = protect_passphrase(text.to_owned()).unwrap();
512        assert!(!format!("{protected:?}").contains(text));
513    }
514
515    #[test]
516    fn vault_composition_persists_at_the_fixed_path() {
517        let root =
518            std::env::temp_dir().join(format!("kcode-k1-daemon-vault-test-{}", std::process::id()));
519        let _ = std::fs::remove_dir_all(&root);
520        let state = state_root(&root);
521        assert_eq!(state.join("vault"), root.join("state/vault"));
522        let parts = || {
523            let ordering = Arc::new(K1TxnOrdering::open(&state.join("ordering")).unwrap());
524            let peering =
525                Arc::new(K1Peering::open(&state.join("peering"), ordering.clone()).unwrap());
526            (ordering, peering)
527        };
528        let password = || SecretString::from("fake-test-password-never-real");
529        let (ordering, peering) = parts();
530        let vault = open_vault(&state, password(), ordering.clone(), peering.clone()).unwrap();
531        vault
532            .set(
533                "fake-provider-secret",
534                SecretString::from("conspicuous-fake-value-never-real"),
535            )
536            .unwrap();
537        drop((vault, peering, ordering));
538        let (ordering, peering) = parts();
539        let vault = open_vault(&state, password(), ordering.clone(), peering.clone()).unwrap();
540        drop((vault, peering, ordering));
541        let (ordering, peering) = parts();
542        assert!(
543            open_vault(
544                &state,
545                SecretString::from("wrong-fake-password-never-real"),
546                ordering,
547                peering
548            )
549            .is_err()
550        );
551        std::fs::remove_dir_all(root).unwrap();
552    }
553
554    #[test]
555    fn audio_model_ids_are_fixed_distinct_and_in_order() {
556        assert_eq!(GEMINI_MODEL_BYTES, *b"gemini-3.1-pro-preview..........");
557        assert_eq!(TERRA_MODEL_BYTES, *b"gpt-5.6-terra...................");
558        assert_eq!(GEMINI_MODEL_BYTES.len(), 32);
559        assert_eq!(TERRA_MODEL_BYTES.len(), 32);
560        let models = audio_models();
561        assert_eq!(models[0].as_bytes(), &GEMINI_MODEL_BYTES);
562        assert_eq!(models[1].as_bytes(), &TERRA_MODEL_BYTES);
563        assert_ne!(models[0], models[1]);
564    }
565
566    #[test]
567    fn only_the_fixed_gemini_vault_key_is_selected() {
568        assert_eq!(GEMINI_API_KEY, "gemini-api-key");
569    }
570
571    #[cfg(unix)]
572    #[test]
573    fn executable_resolver_returns_an_absolute_path_for_a_fake_ffmpeg() {
574        use std::os::unix::fs::PermissionsExt as _;
575        let root = std::env::temp_dir().join(format!(
576            "kcode-k1-daemon-ffmpeg-test-{}",
577            std::process::id()
578        ));
579        let _ = std::fs::remove_dir_all(&root);
580        std::fs::create_dir(&root).unwrap();
581        let fake = root.join("ffmpeg");
582        std::fs::write(&fake, b"#!/bin/sh\nexit 0\n").unwrap();
583        std::fs::set_permissions(&fake, std::fs::Permissions::from_mode(0o700)).unwrap();
584        assert_eq!(
585            resolve_executable("ffmpeg", [root.clone()]).unwrap(),
586            std::fs::canonicalize(&fake).unwrap()
587        );
588        std::fs::remove_dir_all(root).unwrap();
589    }
590
591    #[test]
592    fn invite_link_and_backend_origins_remain_distinct() {
593        assert_eq!(
594            INVITE_LINK_URL,
595            "http://localhost:4321/lib/kcode-k1-ui/*/account.html"
596        );
597        assert_eq!(PUBLIC_ORIGIN, "http://localhost:4450");
598        assert_ne!(INVITE_LINK_URL, PUBLIC_ORIGIN);
599    }
600
601    #[test]
602    fn existing_child_message_is_preserved_under_daemon_context() {
603        let body = contextualize_error_body(
604            br#"{"error":"group_failed","message":"load group: child failure","detail":7}"#,
605        )
606        .unwrap();
607        let payload: Value = serde_json::from_slice(&body).unwrap();
608        assert_eq!(payload["error"], "group_failed");
609        assert_eq!(payload["detail"], 7);
610        assert_eq!(
611            payload["message"],
612            "serve API request: load group: child failure"
613        );
614    }
615
616    #[test]
617    fn missing_child_message_is_derived_from_stable_code() {
618        let body = contextualize_error_body(br#"{"error":"invalid_signature"}"#).unwrap();
619        let payload: Value = serde_json::from_slice(&body).unwrap();
620        assert_eq!(payload["error"], "invalid_signature");
621        assert_eq!(
622            payload["message"],
623            "serve API request: error code invalid_signature"
624        );
625    }
626
627    #[test]
628    fn supplied_root_maps_only_to_state() {
629        assert_eq!(
630            state_root(Path::new("/trusted/k1")),
631            PathBuf::from("/trusted/k1/state")
632        );
633    }
634}