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