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