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}