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