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