Skip to main content

kcode_kennedy_app/
lib.rs

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    /// Create and manage generic named secrets in Kennedy's encrypted vault.
92    Secrets {
93        #[command(subcommand)]
94        command: SecretsCommand,
95    },
96    /// Estimate the token footprint of all current Kmap node text.
97    KmapSize,
98}
99
100#[derive(Subcommand, Debug)]
101enum SecretsCommand {
102    /// Prompt for and store a named secret, replacing any previous value.
103    Set { name: String },
104    /// Remove a named secret without displaying its value.
105    Remove { name: String },
106    /// List configured secret names without displaying their values.
107    List,
108    /// Re-encrypt the vault with a new passphrase.
109    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    // Bind the public Kennedy address before opening any persistent state.
153    // Offline maintenance checks this address before copying the data tree.
154    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, &current).unwrap();
890        let migrated = rusqlite::Connection::open(&current).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, &current).unwrap();
901        let value: String = rusqlite::Connection::open(&current)
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}