1#![forbid(unsafe_code)]
2
3use std::{
4 path::{Path, PathBuf},
5 str::FromStr,
6 sync::Arc,
7};
8
9use anyhow::Context;
10use clap::{Parser, Subcommand};
11use kcode_credential_vault::{CredentialVault, ExposeSecret, SecretString};
12use kcode_kweb_db::{Config as KwebConfig, NoopGossip, WriterId};
13use kcode_speaker_system::SpeechClassifier;
14use zeroize::{Zeroize, Zeroizing};
15
16const OPENAI_API_KEY_SECRET: &str = "openai-api-key";
17const GEMINI_API_KEY_SECRET: &str = "gemini-api-key";
18const TELEGRAM_BOT_TOKEN_SECRET: &str = "telegram-bot-token";
19const CRATES_IO_KEY_SECRET: &str = "cratesio-key";
20const KWEB_WRITER_SIGNING_KEY_SECRET: &str = "kweb-writer-signing-key";
21const KWEB_WRITERS_SECRET: &str = "kweb-writers-by-priority";
22const SPEECH_CLASSIFICATION_DATABASE_PATH: &str = "./data/kennedy-speech-classification.sqlite3";
23
24#[derive(Parser, Debug)]
25struct Args {
26 #[arg(long, global = true, default_value = "./data/kennedy-secrets.age")]
27 vault_path: PathBuf,
28 #[command(subcommand)]
29 command: Option<Command>,
30 #[arg(long, global = true, default_value = "127.0.0.1:4321")]
31 kweb_bind: String,
32 #[arg(long, global = true, default_value = "./data/kweb")]
33 kweb_root: PathBuf,
34 #[arg(
35 long,
36 global = true,
37 default_value = "./data/kennedy-conversations.sqlite3"
38 )]
39 conversation_history_database: PathBuf,
40 #[arg(long, global = true, default_value = "./data/sessions/in-progress")]
41 session_directory: PathBuf,
42 #[arg(long, global = true, default_value = "./data/session-history.txt")]
43 session_history_file: PathBuf,
44 #[arg(long, global = true, default_value = "./data/kennedy-telegram.sqlite3")]
45 telegram_database: PathBuf,
46 #[arg(long, global = true, default_value = "./data/kennedy-users.sqlite3")]
47 user_database: PathBuf,
48 #[arg(
49 long,
50 alias = "audio-ingress-database",
51 global = true,
52 default_value = "./data/kennedy-audio.sqlite3",
53 help = "Optional pre-library AudioIngress database used only for one-time migration"
54 )]
55 legacy_audio_ingress_database: PathBuf,
56 #[arg(
57 long,
58 alias = "audio-ingress-media",
59 global = true,
60 default_value = "./data/audio-ingress-media",
61 help = "AudioIngress-owned persistence root (database and original audio)"
62 )]
63 audio_ingress_directory: PathBuf,
64 #[arg(
65 long,
66 global = true,
67 default_value = "./data/intelligence-usage",
68 help = "One-file-per-call intelligence usage receipt directory"
69 )]
70 intelligence_usage_directory: PathBuf,
71 #[arg(long, default_value = "./data/kcode/kcode-rust-libs")]
72 rust_libs_root: PathBuf,
73 #[arg(long, default_value = "./data/kcode/kcode-web-libs")]
74 web_libs_root: PathBuf,
75 #[arg(long, default_value = "./data/kcode/kcode-web-libs-published")]
76 web_libs_published_root: PathBuf,
77 #[arg(long, default_value = "./data/kcode/kcode-rust-bins")]
78 rust_bins_root: PathBuf,
79 #[arg(long, default_value = "./data/kcode/kcode-rust-bin-artifacts")]
80 rust_bin_artifacts_root: PathBuf,
81 #[arg(long, default_value = "@taek42")]
82 telegram_bootstrap_username: String,
83 #[arg(long, default_value_t = 20 * 1024 * 1024)]
84 telegram_max_voice_bytes: usize,
85 #[arg(long, default_value_t = 8 * 1024 * 1024 * 1024)]
86 audio_ingress_max_upload_bytes: usize,
87}
88
89#[derive(Subcommand, Debug)]
90enum Command {
91 Secrets {
93 #[command(subcommand)]
94 command: SecretsCommand,
95 },
96 KmapSize,
98}
99
100#[derive(Subcommand, Debug)]
101enum SecretsCommand {
102 Set { name: String },
104 Remove { name: String },
106 List,
108 ChangePassphrase,
110}
111
112#[tokio::main]
113pub async fn main() -> anyhow::Result<()> {
114 tracing_subscriber::fmt()
115 .with_env_filter(
116 tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| {
117 "kennedy_server=info,kcode_kennedy_app=info,kcode_kennedy_orchestration=info,kcode_kennedy_telegram_runtime=info,kcode_kennedy_roots=info,kcode_kweb_db=info,kcode_codex_runtime=info,kcode_session_history=info,kcode_tg_kennedy_bot=info,tower_http=info".into()
118 }),
119 )
120 .init();
121 rustls::crypto::ring::default_provider()
122 .install_default()
123 .map_err(|_| anyhow::anyhow!("installing TLS crypto provider"))?;
124 let mut args = Args::parse();
125 let vault_path = args.vault_path.clone();
126 match args.command.take() {
127 Some(Command::Secrets { command }) => {
128 let _maintenance_guard = tokio::net::TcpListener::bind(&args.kweb_bind)
129 .await
130 .with_context(|| {
131 format!(
132 "binding maintenance lock {}; stop the running Kennedy server before changing its credential vault",
133 args.kweb_bind
134 )
135 })?;
136 manage_secrets(command, &vault_path)
137 }
138 Some(Command::KmapSize) => {
139 let _maintenance_guard =
140 maintenance_guard(&args.kweb_bind, "measuring the Kweb").await?;
141 let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
142 let vault = CredentialVault::unlock(&vault_path, passphrase)?;
143 let size = kcode_kmap_size::measure(&args.kweb_root, kweb_config(&vault)?)?;
144 println!("{}", kcode_kmap_size::render(&size));
145 Ok(())
146 }
147 None => run_server(args, vault_path).await,
148 }
149}
150
151async fn run_server(args: Args, vault_path: PathBuf) -> anyhow::Result<()> {
152 let kweb_listener = tokio::net::TcpListener::bind(&args.kweb_bind)
155 .await
156 .with_context(|| format!("binding Kweb listener {}", args.kweb_bind))?;
157 ensure_runtime_parent_directories(&args, &vault_path)?;
158 let vault = if vault_path.exists() {
159 let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
160 CredentialVault::unlock(&vault_path, passphrase)?
161 } else {
162 tracing::warn!(path=%vault_path.display(), "Kennedy credential vault does not exist; secret-backed features are unavailable");
163 CredentialVault::empty()
164 };
165 let openai_api_key = resolve_optional_secret(
166 &vault,
167 OPENAI_API_KEY_SECRET,
168 "OpenAI transcription, media annotation, agents, and image generation/editing",
169 )?;
170 let gemini_api_key = resolve_optional_secret(
171 &vault,
172 GEMINI_API_KEY_SECRET,
173 "Gemini search, media annotation, agents, audio transcription, and image generation/editing",
174 )?;
175 let telegram_bot_token =
176 resolve_optional_secret(&vault, TELEGRAM_BOT_TOKEN_SECRET, "Telegram relay")?
177 .map(kcode_tg_kennedy_bot::BotToken::new)
178 .transpose()?;
179 let crates_io_key =
180 resolve_required_secret(&vault, CRATES_IO_KEY_SECRET, "Rust library publication")?;
181 let kweb_config = kweb_config(&vault)?;
182 let codex_catalog_cache =
183 kcode_codex_runtime::CatalogCache::new(kcode_codex_runtime::DEFAULT_CODEX_EXECUTABLE);
184 let (kmap, system_roots) =
185 kcode_kennedy_roots::open(&args.kweb_root, kweb_config, &args.user_database)?;
186 let speech_classifier = SpeechClassifier::open(SPEECH_CLASSIFICATION_DATABASE_PATH)
187 .with_context(|| {
188 format!("opening speaker-classification database {SPEECH_CLASSIFICATION_DATABASE_PATH}")
189 })?;
190 let speech_classifier = Arc::new(speech_classifier);
191 let dev_tools = kcode_dev_tools::Service::open(kcode_dev_tools::Config {
192 rust_libraries_root: args.rust_libs_root.clone(),
193 web_libraries_root: args.web_libs_root.clone(),
194 web_publications_root: args.web_libs_published_root.clone(),
195 rust_binaries_root: args.rust_bins_root.clone(),
196 rust_binary_publications_root: args.rust_bin_artifacts_root.clone(),
197 crates_io_registry_token: crates_io_key,
198 })
199 .map_err(anyhow::Error::new)
200 .with_context(|| {
201 format!(
202 "opening managed Kcode development roots under {}",
203 args.rust_libs_root
204 .parent()
205 .unwrap_or(Path::new("."))
206 .display()
207 )
208 })?;
209 let web_publications_root = dev_tools.web_publications_root().to_path_buf();
210 let telegram_identity = std::sync::Arc::new(kcode_telegram_identity::Directory::open(
211 &args.user_database,
212 &args.telegram_bootstrap_username,
213 )?);
214 let history_service =
215 kcode_session_history::SessionHistory::open(kcode_session_history::Config {
216 directory: args.session_directory,
217 completed_list: args.session_history_file,
218 provider_cost_compatibility: Some(
219 kcode_intelligence_chatend::provider_cost_compatibility(),
220 ),
221 })?;
222 let (intelligence_service, intelligence_runtime) =
223 kcode_intelligence_router::open(kcode_intelligence_router::Config {
224 openai_api_key,
225 gemini_api_key,
226 codex_catalog_cache,
227 receipt_directory: args.intelligence_usage_directory,
228 })
229 .await?;
230 let agent_runtime = kcode_agent_runtime::AgentRuntime::new(intelligence_service.clone());
231 let telegram_runtime = kcode_tg_kennedy_bot::open(kcode_tg_kennedy_bot::Config {
232 database: args.telegram_database,
233 bot_token: telegram_bot_token,
234 identity_sink: telegram_identity.clone(),
235 max_voice_bytes: args.telegram_max_voice_bytes,
236 })
237 .await?;
238 let telegram_service = telegram_runtime.service();
239 let chunk_intelligence = intelligence_service.clone();
240 let transcribe_chunk: kcode_audio_ingress::AudioChunkCall = Arc::new(move |request| {
241 let intelligence = chunk_intelligence.clone();
242 Box::pin(async move {
243 let user = intelligence
244 .for_user(request.user_id)
245 .map_err(audio_intelligence_error)?;
246 let media = kcode_intelligence_router::Media::audio(
247 request.audio_ogg,
248 "audio-chunk.ogg",
249 "audio/ogg",
250 )
251 .map_err(audio_intelligence_error)?;
252 user.analyze_audio(kcode_intelligence_router::AudioAnalysisRequest {
253 operation: "transcribe_chunk".into(),
254 prompt: request.prompt,
255 model: request.model,
256 media,
257 schema: request.schema,
258 max_output_tokens: request.max_output_tokens,
259 temperature: None,
260 operation_id: uuid::Uuid::new_v4(),
261 parent_operation_id: None,
262 })
263 .await
264 .map(|response| response.value.text)
265 .map_err(audio_intelligence_error)
266 })
267 });
268 let text_intelligence = intelligence_service.clone();
269 let generate_text: kcode_audio_ingress::TextGenerationCall = Arc::new(move |request| {
270 let intelligence = text_intelligence.clone();
271 Box::pin(async move {
272 let reasoning_effort = match request.reasoning_effort.as_str() {
273 "xhigh" => kcode_intelligence_router::ReasoningEffort::XHigh,
274 _ => {
275 return Err(kcode_audio_ingress::IntelligenceError::new(
276 "AudioIngress requested an unsupported reasoning effort.",
277 false,
278 ));
279 }
280 };
281 let user = intelligence
282 .for_user(request.user_id)
283 .map_err(audio_intelligence_error)?;
284 user.generate_text(kcode_intelligence_router::TextGenerationRequest {
285 operation: request.operation,
286 prompt: request.prompt,
287 model: request.model,
288 reasoning_effort,
289 timeout: request.timeout,
290 operation_id: uuid::Uuid::new_v4(),
291 parent_operation_id: None,
292 })
293 .await
294 .map(|response| response.value.text)
295 .map_err(audio_intelligence_error)
296 })
297 });
298 let audio_transcriber =
299 kcode_audio_ingress::AudioTranscriber::new(transcribe_chunk, generate_text);
300 let audio_state_database = args.audio_ingress_directory.join("state.sqlite3");
301 migrate_audio_ingress_database(&args.legacy_audio_ingress_database, &audio_state_database)?;
302 let audio = kcode_audio_ingress::AudioIngress::open(
303 &args.audio_ingress_directory,
304 audio_transcriber,
305 Arc::clone(&speech_classifier),
306 )
307 .await?;
308 let audio_coordinator = kcode_audio_session_ingress::Coordinator::new(
309 audio,
310 history_service.clone(),
311 kcode_audio_session_ingress::Config {
312 user_id: system_roots.user.to_string(),
313 effective_context_tokens: intelligence_runtime.context_window_tokens,
314 },
315 )?;
316 let http_router = kcode_http_api::router(kcode_http_api::Config {
317 kmap: kmap.clone(),
318 user_root_node_id: system_roots.user,
319 kennedy_root_node_id: system_roots.kennedy,
320 telegram: telegram_service.clone(),
321 session_history: history_service.clone(),
322 audio_ingress: audio_coordinator.clone(),
323 audio_max_upload_bytes: args.audio_ingress_max_upload_bytes,
324 web_publications_root,
325 })?;
326 let orchestration_config = kcode_kennedy_orchestration::Config {
327 user_root_node_id: system_roots.user.to_string(),
328 kennedy_root_node_id: system_roots.kennedy.to_string(),
329 telegram_max_media_bytes: args.telegram_max_voice_bytes,
330 runtime_model: kcode_kennedy_orchestration::RuntimeModel::from_intelligence(
331 intelligence_runtime,
332 ),
333 };
334 let telegram_sessions = kcode_telegram_session_coordinator::Service::new(
335 telegram_service.clone(),
336 telegram_identity.clone(),
337 );
338 let session_service =
339 kcode_kennedy_sessions::Service::new(kcode_kennedy_sessions::Capabilities {
340 kmap: kmap.clone(),
341 intelligence: intelligence_service.clone(),
342 agents: agent_runtime,
343 history: history_service.clone(),
344 speech_classifier,
345 dev_tools: dev_tools.clone(),
346 telegram: telegram_sessions,
347 });
348 let orchestration_api = kcode_kennedy_orchestration::Api::new(
349 &orchestration_config,
350 kcode_kennedy_orchestration::LocalServices {
351 kmap: kmap.clone(),
352 intelligence: intelligence_service,
353 history: history_service.clone(),
354 audio: audio_coordinator,
355 directory: telegram_identity.clone(),
356 dev_tools,
357 telegram: telegram_service,
358 },
359 );
360 let orchestration_worker = kcode_kennedy_orchestration::build(
361 orchestration_config,
362 orchestration_api,
363 session_service,
364 );
365 let directory_roots = kcode_kennedy_roots::DirectoryRoots::new(
366 kmap,
367 telegram_identity,
368 args.telegram_bootstrap_username.clone(),
369 system_roots.user,
370 orchestration_worker.writer().clone(),
371 );
372 let telegram_session_runtime = Arc::new(kcode_kennedy_telegram_runtime::Runtime::new(
373 kcode_kennedy_telegram_runtime::Config {
374 telegram_max_media_bytes: args.telegram_max_voice_bytes,
375 telegram_web_user_handle: args.telegram_bootstrap_username,
376 },
377 orchestration_worker.clone(),
378 directory_roots,
379 ));
380 tokio::try_join!(
381 serve_http(kweb_listener, http_router),
382 telegram_runtime.run(),
383 kcode_kennedy_orchestration::run(orchestration_worker),
384 telegram_session_runtime.run(),
385 )?;
386 Ok(())
387}
388
389async fn serve_http(listener: tokio::net::TcpListener, router: axum::Router) -> anyhow::Result<()> {
390 tracing::info!(address=%listener.local_addr()?, "Kennedy main HTTP server ready");
391 axum::serve(listener, router).await?;
392 Ok(())
393}
394
395fn audio_intelligence_error(
396 error: kcode_intelligence_router::Error,
397) -> kcode_audio_ingress::IntelligenceError {
398 let retryable = error.retryable();
399 kcode_audio_ingress::IntelligenceError::new(error.message(), retryable)
400}
401
402fn ensure_runtime_parent_directories(args: &Args, vault_path: &Path) -> anyhow::Result<()> {
403 for path in [
404 vault_path,
405 &args.kweb_root,
406 &args.conversation_history_database,
407 &args.session_directory,
408 &args.session_history_file,
409 &args.telegram_database,
410 &args.user_database,
411 Path::new(SPEECH_CLASSIFICATION_DATABASE_PATH),
412 &args.legacy_audio_ingress_database,
413 &args.audio_ingress_directory,
414 &args.intelligence_usage_directory,
415 &args.rust_libs_root,
416 &args.web_libs_root,
417 &args.web_libs_published_root,
418 &args.rust_bins_root,
419 &args.rust_bin_artifacts_root,
420 ] {
421 let Some(parent) = path.parent().filter(|value| !value.as_os_str().is_empty()) else {
422 continue;
423 };
424 if parent.exists() {
425 continue;
426 }
427 let mut builder = std::fs::DirBuilder::new();
428 builder.recursive(true);
429 #[cfg(unix)]
430 {
431 use std::os::unix::fs::DirBuilderExt;
432 builder.mode(0o700);
433 }
434 builder
435 .create(parent)
436 .with_context(|| format!("creating runtime data directory {}", parent.display()))?;
437 }
438 Ok(())
439}
440
441fn migrate_audio_ingress_database(legacy: &Path, current: &Path) -> anyhow::Result<()> {
442 if current.exists() || !legacy.exists() {
443 return Ok(());
444 }
445 if let Some(parent) = current.parent() {
446 std::fs::create_dir_all(parent)
447 .with_context(|| format!("creating AudioIngress root {}", parent.display()))?;
448 }
449 let source = rusqlite::Connection::open(legacy)
450 .with_context(|| format!("opening legacy AudioIngress database {}", legacy.display()))?;
451 source
452 .execute_batch("PRAGMA wal_checkpoint(TRUNCATE);")
453 .context("checkpointing legacy AudioIngress database")?;
454 source
455 .backup(rusqlite::MAIN_DB, current, None)
456 .context("copying legacy AudioIngress database into its persistence root")?;
457 let destination = rusqlite::Connection::open(current)
458 .with_context(|| format!("opening AudioIngress database {}", current.display()))?;
459 destination
460 .execute_batch("PRAGMA wal_checkpoint(TRUNCATE);")
461 .context("syncing migrated AudioIngress database")?;
462 tracing::info!(
463 source = %legacy.display(),
464 destination = %current.display(),
465 "Migrated AudioIngress database into its owned persistence root"
466 );
467 Ok(())
468}
469
470pub async fn maintenance_guard(
471 bind: &str,
472 purpose: &str,
473) -> anyhow::Result<tokio::net::TcpListener> {
474 tokio::net::TcpListener::bind(bind).await.with_context(|| {
475 format!("binding maintenance lock {bind}; stop the running Kennedy server before {purpose}")
476 })
477}
478
479fn kweb_config(vault: &CredentialVault) -> anyhow::Result<KwebConfig> {
480 let encoded_key = resolve_required_secret(
481 vault,
482 KWEB_WRITER_SIGNING_KEY_SECRET,
483 "Kweb mutation signing",
484 )?;
485 let mut signing_key = Zeroizing::new([0_u8; 32]);
486 let decoded = hex::decode(encoded_key.trim())
487 .context("Kweb writer signing key must be 64 lowercase hexadecimal characters")?;
488 *signing_key = decoded
489 .try_into()
490 .map_err(|_| anyhow::anyhow!("Kweb writer signing key must decode to exactly 32 bytes"))?;
491 let encoded_writers =
492 resolve_required_secret(vault, KWEB_WRITERS_SECRET, "Kweb writer authorization")?;
493 let writers_by_priority = encoded_writers
494 .split(',')
495 .map(str::trim)
496 .filter(|value| !value.is_empty())
497 .map(WriterId::from_str)
498 .collect::<Result<Vec<_>, _>>()
499 .map_err(anyhow::Error::new)
500 .context("decoding the ordered Kweb writer whitelist")?;
501 anyhow::ensure!(
502 !writers_by_priority.is_empty(),
503 "the Kweb writer whitelist is empty"
504 );
505 Ok(KwebConfig {
506 signing_key: *signing_key,
507 writers_by_priority,
508 gossip: Arc::new(NoopGossip),
509 })
510}
511
512fn resolve_optional_secret(
513 vault: &CredentialVault,
514 configured_name: &str,
515 purpose: &str,
516) -> anyhow::Result<Option<String>> {
517 let name = configured_name.trim();
518 if name.is_empty() {
519 return Ok(None);
520 }
521 let secret = vault.secret(name)?;
522 if secret.is_none() {
523 tracing::warn!(secret_name=name, %purpose, "configured Kennedy secret is not present in the vault");
524 }
525 Ok(secret.map(|value| value.expose_secret().to_owned()))
526}
527
528fn resolve_required_secret(
529 vault: &CredentialVault,
530 configured_name: &str,
531 purpose: &str,
532) -> anyhow::Result<String> {
533 let name = configured_name.trim();
534 if name.is_empty() {
535 anyhow::bail!("{purpose} requires a configured Kennedy secret name");
536 }
537 vault
538 .secret(name)?
539 .map(|value| value.expose_secret().to_owned())
540 .with_context(|| {
541 format!(
542 "{purpose} requires Kennedy secret '{name}'; store it with `kennedy-server secrets set {name}`"
543 )
544 })
545}
546
547fn manage_secrets(command: SecretsCommand, vault_path: &Path) -> anyhow::Result<()> {
548 match command {
549 SecretsCommand::Set { name } => {
550 let (mut vault, passphrase) = unlock_for_edit(vault_path)?;
551 let value = prompt_confirmed_value(&format!("Value for {name}: "))?;
552 vault.set(&name, value)?;
553 vault.save(vault_path, &passphrase)?;
554 println!("Stored Kennedy secret '{name}'.");
555 }
556 SecretsCommand::Remove { name } => {
557 if !vault_path.exists() {
558 println!("No Kennedy credential vault exists yet.");
559 return Ok(());
560 }
561 let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
562 let mut vault = CredentialVault::unlock(vault_path, passphrase.clone())?;
563 if vault.remove(&name)? {
564 vault.save(vault_path, &passphrase)?;
565 println!("Removed Kennedy secret '{name}'.");
566 } else {
567 println!("Kennedy secret '{name}' was not configured.");
568 }
569 }
570 SecretsCommand::List => {
571 if !vault_path.exists() {
572 println!("No Kennedy credential vault exists yet.");
573 return Ok(());
574 }
575 let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
576 let vault = CredentialVault::unlock(vault_path, passphrase)?;
577 let names = vault.names().collect::<Vec<_>>();
578 if names.is_empty() {
579 println!("The Kennedy credential vault contains no secrets.");
580 } else {
581 println!("Configured Kennedy secrets:");
582 for name in names {
583 println!("- {name}");
584 }
585 }
586 }
587 SecretsCommand::ChangePassphrase => {
588 if !vault_path.exists() {
589 println!("No Kennedy credential vault exists yet.");
590 return Ok(());
591 }
592 let old = prompt_passphrase("Unlock Kennedy credential vault: ")?;
593 let vault = CredentialVault::unlock(vault_path, old)?;
594 let new = prompt_new_vault_passphrase()?;
595 vault.save(vault_path, &new)?;
596 println!("Changed the Kennedy credential vault passphrase.");
597 }
598 }
599 Ok(())
600}
601
602fn unlock_for_edit(path: &Path) -> anyhow::Result<(CredentialVault, SecretString)> {
603 if path.exists() {
604 let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
605 let vault = CredentialVault::unlock(path, passphrase.clone())?;
606 Ok((vault, passphrase))
607 } else {
608 let passphrase = prompt_new_vault_passphrase()?;
609 Ok((CredentialVault::empty(), passphrase))
610 }
611}
612
613fn prompt_passphrase(prompt: &str) -> anyhow::Result<SecretString> {
614 let mut value = rpassword::prompt_password(prompt)?;
615 if value.is_empty() {
616 value.zeroize();
617 anyhow::bail!("the credential vault passphrase cannot be empty");
618 }
619 Ok(SecretString::from(value))
620}
621
622fn prompt_new_vault_passphrase() -> anyhow::Result<SecretString> {
623 let mut first = rpassword::prompt_password("Create Kennedy credential vault passphrase: ")?;
624 let mut second = rpassword::prompt_password("Confirm credential vault passphrase: ")?;
625 if first.is_empty() || first != second {
626 first.zeroize();
627 second.zeroize();
628 anyhow::bail!("credential vault passphrases were empty or did not match");
629 }
630 second.zeroize();
631 Ok(SecretString::from(first))
632}
633
634fn prompt_confirmed_value(prompt: &str) -> anyhow::Result<String> {
635 let mut first = rpassword::prompt_password(prompt)?;
636 let mut second = rpassword::prompt_password("Confirm secret value: ")?;
637 if first.is_empty() || first != second {
638 first.zeroize();
639 second.zeroize();
640 anyhow::bail!("secret values were empty or did not match");
641 }
642 second.zeroize();
643 Ok(first)
644}
645
646#[cfg(test)]
647mod tests {
648 use super::*;
649
650 #[test]
651 fn secret_names_are_stable_code_defaults() {
652 assert_eq!(OPENAI_API_KEY_SECRET, "openai-api-key");
653 assert_eq!(GEMINI_API_KEY_SECRET, "gemini-api-key");
654 assert_eq!(TELEGRAM_BOT_TOKEN_SECRET, "telegram-bot-token");
655 assert_eq!(CRATES_IO_KEY_SECRET, "cratesio-key");
656 assert_eq!(KWEB_WRITER_SIGNING_KEY_SECRET, "kweb-writer-signing-key");
657 assert_eq!(KWEB_WRITERS_SECRET, "kweb-writers-by-priority");
658 }
659
660 #[test]
661 fn persistent_path_defaults_are_under_data() {
662 let args = Args::try_parse_from(["kennedy-server"]).unwrap();
663 for path in [
664 &args.vault_path,
665 &args.kweb_root,
666 &args.conversation_history_database,
667 &args.session_directory,
668 &args.session_history_file,
669 &args.telegram_database,
670 &args.user_database,
671 &args.legacy_audio_ingress_database,
672 &args.audio_ingress_directory,
673 &args.intelligence_usage_directory,
674 &args.rust_libs_root,
675 &args.web_libs_root,
676 &args.web_libs_published_root,
677 &args.rust_bins_root,
678 &args.rust_bin_artifacts_root,
679 ] {
680 assert!(
681 path.starts_with("./data"),
682 "persistent default is outside data/: {}",
683 path.display()
684 );
685 }
686 for path in [
687 &args.rust_libs_root,
688 &args.web_libs_root,
689 &args.web_libs_published_root,
690 &args.rust_bins_root,
691 &args.rust_bin_artifacts_root,
692 ] {
693 assert!(
694 path.starts_with("./data/kcode"),
695 "managed Kcode default is outside data/kcode/: {}",
696 path.display()
697 );
698 }
699 }
700
701 #[test]
702 fn native_orchestration_remains_a_rust_backend_concern() {
703 assert_eq!(
704 std::any::type_name::<kcode_kennedy_orchestration::Session>(),
705 "kcode_kennedy_sessions::Session"
706 );
707 }
708
709 #[tokio::test]
710 async fn unified_dev_tools_service_opens_all_roots_and_routes_three_source_kinds() {
711 let directory = std::env::temp_dir().join(format!(
712 "kennedy-dev-tools-open-test-{}",
713 uuid::Uuid::new_v4()
714 ));
715 let rust_libraries = directory.join("kcode-rust-libs");
716 let web_libraries = directory.join("kcode-web-libs");
717 let web_publications = directory.join("kcode-web-libs-published");
718 let rust_binaries = directory.join("kcode-rust-bins");
719 let rust_binary_artifacts = directory.join("kcode-rust-bin-artifacts");
720 let service = kcode_dev_tools::Service::open(kcode_dev_tools::Config {
721 rust_libraries_root: rust_libraries.clone(),
722 web_libraries_root: web_libraries.clone(),
723 web_publications_root: web_publications.clone(),
724 rust_binaries_root: rust_binaries.clone(),
725 rust_binary_publications_root: rust_binary_artifacts.clone(),
726 crates_io_registry_token: "test-token".into(),
727 })
728 .unwrap();
729
730 assert_eq!(
731 service.web_libraries_root(),
732 std::fs::canonicalize(&web_libraries).unwrap()
733 );
734 assert_eq!(
735 service.web_publications_root(),
736 std::fs::canonicalize(&web_publications).unwrap()
737 );
738 for path in [
739 rust_libraries,
740 web_libraries,
741 web_publications,
742 rust_binaries,
743 rust_binary_artifacts,
744 ] {
745 assert!(
746 path.is_dir(),
747 "managed root was not created: {}",
748 path.display()
749 );
750 }
751 for (create, open, write, name, path, kind) in [
752 (
753 kcode_dev_tools::CREATE_RUST_LIB_TOOL,
754 kcode_dev_tools::OPEN_RUST_LIB_TOOL,
755 kcode_dev_tools::WRITE_FILE_FREEFORM_RUST_LIB_TOOL,
756 "kennedy-test-lib",
757 "src/extra.rs",
758 kcode_dev_tools::ManagedSourceKind::RustLibrary,
759 ),
760 (
761 kcode_dev_tools::CREATE_WEB_LIB_TOOL,
762 kcode_dev_tools::OPEN_WEB_LIB_TOOL,
763 kcode_dev_tools::WRITE_FILE_FREEFORM_WEB_LIB_TOOL,
764 "kennedy-test-web",
765 "extra.js",
766 kcode_dev_tools::ManagedSourceKind::WebLibrary,
767 ),
768 (
769 kcode_dev_tools::CREATE_RUST_BIN_TOOL,
770 kcode_dev_tools::OPEN_RUST_BIN_TOOL,
771 kcode_dev_tools::WRITE_FILE_FREEFORM_RUST_BIN_TOOL,
772 "kennedy-test-bin",
773 "src/extra.rs",
774 kcode_dev_tools::ManagedSourceKind::RustBinary,
775 ),
776 ] {
777 let created = service
778 .execute(
779 "create-session",
780 create,
781 serde_json::json!({"name":name}),
782 Vec::new(),
783 )
784 .await
785 .unwrap();
786 assert_eq!(created.snapshot.unwrap().kind, kind);
787 let written = service
788 .execute(
789 "create-session",
790 write,
791 serde_json::json!({
792 "name":name,
793 "path":path,
794 "contents":"// Kennedy managed source\n",
795 }),
796 Vec::new(),
797 )
798 .await
799 .unwrap();
800 assert_eq!(written.snapshot.unwrap().kind, kind);
801
802 let open_result = service
803 .execute(
804 "open-session",
805 open,
806 serde_json::json!({"name":name}),
807 Vec::new(),
808 )
809 .await
810 .unwrap();
811 assert_eq!(open_result.snapshot.unwrap().kind, kind);
812 }
813 let asset = service
814 .execute(
815 "create-session",
816 kcode_dev_tools::ATTACH_OBJECT_WEB_LIB_TOOL,
817 serde_json::json!({
818 "name":"kennedy-test-web",
819 "path":"assets/fonts/display.woff2",
820 "objectId":"pending:1",
821 }),
822 vec![vec![0, 159, 146, 150, 255]],
823 )
824 .await
825 .unwrap();
826 let snapshot = asset.snapshot.unwrap();
827 assert_eq!(
828 snapshot.kind,
829 kcode_dev_tools::ManagedSourceKind::WebLibrary
830 );
831 assert!(snapshot.text.contains("Asset: assets/fonts/display.woff2"));
832 assert!(snapshot.text.contains("Bytes: 5"));
833 assert!(!snapshot.text.contains("SHA-256:"));
834 assert_eq!(service.release("create-session").await.unwrap(), 3);
835 assert_eq!(service.release("open-session").await.unwrap(), 3);
836 drop(service);
837 std::fs::remove_dir_all(directory).unwrap();
838 }
839
840 #[test]
841 fn missing_optional_secret_disables_only_its_feature() {
842 let vault = CredentialVault::empty();
843 assert!(
844 resolve_optional_secret(&vault, "openai-api-key", "transcription")
845 .unwrap()
846 .is_none()
847 );
848 assert!(
849 resolve_optional_secret(&vault, "", "disabled")
850 .unwrap()
851 .is_none()
852 );
853 }
854
855 #[test]
856 fn required_secret_must_be_present() {
857 let mut vault = CredentialVault::empty();
858 let error =
859 resolve_required_secret(&vault, CRATES_IO_KEY_SECRET, "publication").unwrap_err();
860 assert!(error.to_string().contains(CRATES_IO_KEY_SECRET));
861
862 vault
863 .set(CRATES_IO_KEY_SECRET, "test-crates-io-key".into())
864 .unwrap();
865 assert_eq!(
866 resolve_required_secret(&vault, CRATES_IO_KEY_SECRET, "publication").unwrap(),
867 "test-crates-io-key"
868 );
869 }
870
871 #[test]
872 fn legacy_audio_database_is_copied_once_into_the_persistence_root() {
873 let directory = std::env::temp_dir().join(format!(
874 "kennedy-audio-migration-test-{}",
875 uuid::Uuid::new_v4()
876 ));
877 std::fs::create_dir(&directory).unwrap();
878 let legacy = directory.join("legacy.sqlite3");
879 let current = directory.join("audio-ingress/state.sqlite3");
880 let source = rusqlite::Connection::open(&legacy).unwrap();
881 source
882 .execute_batch("CREATE TABLE marker(value TEXT NOT NULL);")
883 .unwrap();
884 source
885 .execute("INSERT INTO marker(value) VALUES('legacy')", [])
886 .unwrap();
887 drop(source);
888
889 migrate_audio_ingress_database(&legacy, ¤t).unwrap();
890 let migrated = rusqlite::Connection::open(¤t).unwrap();
891 let value: String = migrated
892 .query_row("SELECT value FROM marker", [], |row| row.get(0))
893 .unwrap();
894 assert_eq!(value, "legacy");
895 migrated
896 .execute("UPDATE marker SET value='current'", [])
897 .unwrap();
898 drop(migrated);
899
900 migrate_audio_ingress_database(&legacy, ¤t).unwrap();
901 let value: String = rusqlite::Connection::open(¤t)
902 .unwrap()
903 .query_row("SELECT value FROM marker", [], |row| row.get(0))
904 .unwrap();
905 assert_eq!(value, "current");
906 std::fs::remove_dir_all(directory).unwrap();
907 }
908
909 #[tokio::test]
910 async fn occupied_kweb_address_prevents_server_from_opening_persistent_state() {
911 let directory =
912 std::env::temp_dir().join(format!("kennedy-server-lock-test-{}", uuid::Uuid::new_v4()));
913 std::fs::create_dir(&directory).unwrap();
914 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
915 let bind = listener.local_addr().unwrap().to_string();
916 let vault = directory.join("vault.age");
917 let kmap = directory.join("kweb");
918 let conversations = directory.join("conversations.sqlite3");
919 let telegram = directory.join("telegram.sqlite3");
920 let users = directory.join("users.sqlite3");
921 let audio = directory.join("audio.sqlite3");
922 let audio_media = directory.join("audio-media");
923 let args = Args {
924 vault_path: vault.clone(),
925 command: None,
926 kweb_bind: bind,
927 kweb_root: kmap.clone(),
928 conversation_history_database: conversations.clone(),
929 session_directory: directory.join("sessions"),
930 session_history_file: directory.join("session-history.txt"),
931 telegram_database: telegram.clone(),
932 user_database: users.clone(),
933 legacy_audio_ingress_database: audio.clone(),
934 audio_ingress_directory: audio_media.clone(),
935 intelligence_usage_directory: directory.join("intelligence-usage"),
936 rust_libs_root: directory.join("rust-libs"),
937 web_libs_root: directory.join("kcode-web-libs"),
938 web_libs_published_root: directory.join("kcode-web-libs-published"),
939 rust_bins_root: directory.join("kcode-rust-bins"),
940 rust_bin_artifacts_root: directory.join("kcode-rust-bin-artifacts"),
941 telegram_bootstrap_username: "@test".to_owned(),
942 telegram_max_voice_bytes: 1024,
943 audio_ingress_max_upload_bytes: 1024,
944 };
945
946 let error = run_server(args, vault.clone()).await.unwrap_err();
947 assert!(error.to_string().contains("binding Kweb listener"));
948 assert!(!vault.exists());
949 assert!(!kmap.exists());
950 assert!(!conversations.exists());
951 assert!(!telegram.exists());
952 assert!(!users.exists());
953 assert!(!audio.exists());
954 assert!(!audio_media.exists());
955 std::fs::remove_dir_all(directory).unwrap();
956 }
957}