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 kcode_credential_vault::{CredentialVault, ExposeSecret, SecretString};
11use kcode_kennedy_cli::{Args, Command, SecretsCommand};
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#[tokio::main]
25pub async fn main() -> anyhow::Result<()> {
26    tracing_subscriber::fmt()
27        .with_env_filter(
28            tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| {
29                "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()
30            }),
31        )
32        .init();
33    rustls::crypto::ring::default_provider()
34        .install_default()
35        .map_err(|_| anyhow::anyhow!("installing TLS crypto provider"))?;
36    let mut args = kcode_kennedy_cli::parse();
37    let vault_path = args.vault_path.clone();
38    match args.command.take() {
39        Some(Command::Secrets { command }) => {
40            let _maintenance_guard = tokio::net::TcpListener::bind(&args.kweb_bind)
41                .await
42                .with_context(|| {
43                    format!(
44                        "binding maintenance lock {}; stop the running Kennedy server before changing its credential vault",
45                        args.kweb_bind
46                    )
47                })?;
48            manage_secrets(command, &vault_path)
49        }
50        Some(Command::KmapSize) => {
51            let _maintenance_guard =
52                maintenance_guard(&args.kweb_bind, "measuring the Kweb").await?;
53            let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
54            let vault = CredentialVault::unlock(&vault_path, passphrase)?;
55            let size = kcode_kmap_size::measure(&args.kweb_root, kweb_config(&vault)?)?;
56            println!("{}", kcode_kmap_size::render(&size));
57            Ok(())
58        }
59        None => run_server(args, vault_path).await,
60    }
61}
62
63async fn run_server(args: Args, vault_path: PathBuf) -> anyhow::Result<()> {
64    // Bind the public Kennedy address before opening any persistent state.
65    // Offline maintenance checks this address before copying the data tree.
66    let kweb_listener = tokio::net::TcpListener::bind(&args.kweb_bind)
67        .await
68        .with_context(|| format!("binding Kweb listener {}", args.kweb_bind))?;
69    ensure_runtime_parent_directories(&args, &vault_path)?;
70    let vault = if vault_path.exists() {
71        let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
72        CredentialVault::unlock(&vault_path, passphrase)?
73    } else {
74        tracing::warn!(path=%vault_path.display(), "Kennedy credential vault does not exist; secret-backed features are unavailable");
75        CredentialVault::empty()
76    };
77    let openai_api_key = resolve_optional_secret(
78        &vault,
79        OPENAI_API_KEY_SECRET,
80        "OpenAI transcription, media annotation, agents, and image generation/editing",
81    )?;
82    let gemini_api_key = resolve_optional_secret(
83        &vault,
84        GEMINI_API_KEY_SECRET,
85        "Gemini search, media annotation, agents, audio transcription, and image generation/editing",
86    )?;
87    let telegram_bot_token =
88        resolve_optional_secret(&vault, TELEGRAM_BOT_TOKEN_SECRET, "Telegram relay")?
89            .map(kcode_tg_kennedy_bot::BotToken::new)
90            .transpose()?;
91    let crates_io_key =
92        resolve_required_secret(&vault, CRATES_IO_KEY_SECRET, "Rust library publication")?;
93    let kweb_config = kweb_config(&vault)?;
94    let codex_catalog_cache =
95        kcode_codex_runtime::CatalogCache::new(kcode_codex_runtime::DEFAULT_CODEX_EXECUTABLE);
96    let (kmap, system_roots) =
97        kcode_kennedy_roots::open(&args.kweb_root, kweb_config, &args.user_database)?;
98    let (kmap_commands, kmap_command_runtime) =
99        kcode_kmap_command_lane::open(&args.user_database, kmap.clone())?;
100    let credits = kcode_credits::Credits::open(&args.credits_database)?;
101    let task_board = kcode_task_board::TaskBoard::open(&args.task_board_database, credits.clone())?;
102    let speech_classifier = SpeechClassifier::open(SPEECH_CLASSIFICATION_DATABASE_PATH)
103        .with_context(|| {
104            format!("opening speaker-classification database {SPEECH_CLASSIFICATION_DATABASE_PATH}")
105        })?;
106    let speech_classifier = Arc::new(speech_classifier);
107    let dev_tools = kcode_dev_tools::Service::open(kcode_dev_tools::Config {
108        rust_libraries_root: args.rust_libs_root.clone(),
109        web_libraries_root: args.web_libs_root.clone(),
110        web_publications_root: args.web_libs_published_root.clone(),
111        rust_binaries_root: args.rust_bins_root.clone(),
112        rust_binary_publications_root: args.rust_bin_artifacts_root.clone(),
113        crates_io_registry_token: crates_io_key,
114    })
115    .map_err(anyhow::Error::new)
116    .with_context(|| {
117        format!(
118            "opening managed Kcode development roots under {}",
119            args.rust_libs_root
120                .parent()
121                .unwrap_or(Path::new("."))
122                .display()
123        )
124    })?;
125    let web_publications_root = dev_tools.web_publications_root().to_path_buf();
126    let telegram_identity = std::sync::Arc::new(kcode_telegram_identity::Directory::open(
127        &args.user_database,
128        &args.telegram_bootstrap_username,
129    )?);
130    let history_service =
131        kcode_session_history::SessionHistory::open(kcode_session_history::Config {
132            directory: args.session_directory,
133            completed_list: args.session_history_file,
134            provider_cost_compatibility: Some(
135                kcode_intelligence_chatend::provider_cost_compatibility(),
136            ),
137        })?;
138    let (intelligence_service, intelligence_runtime) =
139        kcode_intelligence_router::open(kcode_intelligence_router::Config {
140            openai_api_key,
141            gemini_api_key,
142            codex_catalog_cache,
143            receipt_directory: args.intelligence_usage_directory,
144        })
145        .await?;
146    let agent_runtime = kcode_agent_runtime::AgentRuntime::new(intelligence_service.clone());
147    let telegram_runtime = kcode_tg_kennedy_bot::open(kcode_tg_kennedy_bot::Config {
148        database: args.telegram_database,
149        bot_token: telegram_bot_token,
150        identity_sink: telegram_identity.clone(),
151        max_voice_bytes: args.telegram_max_voice_bytes,
152    })
153    .await?;
154    let telegram_service = telegram_runtime.service();
155    let chunk_intelligence = intelligence_service.clone();
156    let transcribe_chunk: kcode_audio_ingress::AudioChunkCall = Arc::new(move |request| {
157        let intelligence = chunk_intelligence.clone();
158        Box::pin(async move {
159            let user = intelligence
160                .for_user(request.user_id)
161                .map_err(audio_intelligence_error)?;
162            let media = kcode_intelligence_router::Media::audio(
163                request.audio_ogg,
164                "audio-chunk.ogg",
165                "audio/ogg",
166            )
167            .map_err(audio_intelligence_error)?;
168            user.analyze_audio(kcode_intelligence_router::AudioAnalysisRequest {
169                operation: "transcribe_chunk".into(),
170                prompt: request.prompt,
171                model: request.model,
172                media,
173                schema: request.schema,
174                max_output_tokens: request.max_output_tokens,
175                temperature: None,
176                operation_id: uuid::Uuid::new_v4(),
177                parent_operation_id: None,
178            })
179            .await
180            .map(|response| response.value.text)
181            .map_err(audio_intelligence_error)
182        })
183    });
184    let text_intelligence = intelligence_service.clone();
185    let generate_text: kcode_audio_ingress::TextGenerationCall = Arc::new(move |request| {
186        let intelligence = text_intelligence.clone();
187        Box::pin(async move {
188            let reasoning_effort = match request.reasoning_effort.as_str() {
189                "xhigh" => kcode_intelligence_router::ReasoningEffort::XHigh,
190                _ => {
191                    return Err(kcode_audio_ingress::IntelligenceError::new(
192                        "AudioIngress requested an unsupported reasoning effort.",
193                        false,
194                    ));
195                }
196            };
197            let user = intelligence
198                .for_user(request.user_id)
199                .map_err(audio_intelligence_error)?;
200            user.generate_text(kcode_intelligence_router::TextGenerationRequest {
201                operation: request.operation,
202                prompt: request.prompt,
203                model: request.model,
204                reasoning_effort,
205                timeout: request.timeout,
206                operation_id: uuid::Uuid::new_v4(),
207                parent_operation_id: None,
208            })
209            .await
210            .map(|response| response.value.text)
211            .map_err(audio_intelligence_error)
212        })
213    });
214    let audio_transcriber =
215        kcode_audio_ingress::AudioTranscriber::new(transcribe_chunk, generate_text);
216    let audio = kcode_audio_ingress::AudioIngress::open(
217        &args.audio_ingress_directory,
218        audio_transcriber,
219        Arc::clone(&speech_classifier),
220    )
221    .await?;
222    let audio_coordinator = kcode_audio_session_ingress::Coordinator::new(
223        audio,
224        history_service.clone(),
225        kcode_audio_session_ingress::Config {
226            user_id: system_roots.user.to_string(),
227            effective_context_tokens: intelligence_runtime.context_window_tokens,
228        },
229    )?;
230    let http_router = kcode_http_api::router(kcode_http_api::Config {
231        kmap: kmap.clone(),
232        kmap_commands,
233        user_root_node_id: system_roots.user,
234        kennedy_root_node_id: system_roots.kennedy,
235        telegram: telegram_service.clone(),
236        session_history: history_service.clone(),
237        audio_ingress: audio_coordinator.clone(),
238        audio_max_upload_bytes: args.audio_ingress_max_upload_bytes,
239        task_board: task_board.clone(),
240        credits,
241        web_publications_root,
242    })?;
243    let orchestration_config = kcode_kennedy_orchestration::Config {
244        user_root_node_id: system_roots.user.to_string(),
245        kennedy_root_node_id: system_roots.kennedy.to_string(),
246        telegram_max_media_bytes: args.telegram_max_voice_bytes,
247        runtime_model: kcode_kennedy_orchestration::RuntimeModel::from_intelligence(
248            intelligence_runtime,
249        ),
250    };
251    let telegram_sessions = kcode_telegram_session_coordinator::Service::new(
252        telegram_service.clone(),
253        telegram_identity.clone(),
254    );
255    let session_service =
256        kcode_kennedy_sessions::Service::new(kcode_kennedy_sessions::Capabilities {
257            load_fixed_connections: args.fixed,
258            kmap: kmap.clone(),
259            intelligence: intelligence_service.clone(),
260            agents: agent_runtime,
261            history: history_service.clone(),
262            speech_classifier,
263            dev_tools: dev_tools.clone(),
264            telegram: telegram_sessions,
265        })
266        .with_task_board(task_board);
267    let orchestration_api = kcode_kennedy_orchestration::Api::new(
268        &orchestration_config,
269        kcode_kennedy_orchestration::LocalServices {
270            kmap: kmap.clone(),
271            intelligence: intelligence_service,
272            history: history_service.clone(),
273            audio: audio_coordinator,
274            directory: telegram_identity.clone(),
275            dev_tools,
276            telegram: telegram_service,
277        },
278    );
279    let orchestration_worker = kcode_kennedy_orchestration::build(
280        orchestration_config,
281        orchestration_api,
282        session_service,
283    );
284    let directory_roots = kcode_kennedy_roots::DirectoryRoots::new(
285        kmap,
286        telegram_identity,
287        args.telegram_bootstrap_username.clone(),
288        system_roots.user,
289        orchestration_worker.writer().clone(),
290    );
291    let telegram_session_runtime = Arc::new(kcode_kennedy_telegram_runtime::Runtime::new(
292        kcode_kennedy_telegram_runtime::Config {
293            telegram_max_media_bytes: args.telegram_max_voice_bytes,
294            telegram_web_user_handle: args.telegram_bootstrap_username,
295        },
296        orchestration_worker.clone(),
297        directory_roots,
298    ));
299    tokio::try_join!(
300        async {
301            kcode_http_api::serve(kweb_listener, http_router)
302                .await
303                .map_err(anyhow::Error::new)
304        },
305        telegram_runtime.run(),
306        kcode_kennedy_orchestration::run(orchestration_worker),
307        telegram_session_runtime.run(),
308        async { kmap_command_runtime.await.map_err(anyhow::Error::new) },
309    )?;
310    Ok(())
311}
312
313fn audio_intelligence_error(
314    error: kcode_intelligence_router::Error,
315) -> kcode_audio_ingress::IntelligenceError {
316    let retryable = error.retryable();
317    kcode_audio_ingress::IntelligenceError::new(error.message(), retryable)
318}
319
320fn ensure_runtime_parent_directories(args: &Args, vault_path: &Path) -> anyhow::Result<()> {
321    for path in [
322        vault_path,
323        &args.kweb_root,
324        &args.conversation_history_database,
325        &args.session_directory,
326        &args.session_history_file,
327        &args.telegram_database,
328        &args.user_database,
329        &args.task_board_database,
330        &args.credits_database,
331        Path::new(SPEECH_CLASSIFICATION_DATABASE_PATH),
332        &args.audio_ingress_directory,
333        &args.intelligence_usage_directory,
334        &args.rust_libs_root,
335        &args.web_libs_root,
336        &args.web_libs_published_root,
337        &args.rust_bins_root,
338        &args.rust_bin_artifacts_root,
339    ] {
340        let Some(parent) = path.parent().filter(|value| !value.as_os_str().is_empty()) else {
341            continue;
342        };
343        if parent.exists() {
344            continue;
345        }
346        let mut builder = std::fs::DirBuilder::new();
347        builder.recursive(true);
348        #[cfg(unix)]
349        {
350            use std::os::unix::fs::DirBuilderExt;
351            builder.mode(0o700);
352        }
353        builder
354            .create(parent)
355            .with_context(|| format!("creating runtime data directory {}", parent.display()))?;
356    }
357    Ok(())
358}
359
360pub async fn maintenance_guard(
361    bind: &str,
362    purpose: &str,
363) -> anyhow::Result<tokio::net::TcpListener> {
364    tokio::net::TcpListener::bind(bind).await.with_context(|| {
365        format!("binding maintenance lock {bind}; stop the running Kennedy server before {purpose}")
366    })
367}
368
369fn kweb_config(vault: &CredentialVault) -> anyhow::Result<KwebConfig> {
370    let encoded_key = resolve_required_secret(
371        vault,
372        KWEB_WRITER_SIGNING_KEY_SECRET,
373        "Kweb mutation signing",
374    )?;
375    let mut signing_key = Zeroizing::new([0_u8; 32]);
376    let decoded = hex::decode(encoded_key.trim())
377        .context("Kweb writer signing key must be 64 lowercase hexadecimal characters")?;
378    *signing_key = decoded
379        .try_into()
380        .map_err(|_| anyhow::anyhow!("Kweb writer signing key must decode to exactly 32 bytes"))?;
381    let encoded_writers =
382        resolve_required_secret(vault, KWEB_WRITERS_SECRET, "Kweb writer authorization")?;
383    let writers_by_priority = encoded_writers
384        .split(',')
385        .map(str::trim)
386        .filter(|value| !value.is_empty())
387        .map(WriterId::from_str)
388        .collect::<Result<Vec<_>, _>>()
389        .map_err(anyhow::Error::new)
390        .context("decoding the ordered Kweb writer whitelist")?;
391    anyhow::ensure!(
392        !writers_by_priority.is_empty(),
393        "the Kweb writer whitelist is empty"
394    );
395    Ok(KwebConfig {
396        signing_key: *signing_key,
397        writers_by_priority,
398        gossip: Arc::new(NoopGossip),
399    })
400}
401
402fn resolve_optional_secret(
403    vault: &CredentialVault,
404    configured_name: &str,
405    purpose: &str,
406) -> anyhow::Result<Option<String>> {
407    let name = configured_name.trim();
408    if name.is_empty() {
409        return Ok(None);
410    }
411    let secret = vault.secret(name)?;
412    if secret.is_none() {
413        tracing::warn!(secret_name=name, %purpose, "configured Kennedy secret is not present in the vault");
414    }
415    Ok(secret.map(|value| value.expose_secret().to_owned()))
416}
417
418fn resolve_required_secret(
419    vault: &CredentialVault,
420    configured_name: &str,
421    purpose: &str,
422) -> anyhow::Result<String> {
423    let name = configured_name.trim();
424    if name.is_empty() {
425        anyhow::bail!("{purpose} requires a configured Kennedy secret name");
426    }
427    vault
428        .secret(name)?
429        .map(|value| value.expose_secret().to_owned())
430        .with_context(|| {
431            format!(
432                "{purpose} requires Kennedy secret '{name}'; store it with `kennedy-server secrets set {name}`"
433            )
434        })
435}
436
437fn manage_secrets(command: SecretsCommand, vault_path: &Path) -> anyhow::Result<()> {
438    match command {
439        SecretsCommand::Set { name } => {
440            let (mut vault, passphrase) = unlock_for_edit(vault_path)?;
441            let value = prompt_confirmed_value(&format!("Value for {name}: "))?;
442            vault.set(&name, value)?;
443            vault.save(vault_path, &passphrase)?;
444            println!("Stored Kennedy secret '{name}'.");
445        }
446        SecretsCommand::Remove { name } => {
447            if !vault_path.exists() {
448                println!("No Kennedy credential vault exists yet.");
449                return Ok(());
450            }
451            let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
452            let mut vault = CredentialVault::unlock(vault_path, passphrase.clone())?;
453            if vault.remove(&name)? {
454                vault.save(vault_path, &passphrase)?;
455                println!("Removed Kennedy secret '{name}'.");
456            } else {
457                println!("Kennedy secret '{name}' was not configured.");
458            }
459        }
460        SecretsCommand::List => {
461            if !vault_path.exists() {
462                println!("No Kennedy credential vault exists yet.");
463                return Ok(());
464            }
465            let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
466            let vault = CredentialVault::unlock(vault_path, passphrase)?;
467            let names = vault.names().collect::<Vec<_>>();
468            if names.is_empty() {
469                println!("The Kennedy credential vault contains no secrets.");
470            } else {
471                println!("Configured Kennedy secrets:");
472                for name in names {
473                    println!("- {name}");
474                }
475            }
476        }
477        SecretsCommand::ChangePassphrase => {
478            if !vault_path.exists() {
479                println!("No Kennedy credential vault exists yet.");
480                return Ok(());
481            }
482            let old = prompt_passphrase("Unlock Kennedy credential vault: ")?;
483            let vault = CredentialVault::unlock(vault_path, old)?;
484            let new = prompt_new_vault_passphrase()?;
485            vault.save(vault_path, &new)?;
486            println!("Changed the Kennedy credential vault passphrase.");
487        }
488    }
489    Ok(())
490}
491
492fn unlock_for_edit(path: &Path) -> anyhow::Result<(CredentialVault, SecretString)> {
493    if path.exists() {
494        let passphrase = prompt_passphrase("Unlock Kennedy credential vault: ")?;
495        let vault = CredentialVault::unlock(path, passphrase.clone())?;
496        Ok((vault, passphrase))
497    } else {
498        let passphrase = prompt_new_vault_passphrase()?;
499        Ok((CredentialVault::empty(), passphrase))
500    }
501}
502
503fn prompt_passphrase(prompt: &str) -> anyhow::Result<SecretString> {
504    let mut value = rpassword::prompt_password(prompt)?;
505    if value.is_empty() {
506        value.zeroize();
507        anyhow::bail!("the credential vault passphrase cannot be empty");
508    }
509    Ok(SecretString::from(value))
510}
511
512fn prompt_new_vault_passphrase() -> anyhow::Result<SecretString> {
513    let mut first = rpassword::prompt_password("Create Kennedy credential vault passphrase: ")?;
514    let mut second = rpassword::prompt_password("Confirm Kennedy credential vault passphrase: ")?;
515    if first.is_empty() || first != second {
516        first.zeroize();
517        second.zeroize();
518        anyhow::bail!("credential vault passphrases were empty or did not match");
519    }
520    second.zeroize();
521    Ok(SecretString::from(first))
522}
523
524fn prompt_confirmed_value(prompt: &str) -> anyhow::Result<String> {
525    let mut first = rpassword::prompt_password(prompt)?;
526    let mut second = rpassword::prompt_password("Confirm secret value: ")?;
527    if first.is_empty() || first != second {
528        first.zeroize();
529        second.zeroize();
530        anyhow::bail!("secret values were empty or did not match");
531    }
532    second.zeroize();
533    Ok(first)
534}
535
536#[cfg(test)]
537mod tests {
538    use super::*;
539
540    #[test]
541    fn secret_names_are_stable_code_defaults() {
542        assert_eq!(OPENAI_API_KEY_SECRET, "openai-api-key");
543        assert_eq!(GEMINI_API_KEY_SECRET, "gemini-api-key");
544        assert_eq!(TELEGRAM_BOT_TOKEN_SECRET, "telegram-bot-token");
545        assert_eq!(CRATES_IO_KEY_SECRET, "cratesio-key");
546        assert_eq!(KWEB_WRITER_SIGNING_KEY_SECRET, "kweb-writer-signing-key");
547        assert_eq!(KWEB_WRITERS_SECRET, "kweb-writers-by-priority");
548    }
549
550    #[test]
551    fn native_orchestration_remains_a_rust_backend_concern() {
552        assert_eq!(
553            std::any::type_name::<kcode_kennedy_orchestration::Session>(),
554            "kcode_kennedy_sessions::Session"
555        );
556    }
557
558    #[test]
559    fn marker_release_selects_exact_orchestration_and_sessions_versions() {
560        let manifest = include_str!("../Cargo.toml");
561        assert!(manifest.contains("kcode-kennedy-orchestration = \"=0.3.2\""));
562        assert!(manifest.contains("kcode-kennedy-sessions = \"0.2.1\""));
563    }
564
565    #[tokio::test]
566    async fn unified_dev_tools_service_opens_all_roots_and_routes_three_source_kinds() {
567        let directory = std::env::temp_dir().join(format!(
568            "kennedy-dev-tools-open-test-{}",
569            uuid::Uuid::new_v4()
570        ));
571        let rust_libraries = directory.join("kcode-rust-libs");
572        let web_libraries = directory.join("kcode-web-libs");
573        let web_publications = directory.join("kcode-web-libs-published");
574        let rust_binaries = directory.join("kcode-rust-bins");
575        let rust_binary_artifacts = directory.join("kcode-rust-bin-artifacts");
576        let service = kcode_dev_tools::Service::open(kcode_dev_tools::Config {
577            rust_libraries_root: rust_libraries.clone(),
578            web_libraries_root: web_libraries.clone(),
579            web_publications_root: web_publications.clone(),
580            rust_binaries_root: rust_binaries.clone(),
581            rust_binary_publications_root: rust_binary_artifacts.clone(),
582            crates_io_registry_token: "test-token".into(),
583        })
584        .unwrap();
585
586        assert_eq!(
587            service.web_libraries_root(),
588            std::fs::canonicalize(&web_libraries).unwrap()
589        );
590        assert_eq!(
591            service.web_publications_root(),
592            std::fs::canonicalize(&web_publications).unwrap()
593        );
594        for path in [
595            rust_libraries,
596            web_libraries,
597            web_publications,
598            rust_binaries,
599            rust_binary_artifacts,
600        ] {
601            assert!(
602                path.is_dir(),
603                "managed root was not created: {}",
604                path.display()
605            );
606        }
607        for (create, open, write, name, path, kind) in [
608            (
609                kcode_dev_tools::CREATE_RUST_LIB_TOOL,
610                kcode_dev_tools::OPEN_RUST_LIB_TOOL,
611                kcode_dev_tools::WRITE_FILE_FREEFORM_RUST_LIB_TOOL,
612                "kennedy-test-lib",
613                "src/extra.rs",
614                kcode_dev_tools::ManagedSourceKind::RustLibrary,
615            ),
616            (
617                kcode_dev_tools::CREATE_WEB_LIB_TOOL,
618                kcode_dev_tools::OPEN_WEB_LIB_TOOL,
619                kcode_dev_tools::WRITE_FILE_FREEFORM_WEB_LIB_TOOL,
620                "kennedy-test-web",
621                "extra.js",
622                kcode_dev_tools::ManagedSourceKind::WebLibrary,
623            ),
624            (
625                kcode_dev_tools::CREATE_RUST_BIN_TOOL,
626                kcode_dev_tools::OPEN_RUST_BIN_TOOL,
627                kcode_dev_tools::WRITE_FILE_FREEFORM_RUST_BIN_TOOL,
628                "kennedy-test-bin",
629                "src/extra.rs",
630                kcode_dev_tools::ManagedSourceKind::RustBinary,
631            ),
632        ] {
633            let created = service
634                .execute(
635                    "create-session",
636                    create,
637                    serde_json::json!({"name":name}),
638                    Vec::new(),
639                )
640                .await
641                .unwrap();
642            assert_eq!(created.snapshot.unwrap().kind, kind);
643            let written = service
644                .execute(
645                    "create-session",
646                    write,
647                    serde_json::json!({
648                        "name":name,
649                        "path":path,
650                        "contents":"// Kennedy managed source\n",
651                    }),
652                    Vec::new(),
653                )
654                .await
655                .unwrap();
656            assert_eq!(written.snapshot.unwrap().kind, kind);
657
658            let open_result = service
659                .execute(
660                    "open-session",
661                    open,
662                    serde_json::json!({"name":name}),
663                    Vec::new(),
664                )
665                .await
666                .unwrap();
667            assert_eq!(open_result.snapshot.unwrap().kind, kind);
668        }
669        let asset = service
670            .execute(
671                "create-session",
672                kcode_dev_tools::ATTACH_OBJECT_WEB_LIB_TOOL,
673                serde_json::json!({
674                    "name":"kennedy-test-web",
675                    "path":"assets/fonts/display.woff2",
676                    "objectId":"pending:1",
677                }),
678                vec![vec![0, 159, 146, 150, 255]],
679            )
680            .await
681            .unwrap();
682        let snapshot = asset.snapshot.unwrap();
683        assert_eq!(
684            snapshot.kind,
685            kcode_dev_tools::ManagedSourceKind::WebLibrary
686        );
687        assert!(snapshot.text.contains("Asset: assets/fonts/display.woff2"));
688        assert!(snapshot.text.contains("Bytes: 5"));
689        assert!(!snapshot.text.contains("SHA-256:"));
690        assert_eq!(service.release("create-session").await.unwrap(), 3);
691        assert_eq!(service.release("open-session").await.unwrap(), 3);
692        drop(service);
693        std::fs::remove_dir_all(directory).unwrap();
694    }
695
696    #[test]
697    fn missing_optional_secret_disables_only_its_feature() {
698        let vault = CredentialVault::empty();
699        assert!(
700            resolve_optional_secret(&vault, "openai-api-key", "transcription")
701                .unwrap()
702                .is_none()
703        );
704        assert!(
705            resolve_optional_secret(&vault, "", "disabled")
706                .unwrap()
707                .is_none()
708        );
709    }
710
711    #[test]
712    fn required_secret_must_be_present() {
713        let mut vault = CredentialVault::empty();
714        let error =
715            resolve_required_secret(&vault, CRATES_IO_KEY_SECRET, "publication").unwrap_err();
716        assert!(error.to_string().contains(CRATES_IO_KEY_SECRET));
717
718        vault
719            .set(CRATES_IO_KEY_SECRET, "test-crates-io-key".into())
720            .unwrap();
721        assert_eq!(
722            resolve_required_secret(&vault, CRATES_IO_KEY_SECRET, "publication").unwrap(),
723            "test-crates-io-key"
724        );
725    }
726
727    #[tokio::test]
728    async fn occupied_kweb_address_prevents_server_from_opening_persistent_state() {
729        let directory =
730            std::env::temp_dir().join(format!("kennedy-server-lock-test-{}", uuid::Uuid::new_v4()));
731        std::fs::create_dir(&directory).unwrap();
732        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
733        let bind = listener.local_addr().unwrap().to_string();
734        let vault = directory.join("vault.age");
735        let kmap = directory.join("kweb");
736        let conversations = directory.join("conversations.sqlite3");
737        let telegram = directory.join("telegram.sqlite3");
738        let users = directory.join("users.sqlite3");
739        let tasks = directory.join("tasks.sqlite3");
740        let credits = directory.join("credits.sqlite3");
741        let audio_media = directory.join("audio-media");
742        let args = Args {
743            vault_path: vault.clone(),
744            command: None,
745            kweb_bind: bind,
746            kweb_root: kmap.clone(),
747            conversation_history_database: conversations.clone(),
748            session_directory: directory.join("sessions"),
749            session_history_file: directory.join("session-history.txt"),
750            telegram_database: telegram.clone(),
751            user_database: users.clone(),
752            task_board_database: tasks.clone(),
753            credits_database: credits.clone(),
754            audio_ingress_directory: audio_media.clone(),
755            intelligence_usage_directory: directory.join("intelligence-usage"),
756            rust_libs_root: directory.join("rust-libs"),
757            web_libs_root: directory.join("kcode-web-libs"),
758            web_libs_published_root: directory.join("kcode-web-libs-published"),
759            rust_bins_root: directory.join("kcode-rust-bins"),
760            rust_bin_artifacts_root: directory.join("kcode-rust-bin-artifacts"),
761            telegram_bootstrap_username: "@test".to_owned(),
762            telegram_max_voice_bytes: 1024,
763            audio_ingress_max_upload_bytes: 1024,
764            fixed: false,
765        };
766
767        let error = run_server(args, vault.clone()).await.unwrap_err();
768        assert!(error.to_string().contains("binding Kweb listener"));
769        assert!(!vault.exists());
770        assert!(!kmap.exists());
771        assert!(!conversations.exists());
772        assert!(!telegram.exists());
773        assert!(!users.exists());
774        assert!(!tasks.exists());
775        assert!(!credits.exists());
776        assert!(!audio_media.exists());
777        std::fs::remove_dir_all(directory).unwrap();
778    }
779}