mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
Finish vault-backed workflow auth and installer QA fixes
Pass the shared storage dir into worker runs so vault-backed credentials load during real workflow execution, including server-spawned workers. Also finish the QA follow-ups around scripted install behavior, list credential metadata in secret listings, and give the slow OpenAPI conformance test a narrow nextest timeout override.
This commit is contained in:
parent
f3aa30d782
commit
e980008a52
9 changed files with 292 additions and 13 deletions
|
|
@ -11,6 +11,10 @@ leak-timeout = "500ms"
|
|||
filter = "package(fabro-server) & kind(test)"
|
||||
slow-timeout = { period = "5s", terminate-after = 4 }
|
||||
|
||||
[[profile.default.overrides]]
|
||||
filter = "package(fabro-server) & test(all_spec_routes_are_routable)"
|
||||
slow-timeout = { period = "15s", terminate-after = 4 }
|
||||
|
||||
[[profile.default.overrides]]
|
||||
filter = "package(fabro-workflow) & kind(test)"
|
||||
slow-timeout = { period = "2s", terminate-after = 3 }
|
||||
|
|
|
|||
|
|
@ -751,6 +751,10 @@ pub(crate) struct RunWorkerArgs {
|
|||
#[arg(long)]
|
||||
pub(crate) server: String,
|
||||
|
||||
/// Fabro storage directory for loading worker-visible secrets
|
||||
#[arg(long, hide = true)]
|
||||
pub(crate) storage_dir: Option<PathBuf>,
|
||||
|
||||
/// Short-lived bearer token for artifact uploads
|
||||
#[arg(long, hide = true)]
|
||||
pub(crate) artifact_upload_token: Option<String>,
|
||||
|
|
@ -1272,8 +1276,9 @@ pub(crate) struct DoctorArgs {
|
|||
pub(crate) verbose: bool,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy, ValueEnum)]
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, ValueEnum)]
|
||||
pub(crate) enum InstallGitHubStrategyArg {
|
||||
#[value(name = "gh_cli", alias = "gh-cli")]
|
||||
GhCli,
|
||||
App,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -238,6 +238,10 @@ fn merge_server_settings(doc: &mut toml::Value, username: &str) -> Result<()> {
|
|||
|
||||
let listen = ensure_table(server, "listen")?;
|
||||
listen.insert("type".to_string(), toml::Value::String("tcp".to_string()));
|
||||
listen.insert(
|
||||
"address".to_string(),
|
||||
toml::Value::String("127.0.0.1:3000".to_string()),
|
||||
);
|
||||
let listen_tls = ensure_table(listen, "tls")?;
|
||||
let certs_dir = fabro_util::Home::from_env().certs_dir();
|
||||
listen_tls.insert(
|
||||
|
|
@ -1204,10 +1208,11 @@ pub(crate) async fn run_install(
|
|||
let server_was_running = record::active_server_record(&storage_dir).is_some();
|
||||
let fabro_dir = fabro_util::Home::from_env().root().to_path_buf();
|
||||
let config_path = fabro_dir.join(SETTINGS_CONFIG_FILENAME);
|
||||
let config_existed_before_install = config_path.exists();
|
||||
let input_source: Box<dyn InstallInputSource + Send + Sync> =
|
||||
match NonInteractiveInstallInputSource::new(args)? {
|
||||
Some(source) => {
|
||||
source.validate(config_path.exists())?;
|
||||
source.validate(config_existed_before_install)?;
|
||||
Box::new(source)
|
||||
}
|
||||
None => Box::new(InteractiveInstallInputSource),
|
||||
|
|
@ -1370,7 +1375,7 @@ pub(crate) async fn run_install(
|
|||
fabro_util::printerr!(printer, "");
|
||||
|
||||
match input_source
|
||||
.choose_server_config(config_path.exists())
|
||||
.choose_server_config(config_existed_before_install)
|
||||
.await?
|
||||
{
|
||||
ServerConfigSelection::KeepExisting => {
|
||||
|
|
@ -1714,7 +1719,15 @@ mod tests {
|
|||
.and_then(|s| s.listen.as_ref())
|
||||
.expect("server.listen should be set");
|
||||
let tls = match listen {
|
||||
ServerListenLayer::Tcp { tls, .. } => tls.as_ref().expect("server.listen.tls"),
|
||||
ServerListenLayer::Tcp { address, tls } => {
|
||||
assert_eq!(
|
||||
address
|
||||
.as_ref()
|
||||
.map(fabro_types::settings::InterpString::as_source),
|
||||
Some("127.0.0.1:3000".to_string())
|
||||
);
|
||||
tls.as_ref().expect("server.listen.tls")
|
||||
}
|
||||
ServerListenLayer::Unix { .. } => panic!("expected tcp listen"),
|
||||
};
|
||||
let certs_dir = fabro_util::Home::from_env().certs_dir();
|
||||
|
|
|
|||
|
|
@ -89,11 +89,22 @@ pub(crate) async fn dispatch(
|
|||
}
|
||||
RunCommands::RunWorker(RunWorkerArgs {
|
||||
server,
|
||||
storage_dir,
|
||||
artifact_upload_token,
|
||||
run_dir,
|
||||
run_id,
|
||||
mode,
|
||||
}) => runner::execute(run_id, server, artifact_upload_token, run_dir, mode).await,
|
||||
}) => {
|
||||
runner::execute(
|
||||
run_id,
|
||||
server,
|
||||
storage_dir,
|
||||
artifact_upload_token,
|
||||
run_dir,
|
||||
mode,
|
||||
)
|
||||
.await
|
||||
}
|
||||
RunCommands::Diff(args) => diff::run(args, globals, printer).await,
|
||||
RunCommands::Logs(args) => {
|
||||
let styles = Styles::detect_stdout();
|
||||
|
|
|
|||
|
|
@ -6,11 +6,13 @@ use std::time::Duration;
|
|||
|
||||
use anyhow::{Context, Result, anyhow};
|
||||
use async_trait::async_trait;
|
||||
use fabro_config::Storage;
|
||||
use fabro_interview::{ControlInterviewer, WorkerControlEnvelope, WorkerControlMessage};
|
||||
use fabro_store::{EventEnvelope, EventPayload, RunProjection};
|
||||
use fabro_types::settings::run::RunMode;
|
||||
use fabro_types::settings::{InterpString, SettingsLayer};
|
||||
use fabro_types::{EventBody, RunBlobId, RunEvent, RunId, StatusReason};
|
||||
use fabro_vault::Vault;
|
||||
use fabro_workflow::artifact_snapshot::CapturedArtifactInfo;
|
||||
use fabro_workflow::artifact_upload::{ArtifactSink, StageArtifactUploader};
|
||||
use fabro_workflow::event::{Emitter, RunEventSink};
|
||||
|
|
@ -19,7 +21,7 @@ use fabro_workflow::run_control::RunControlState;
|
|||
use fabro_workflow::runtime_store::{RunStoreBackend, RunStoreHandle};
|
||||
#[cfg(unix)]
|
||||
use tokio::signal::unix::{SignalKind, signal};
|
||||
use tokio::sync::{Mutex, mpsc};
|
||||
use tokio::sync::{Mutex, RwLock as AsyncRwLock, mpsc};
|
||||
use tokio::time::sleep;
|
||||
|
||||
use crate::args::RunWorkerMode;
|
||||
|
|
@ -48,6 +50,7 @@ enum WorkerTitlePhase {
|
|||
pub(crate) async fn execute(
|
||||
run_id: RunId,
|
||||
server: String,
|
||||
storage_dir: Option<PathBuf>,
|
||||
artifact_upload_token: Option<String>,
|
||||
run_dir: PathBuf,
|
||||
mode: RunWorkerMode,
|
||||
|
|
@ -76,6 +79,7 @@ pub(crate) async fn execute(
|
|||
let run_control = RunControlState::new();
|
||||
install_signal_handlers(Arc::clone(&run_control), Arc::clone(&cancel_token))?;
|
||||
let github_app = maybe_build_github_credentials(&run_record.settings).await?;
|
||||
let vault = load_worker_vault(storage_dir.as_deref())?;
|
||||
let services = StartServices {
|
||||
run_id,
|
||||
cancel_token: Some(Arc::clone(&cancel_token)),
|
||||
|
|
@ -92,7 +96,7 @@ pub(crate) async fn execute(
|
|||
artifact_sink,
|
||||
run_control: Some(run_control),
|
||||
github_app,
|
||||
vault: None,
|
||||
vault,
|
||||
on_node: None,
|
||||
registry_override: None,
|
||||
};
|
||||
|
|
@ -109,6 +113,21 @@ pub(crate) async fn execute(
|
|||
Ok(())
|
||||
}
|
||||
|
||||
fn load_worker_vault(storage_dir: Option<&Path>) -> Result<Option<Arc<AsyncRwLock<Vault>>>> {
|
||||
let Some(storage_dir) = storage_dir else {
|
||||
return Ok(None);
|
||||
};
|
||||
|
||||
let storage = Storage::new(storage_dir);
|
||||
let vault = Vault::load(storage.secrets_path()).with_context(|| {
|
||||
format!(
|
||||
"failed to load worker vault from {}",
|
||||
storage.root().display()
|
||||
)
|
||||
})?;
|
||||
Ok(Some(Arc::new(AsyncRwLock::new(vault))))
|
||||
}
|
||||
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
enum WorkerControlStreamEvent {
|
||||
Line(String),
|
||||
|
|
@ -553,18 +572,23 @@ mod tests {
|
|||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicBool, Ordering};
|
||||
|
||||
use fabro_auth::{AuthCredential, AuthDetails};
|
||||
use fabro_config::Storage;
|
||||
use fabro_interview::{AnswerValue, ControlInterviewer, Interviewer, Question, QuestionType};
|
||||
use fabro_model::Provider;
|
||||
use fabro_types::run_event::{
|
||||
InterviewCompletedProps, InterviewStartedProps, RunCompletedProps, RunControlEffectProps,
|
||||
RunFailedProps, RunStatusTransitionProps,
|
||||
};
|
||||
use fabro_types::{EventBody, StatusReason, fixtures};
|
||||
use fabro_vault::{SecretType, Vault};
|
||||
use fabro_workflow::artifact_upload::StageArtifactUploader;
|
||||
|
||||
use super::{
|
||||
MissingArtifactUploadTokenUploader, WorkerControlStreamEvent, WorkerTitlePhase,
|
||||
apply_worker_control_line, handle_worker_control_stream_events, initial_worker_title_phase,
|
||||
read_worker_control_stream_blocking, worker_title, worker_title_phase_for_event,
|
||||
load_worker_vault, read_worker_control_stream_blocking, worker_title,
|
||||
worker_title_phase_for_event,
|
||||
};
|
||||
use crate::args::RunWorkerMode;
|
||||
|
||||
|
|
@ -774,4 +798,31 @@ mod tests {
|
|||
assert_eq!(answer.value, AnswerValue::Interrupted);
|
||||
assert!(!cancel_token.load(Ordering::SeqCst));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn load_worker_vault_reads_credentials_from_storage_dir() {
|
||||
let temp = tempfile::tempdir().unwrap();
|
||||
let storage = Storage::new(temp.path());
|
||||
let mut vault = Vault::load(storage.secrets_path()).unwrap();
|
||||
vault
|
||||
.set(
|
||||
"anthropic",
|
||||
&serde_json::to_string(&AuthCredential {
|
||||
provider: Provider::Anthropic,
|
||||
details: AuthDetails::ApiKey {
|
||||
key: "vault-key".to_string(),
|
||||
},
|
||||
})
|
||||
.unwrap(),
|
||||
SecretType::Credential,
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
let loaded = load_worker_vault(Some(temp.path())).unwrap().unwrap();
|
||||
let guard = loaded.read().await;
|
||||
let credential = guard.get("anthropic").unwrap();
|
||||
|
||||
assert!(credential.contains("vault-key"));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -323,7 +323,8 @@ async fn main_inner() -> (String, Result<()>) {
|
|||
#[cfg(test)]
|
||||
mod tests {
|
||||
use args::{
|
||||
Commands, ModelsCommand, ProviderCommand, ProviderNamespace, StoreCommand, StoreNamespace,
|
||||
Commands, InstallGitHubStrategyArg, ModelsCommand, ProviderCommand, ProviderNamespace,
|
||||
StoreCommand, StoreNamespace,
|
||||
};
|
||||
|
||||
use super::*;
|
||||
|
|
@ -378,6 +379,34 @@ mod tests {
|
|||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_install_non_interactive_accepts_gh_cli_strategy() {
|
||||
let cli = Cli::try_parse_from([
|
||||
"fabro",
|
||||
"install",
|
||||
"--non-interactive",
|
||||
"--llm-provider",
|
||||
"anthropic",
|
||||
"--llm-api-key-env",
|
||||
"ANTHROPIC_API_KEY",
|
||||
"--github-strategy",
|
||||
"gh_cli",
|
||||
"--github-username",
|
||||
"brynary",
|
||||
])
|
||||
.expect("should parse");
|
||||
match *cli.command {
|
||||
Commands::Install(args) => {
|
||||
assert!(args.non_interactive);
|
||||
assert_eq!(
|
||||
args.scripted.github_strategy,
|
||||
Some(InstallGitHubStrategyArg::GhCli)
|
||||
);
|
||||
}
|
||||
_ => panic!("unexpected command variant"),
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn parse_provider_login_missing_provider_flag() {
|
||||
let result = Cli::try_parse_from(["fabro", "provider", "login"]);
|
||||
|
|
|
|||
|
|
@ -1,4 +1,8 @@
|
|||
use fabro_auth::{AuthCredential, AuthDetails};
|
||||
use fabro_config::Storage;
|
||||
use fabro_model::Provider;
|
||||
use fabro_test::{fabro_snapshot, test_context};
|
||||
use fabro_vault::{SecretType, Vault};
|
||||
use httpmock::MockServer;
|
||||
use serde_json::Value;
|
||||
|
||||
|
|
@ -86,6 +90,32 @@ fn run_completed_event(run_id: &str) -> serde_json::Value {
|
|||
})
|
||||
}
|
||||
|
||||
fn seed_anthropic_vault(storage_dir: &std::path::Path, base_url: &str) {
|
||||
let mut vault = Vault::load(Storage::new(storage_dir).secrets_path()).unwrap();
|
||||
vault
|
||||
.set(
|
||||
"anthropic",
|
||||
&serde_json::to_string(&AuthCredential {
|
||||
provider: Provider::Anthropic,
|
||||
details: AuthDetails::ApiKey {
|
||||
key: "vault-anthropic-key".to_string(),
|
||||
},
|
||||
})
|
||||
.unwrap(),
|
||||
SecretType::Credential,
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
vault
|
||||
.set(
|
||||
"ANTHROPIC_BASE_URL",
|
||||
base_url,
|
||||
SecretType::Environment,
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
fn run_running_event(run_id: &str, seq: u32) -> serde_json::Value {
|
||||
serde_json::json!({
|
||||
"seq": seq,
|
||||
|
|
@ -234,6 +264,96 @@ fn detach_uses_configured_server_target_without_server_flag() {
|
|||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn run_uses_vault_credentials_for_worker_execution() {
|
||||
let context = test_context!();
|
||||
let run_id = unique_run_id();
|
||||
let llm_server = MockServer::start();
|
||||
seed_anthropic_vault(
|
||||
&context.storage_dir,
|
||||
&format!("{}/v1", llm_server.base_url()),
|
||||
);
|
||||
|
||||
let llm_mock = llm_server.mock(|when, then| {
|
||||
when.method("POST")
|
||||
.path("/v1/messages")
|
||||
.header("x-api-key", "vault-anthropic-key");
|
||||
then.status(200)
|
||||
.header("Content-Type", "application/json")
|
||||
.body(
|
||||
serde_json::json!({
|
||||
"id": "msg_test_123",
|
||||
"model": "claude-haiku-4-5",
|
||||
"content": [
|
||||
{
|
||||
"type": "text",
|
||||
"text": "Hello from vault"
|
||||
}
|
||||
],
|
||||
"stop_reason": "end_turn",
|
||||
"usage": {
|
||||
"input_tokens": 12,
|
||||
"output_tokens": 4
|
||||
}
|
||||
})
|
||||
.to_string(),
|
||||
);
|
||||
});
|
||||
|
||||
context.write_temp(
|
||||
"vault_worker_llm.fabro",
|
||||
"\
|
||||
digraph VaultWorkerLlm {
|
||||
graph [goal=\"Use a vault-backed Anthropic credential\"]
|
||||
rankdir=LR
|
||||
|
||||
start [shape=Mdiamond, label=\"Start\"]
|
||||
exit [shape=Msquare, label=\"Exit\"]
|
||||
draft [shape=tab, label=\"Draft\", prompt=\"Write a short greeting.\"]
|
||||
|
||||
start -> draft -> exit
|
||||
}
|
||||
",
|
||||
);
|
||||
|
||||
let output = context
|
||||
.run_cmd()
|
||||
.env_remove("ANTHROPIC_API_KEY")
|
||||
.env_remove("ANTHROPIC_BASE_URL")
|
||||
.env_remove("OPENAI_API_KEY")
|
||||
.env_remove("OPENAI_BASE_URL")
|
||||
.env_remove("GEMINI_API_KEY")
|
||||
.args([
|
||||
"--run-id",
|
||||
run_id.as_str(),
|
||||
"--auto-approve",
|
||||
"--no-retro",
|
||||
"--sandbox",
|
||||
"local",
|
||||
"--provider",
|
||||
"anthropic",
|
||||
"--model",
|
||||
"claude-haiku-4-5",
|
||||
context
|
||||
.temp_dir
|
||||
.join("vault_worker_llm.fabro")
|
||||
.to_str()
|
||||
.unwrap(),
|
||||
])
|
||||
.output()
|
||||
.expect("command should execute");
|
||||
|
||||
assert!(
|
||||
output.status.success(),
|
||||
"command failed:\nstdout:\n{}\nstderr:\n{}",
|
||||
String::from_utf8_lossy(&output.stdout),
|
||||
String::from_utf8_lossy(&output.stderr)
|
||||
);
|
||||
|
||||
llm_mock.assert();
|
||||
wait_for_event_names(&context.find_run_dir(&run_id), &["run.completed"]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn detach_rejects_storage_dir_flag() {
|
||||
let context = test_context!();
|
||||
|
|
|
|||
|
|
@ -3139,6 +3139,8 @@ fn worker_command(
|
|||
cmd.arg("__run-worker")
|
||||
.arg("--server")
|
||||
.arg(server_target)
|
||||
.arg("--storage-dir")
|
||||
.arg(&storage_dir)
|
||||
.arg("--artifact-upload-token")
|
||||
.arg(artifact_upload_token)
|
||||
.arg("--run-dir")
|
||||
|
|
@ -6354,10 +6356,53 @@ type = "http"
|
|||
|
||||
let response = app.oneshot(req).await.unwrap();
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
assert!(state.vault.read().await.list().is_empty());
|
||||
let listed = state.vault.read().await.list();
|
||||
assert_eq!(listed.len(), 1);
|
||||
assert_eq!(listed[0].name, "openai_codex");
|
||||
assert_eq!(listed[0].secret_type, SecretType::Credential);
|
||||
assert!(state.vault.read().await.get("openai_codex").is_some());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn list_secrets_includes_credential_metadata() {
|
||||
let state = create_app_state();
|
||||
{
|
||||
let mut vault = state.vault.write().await;
|
||||
vault
|
||||
.set(
|
||||
"anthropic",
|
||||
"{\"provider\":\"anthropic\"}",
|
||||
SecretType::Credential,
|
||||
Some("saved auth"),
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
let app = build_router(Arc::clone(&state), AuthMode::Disabled);
|
||||
|
||||
let response = app
|
||||
.oneshot(
|
||||
Request::builder()
|
||||
.method("GET")
|
||||
.uri(api("/secrets"))
|
||||
.body(Body::empty())
|
||||
.unwrap(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(response.status(), StatusCode::OK);
|
||||
let body = body_json(response.into_body()).await;
|
||||
let data = body["data"].as_array().expect("data should be an array");
|
||||
let entry = data
|
||||
.iter()
|
||||
.find(|entry| entry["name"] == "anthropic")
|
||||
.expect("credential metadata should be listed");
|
||||
assert_eq!(entry["type"], "credential");
|
||||
assert_eq!(entry["description"], "saved auth");
|
||||
assert!(entry.get("updated_at").is_some());
|
||||
assert!(entry.get("value").is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn create_secret_rejects_invalid_credential_json() {
|
||||
let state = create_app_state();
|
||||
|
|
|
|||
|
|
@ -135,7 +135,6 @@ impl Vault {
|
|||
let mut data = self
|
||||
.entries
|
||||
.iter()
|
||||
.filter(|(_, entry)| entry.secret_type != SecretType::Credential)
|
||||
.map(|(name, entry)| SecretMetadata {
|
||||
name: name.clone(),
|
||||
secret_type: entry.secret_type,
|
||||
|
|
@ -351,7 +350,7 @@ mod tests {
|
|||
}
|
||||
|
||||
#[test]
|
||||
fn list_hides_credential_entries_loaded_from_disk() {
|
||||
fn list_includes_credential_entries_loaded_from_disk() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let path = dir.path().join("secrets.json");
|
||||
std::fs::write(
|
||||
|
|
@ -376,8 +375,10 @@ mod tests {
|
|||
|
||||
let store = Vault::load(path).unwrap();
|
||||
|
||||
assert_eq!(store.list().len(), 1);
|
||||
assert_eq!(store.list().len(), 2);
|
||||
assert_eq!(store.list()[0].name, "OPENAI_API_KEY");
|
||||
assert_eq!(store.list()[1].name, "openai_codex");
|
||||
assert_eq!(store.list()[1].secret_type, SecretType::Credential);
|
||||
assert_eq!(store.get("openai_codex"), Some("{\"provider\":\"openai\"}"));
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue