Simplify SQLite stores after review

- Delete the dead test-only Vault-based env-secrets migration and point
  the startup migration tests at the production migrate_to_store path
  over a real SQLite-backed SecretStore
- Extract shared legacy-import helpers (timestamped backup rename,
  is_toml_file) into fabro_db::legacy and parse_rfc3339_utc into
  fabro-db, replacing four per-crate copies
- Take one secrets snapshot in migrate_to_store instead of per-name
  queries
- Share one bind order between the MCP store INSERT and UPDATE
  statements
- Return SecretEntry directly from entry_from_row
- Unify the environment/MCP store blocking loaders into a generic
  load_store_blocking helper

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-07-11 14:55:21 -04:00 • committed by Scott Werner
parent ec3933d5de
commit 431399826d
13 changed files with 260 additions and 371 deletions

1
Cargo.lock generated
View file

@ -2561,6 +2561,7 @@ name = "fabro-db"
version = "0.302.0-nightly.1"
dependencies = [
"anyhow",
"chrono",
"sqlx",
"tempfile",
"tokio",

View file

@ -14,6 +14,7 @@ workspace = true
[dependencies]
anyhow.workspace = true
chrono.workspace = true
sqlx.workspace = true
tokio.workspace = true
tracing.workspace = true

View file

@ -0,0 +1,79 @@
//! Helpers shared by the one-time imports that seed SQLite tables from
//! pre-SQLite on-disk stores.
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use chrono::{DateTime, Utc};
use tokio::fs;
/// Failure to move an imported legacy source aside, carrying the backup path
/// the caller needs for its own error variant.
#[derive(Debug)]
pub struct LegacyBackupError {
pub backup_path: PathBuf,
pub source: std::io::Error,
}
/// Move an imported legacy file or directory aside to
/// `<name>.imported-<timestamp>.bak` next to the original. `fallback_name` is
/// used when the source path has no final component.
pub async fn rename_to_legacy_backup(
source: &Path,
fallback_name: &str,
) -> Result<PathBuf, LegacyBackupError> {
let backup_path = legacy_backup_path(source, fallback_name, Utc::now());
fs::rename(source, &backup_path)
.await
.map_err(|source| LegacyBackupError {
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}
fn legacy_backup_path(source: &Path, fallback_name: &str, imported_at: DateTime<Utc>) -> PathBuf {
let timestamp = imported_at.format("%Y%m%dT%H%M%S%fZ");
let mut file_name = source
.file_name()
.map_or_else(|| OsString::from(fallback_name), OsString::from);
file_name.push(format!(".imported-{timestamp}.bak"));
source.with_file_name(file_name)
}
/// True when `path` has a `.toml` extension (legacy per-item store files).
pub fn is_toml_file(path: &Path) -> bool {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| extension == "toml")
}
#[cfg(test)]
mod tests {
use chrono::TimeZone as _;
use super::*;
#[test]
fn backup_path_appends_timestamped_suffix() {
let imported_at = Utc.with_ymd_and_hms(2026, 7, 11, 1, 2, 3).unwrap();
let backup =
legacy_backup_path(Path::new("/data/secrets.json"), "secrets.json", imported_at);
assert_eq!(
backup,
Path::new("/data/secrets.json.imported-20260711T010203000000000Z.bak")
);
}
#[test]
fn backup_path_uses_fallback_when_source_has_no_file_name() {
let imported_at = Utc.with_ymd_and_hms(2026, 7, 11, 1, 2, 3).unwrap();
let backup = legacy_backup_path(Path::new("/"), "mcps", imported_at);
assert!(
backup
.file_name()
.and_then(|name| name.to_str())
.is_some_and(|name| name.starts_with("mcps.imported-"))
);
}
}

View file

@ -3,6 +3,7 @@ use std::path::{Path, PathBuf};
use std::time::Duration;
use anyhow::Context as _;
use chrono::{DateTime, Utc};
use sqlx::migrate::{Migrate as _, Migrator};
use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions, SqliteSynchronous};
use tokio::fs;
@ -10,8 +11,15 @@ use tokio::fs;
use tokio::task::spawn_blocking;
use tracing::info;
pub mod legacy;
pub type DbPool = sqlx::SqlitePool;
/// Parse an RFC 3339 TEXT column value into a UTC timestamp.
pub fn parse_rfc3339_utc(value: &str) -> Result<DateTime<Utc>, chrono::ParseError> {
DateTime::parse_from_rfc3339(value).map(|timestamp| timestamp.with_timezone(&Utc))
}
static MIGRATOR: Migrator = sqlx::migrate!("./migrations");
#[derive(Clone)]

View file

@ -1,15 +1,13 @@
use std::collections::{BTreeMap, HashMap, HashSet};
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use std::str::FromStr;
use std::sync::Arc;
use chrono::{DateTime, Utc};
use fabro_config::{
EnvironmentDockerfileLayer, EnvironmentImageLayer, EnvironmentLayer, EnvironmentLifecycleLayer,
EnvironmentNetworkLayer, EnvironmentResourcesLayer, MergeMap, StickyMap,
};
use fabro_db::DbPool;
use fabro_db::{DbPool, legacy};
use fabro_types::settings::run::{DockerfileSource, EnvironmentProvider, EnvironmentSettings};
use fabro_types::settings::{Duration, InterpString, Size};
use serde::de::DeserializeOwned;
@ -706,7 +704,7 @@ async fn legacy_environment_paths(
.file_type()
.await
.map_err(|source| EnvironmentStoreError::io(&path, source))?;
if file_type.is_file() && is_toml_file(&path) {
if file_type.is_file() && legacy::is_toml_file(&path) {
paths.push(LegacyEnvironmentPath {
id: id_from_path(&path)?,
path,
@ -756,20 +754,9 @@ async fn existing_environment_ids(
async fn rename_imported_legacy_directory(
source_dir: &Path,
) -> Result<PathBuf, EnvironmentStoreError> {
let backup_path = legacy_backup_path(source_dir, Utc::now());
fs::rename(source_dir, &backup_path)
legacy::rename_to_legacy_backup(source_dir, "environments")
.await
.map_err(|source| EnvironmentStoreError::io(&backup_path, source))?;
Ok(backup_path)
}
fn legacy_backup_path(source_dir: &Path, imported_at: DateTime<Utc>) -> PathBuf {
let timestamp = imported_at.format("%Y%m%dT%H%M%S%fZ");
let mut file_name = source_dir
.file_name()
.map_or_else(|| OsString::from("environments"), OsString::from);
file_name.push(format!(".imported-{timestamp}.bak"));
source_dir.with_file_name(file_name)
.map_err(|err| EnvironmentStoreError::io(&err.backup_path, err.source))
}
fn id_from_path(path: &Path) -> Result<EnvironmentId, EnvironmentStoreError> {
@ -786,12 +773,6 @@ fn id_from_path(path: &Path) -> Result<EnvironmentId, EnvironmentStoreError> {
})
}
fn is_toml_file(path: &Path) -> bool {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| extension == "toml")
}
fn encode_json<T: serde::Serialize>(
field: &'static str,
value: &T,

View file

@ -1,12 +1,10 @@
use std::collections::{BTreeMap, HashMap};
use std::ffi::OsString;
use std::io::ErrorKind;
use std::path::{Path, PathBuf};
use std::str::FromStr as _;
use std::sync::RwLock;
use chrono::{DateTime, Utc};
use fabro_db::DbPool;
use fabro_db::{DbPool, legacy};
use fabro_types::settings::run::{McpHttpProtocol, McpServerSettings, McpTransport};
use fabro_types::{
McpServerDefinition, McpServerDraft, McpServerId, McpServerReplace, McpServerRevision,
@ -29,6 +27,8 @@ use crate::model;
/// Reads use a synchronous in-memory catalog because manifest resolution is
/// synchronous. Mutations are serialized within this process and use
/// revision-guarded SQL so SQLite remains authoritative for concurrency.
/// Writes to the same database from other processes are not observed until
/// this store is reloaded.
pub struct McpServerStore {
pool: DbPool,
mutations: Mutex<()>,
@ -444,20 +444,7 @@ async fn update_definition(
expected: &McpServerRevision,
) -> Result<(), McpServerStoreError> {
let row = McpServerSqlRow::from_definition(definition)?;
let result = sqlx::query(UPDATE_DEFINITION_SQL)
.bind(row.revision)
.bind(row.display_name)
.bind(row.description)
.bind(row.transport_type)
.bind(row.protocol)
.bind(row.command_json)
.bind(row.url)
.bind(row.port)
.bind(row.env_json)
.bind(row.headers_json)
.bind(row.startup_timeout_secs)
.bind(row.tool_timeout_secs)
.bind(row.id)
let result = bind_definition(sqlx::query(UPDATE_DEFINITION_SQL), row)
.bind(expected.as_str())
.execute(&mut **transaction)
.await?;
@ -467,12 +454,13 @@ async fn update_definition(
Ok(())
}
/// Binds the shared placeholder order of [`INSERT_DEFINITION_SQL`] and
/// [`UPDATE_DEFINITION_SQL`]: the twelve value columns first, then `id`.
fn bind_definition(
query: Query<'_, Sqlite, SqliteArguments>,
row: McpServerSqlRow,
) -> Query<'_, Sqlite, SqliteArguments> {
query
.bind(row.id)
.bind(row.revision)
.bind(row.display_name)
.bind(row.description)
@ -485,6 +473,7 @@ fn bind_definition(
.bind(row.headers_json)
.bind(row.startup_timeout_secs)
.bind(row.tool_timeout_secs)
.bind(row.id)
}
struct McpServerSqlRow {
@ -590,7 +579,6 @@ fn encode_string_map(
const INSERT_DEFINITION_SQL: &str = r"
INSERT INTO mcp_servers (
id,
revision,
display_name,
description,
@ -602,7 +590,8 @@ INSERT INTO mcp_servers (
env_json,
headers_json,
startup_timeout_secs,
tool_timeout_secs
tool_timeout_secs,
id
)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
ON CONFLICT(id) DO NOTHING
@ -686,7 +675,7 @@ async fn legacy_definition_paths(
.file_type()
.await
.map_err(|source| McpServerStoreError::io(&path, source))?;
if file_type.is_file() && is_toml_file(&path) {
if file_type.is_file() && legacy::is_toml_file(&path) {
paths.push((id_from_path(&path)?, path));
}
}
@ -722,31 +711,14 @@ fn id_from_path(path: &Path) -> Result<McpServerId, McpServerStoreError> {
})
}
fn is_toml_file(path: &Path) -> bool {
path.extension()
.and_then(|extension| extension.to_str())
.is_some_and(|extension| extension == "toml")
}
async fn rename_imported_legacy_directory(
source_dir: &Path,
) -> Result<PathBuf, McpServerStoreError> {
let backup_path = legacy_backup_path(source_dir, Utc::now());
fs::rename(source_dir, &backup_path)
legacy::rename_to_legacy_backup(source_dir, "mcps")
.await
.map_err(|source| McpServerStoreError::LegacyBackup {
.map_err(|err| McpServerStoreError::LegacyBackup {
source_path: source_dir.to_path_buf(),
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}
fn legacy_backup_path(source_dir: &Path, imported_at: DateTime<Utc>) -> PathBuf {
let timestamp = imported_at.format("%Y%m%dT%H%M%S%fZ");
let mut file_name = source_dir
.file_name()
.map_or_else(|| OsString::from("mcps"), OsString::from);
file_name.push(format!(".imported-{timestamp}.bak"));
source_dir.with_file_name(file_name)
backup_path: err.backup_path,
source: err.source,
})
}

View file

@ -15,8 +15,6 @@ use std::path::{Path, PathBuf};
use anyhow::Context as _;
use fabro_config::envfile::{self, EnvFileRemoval};
use fabro_static::{EnvVars, optional_vault_secrets};
#[cfg(test)]
use fabro_vault::Vault;
use fabro_vault::{SecretStore, SecretStoreWrite, SecretType};
pub(crate) const REMOVAL_DEADLINE: &str = "2026-08-18";
@ -36,94 +34,6 @@ impl OptionalServerEnvSecretsMigrationReport {
}
}
#[cfg(test)]
pub(crate) fn migrate(
vault: &mut Vault,
server_env_path: &Path,
env_entries: &HashMap<String, String>,
) -> anyhow::Result<OptionalServerEnvSecretsMigrationReport> {
let server_env_entries = envfile::read_env_file(server_env_path)
.with_context(|| format!("read server env file {}", server_env_path.display()))?;
let mut vault_writes = Vec::new();
let mut env_removals = Vec::new();
let mut warnings = Vec::new();
let mut preserved_env_entries = 0;
for &name in optional_vault_secrets() {
let process_value = env_entries.get(name);
let file_value = server_env_entries.get(name);
if let Some(vault_value) = vault.get(name) {
if let Some(file_value) = file_value {
if file_value == vault_value {
env_removals.push(env_removal(name));
} else {
preserved_env_entries += 1;
warnings.push(format!(
"Preserved {name} in server.env because the vault already contains a different value"
));
}
}
continue;
}
match (process_value, file_value) {
(Some(value), Some(file_value)) => {
vault_writes.push((name, value.clone(), secret_type_for(name)));
if value == file_value {
env_removals.push(env_removal(name));
} else {
preserved_env_entries += 1;
warnings.push(format!(
"Preserved {name} in server.env because process env takes precedence and the file value differs"
));
}
}
(Some(value), None) => {
vault_writes.push((name, value.clone(), secret_type_for(name)));
}
(None, Some(value)) => {
vault_writes.push((name, value.clone(), secret_type_for(name)));
env_removals.push(env_removal(name));
}
(None, None) => {}
}
}
let mut report = OptionalServerEnvSecretsMigrationReport {
migrated_secrets: vault_writes.len(),
removed_env_entries: 0,
preserved_env_entries,
backup_path: None,
warnings,
};
if vault_writes.is_empty() && env_removals.is_empty() {
return Ok(report);
}
for (name, value, secret_type) in vault_writes {
vault
.set(name, &value, secret_type, None)
.with_context(|| format!("write migrated secret {name} to vault"))?;
}
if !env_removals.is_empty() {
let backup_path = backup_server_env_file(server_env_path)?;
let update_report =
envfile::update_env_file_with_report(server_env_path, env_removals, Vec::new())
.with_context(|| {
format!(
"remove migrated optional secrets from {}",
server_env_path.display()
)
})?;
report.removed_env_entries = update_report.removed_keys.len();
report.backup_path = Some(backup_path);
}
Ok(report)
}
pub(crate) async fn migrate_to_store(
store: &SecretStore,
server_env_path: &Path,
@ -131,6 +41,7 @@ pub(crate) async fn migrate_to_store(
) -> anyhow::Result<OptionalServerEnvSecretsMigrationReport> {
let server_env_entries = envfile::read_env_file(server_env_path)
.with_context(|| format!("read server env file {}", server_env_path.display()))?;
let stored = store.snapshot().await?;
let mut writes = Vec::new();
let mut env_removals = Vec::new();
let mut warnings = Vec::new();
@ -139,7 +50,7 @@ pub(crate) async fn migrate_to_store(
for &name in optional_vault_secrets() {
let process_value = env_entries.get(name);
let file_value = server_env_entries.get(name);
if let Some(entry) = store.get(name).await? {
if let Some(entry) = stored.get_entry(name) {
if let Some(file_value) = file_value {
if file_value == &entry.value {
env_removals.push(env_removal(name));

View file

@ -2,8 +2,6 @@ use std::collections::HashMap;
use std::path::Path;
use fabro_vault::SecretStore;
#[cfg(test)]
use fabro_vault::Vault;
#[path = "../migrations/2026051801_legacy_vault_entries.rs"]
mod legacy_vault_entries;
@ -21,15 +19,6 @@ pub(crate) fn migrate_legacy_vault_file(path: &Path) -> anyhow::Result<LegacyVau
legacy_vault_entries::migrate_legacy_vault_file(path)
}
#[cfg(test)]
pub(crate) fn migrate_optional_server_env_secrets_to_vault(
vault: &mut Vault,
server_env_path: &Path,
env_entries: &HashMap<String, String>,
) -> anyhow::Result<OptionalServerEnvSecretsMigrationReport> {
optional_server_env_secrets_to_vault::migrate(vault, server_env_path, env_entries)
}
pub(crate) async fn migrate_optional_server_env_secrets_to_store(
store: &SecretStore,
server_env_path: &Path,

View file

@ -2328,43 +2328,21 @@ fn mcp_server_dir_for_active_config(active_config_path: &std::path::Path) -> Pat
reason = "synchronous app-state assembly may run inside an async runtime; a short-lived OS \
thread avoids nested Tokio runtimes"
)]
fn load_environment_store_blocking(
pool: DbPool,
local_enabled: bool,
) -> anyhow::Result<EnvironmentStore> {
fn load_store_blocking<T, F, Fut>(description: &'static str, load: F) -> anyhow::Result<T>
where
T: Send + 'static,
F: FnOnce() -> Fut + Send + 'static,
Fut: std::future::Future<Output = anyhow::Result<T>>,
{
std::thread::spawn(move || {
let runtime = TokioRuntimeBuilder::new_current_thread()
.enable_all()
.build()
.context("build environment store runtime")?;
runtime
.block_on(EnvironmentStore::load(pool, local_enabled))
.map_err(anyhow::Error::new)
.with_context(|| format!("build {description} runtime"))?;
runtime.block_on(load())
})
.join()
.expect("environment store load thread should not panic")
}
#[expect(
clippy::disallowed_methods,
reason = "synchronous app-state assembly may run inside an async runtime; a short-lived OS \
thread avoids nested Tokio runtimes"
)]
fn load_mcp_server_store_blocking(
pool: DbPool,
legacy_dir: PathBuf,
) -> anyhow::Result<McpServerStore> {
std::thread::spawn(move || {
let runtime = TokioRuntimeBuilder::new_current_thread()
.enable_all()
.build()
.context("build MCP server store runtime")?;
runtime
.block_on(McpServerStore::open(pool, legacy_dir))
.map_err(anyhow::Error::new)
})
.join()
.expect("MCP server store load thread should not panic")
.expect("store load thread should not panic")
}
pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppState>> {
@ -2404,14 +2382,24 @@ pub(crate) fn build_app_state(config: AppStateConfig) -> anyhow::Result<Arc<AppS
.providers
.local
.enabled;
let environment_pool = db_pool.clone();
let environment_store = Arc::new(
load_environment_store_blocking(db_pool.clone(), local_provider_enabled)
.context("load environments")?,
load_store_blocking("environment store", move || async move {
EnvironmentStore::load(environment_pool, local_provider_enabled)
.await
.map_err(anyhow::Error::new)
})
.context("load environments")?,
);
let mcp_server_dir = mcp_server_dir_for_active_config(&active_config_path);
let mcp_server_pool = db_pool.clone();
let mcp_server_store = Arc::new(
load_mcp_server_store_blocking(db_pool.clone(), mcp_server_dir)
.context("load mcp servers")?,
load_store_blocking("MCP server store", move || async move {
McpServerStore::open(mcp_server_pool, mcp_server_dir)
.await
.map_err(anyhow::Error::new)
})
.context("load mcp servers")?,
);
let variables = Arc::new(VariableStore::new(db_pool.clone()));
let secret_store = Arc::new(SecretStore::new(db_pool));

View file

@ -112,8 +112,12 @@ fn test_environment_store(
})
.join()
.expect("environment store setup thread should not panic");
let store = load_environment_store_blocking(pool, local_enabled)
.expect("test environment store should load");
let store = load_store_blocking("environment store", move || async move {
EnvironmentStore::load(pool, local_enabled)
.await
.map_err(anyhow::Error::new)
})
.expect("test environment store should load");
(temp, store)
}
@ -138,8 +142,13 @@ fn test_mcp_server_store() -> (tempfile::TempDir, McpServerStore) {
})
.join()
.expect("MCP server store setup thread should not panic");
let store = load_mcp_server_store_blocking(pool, temp.path().join("mcps"))
.expect("test MCP server store should load");
let mcps_dir = temp.path().join("mcps");
let store = load_store_blocking("MCP server store", move || async move {
McpServerStore::open(pool, mcps_dir)
.await
.map_err(anyhow::Error::new)
})
.expect("test MCP server store should load");
(temp, store)
}

View file

@ -1,8 +1,6 @@
use std::collections::HashMap;
use std::path::Path;
#[cfg(test)]
use anyhow::Context as _;
use fabro_static::EnvVars;
use fabro_types::settings::ServerNamespace;
use fabro_vault::Vault;
@ -54,54 +52,6 @@ pub fn migrate_startup_vault(vault_path: impl AsRef<Path>) {
}
}
#[cfg(test)]
fn load_startup_vault(vault_path: impl AsRef<Path>) -> anyhow::Result<Vault> {
let vault_path = vault_path.as_ref();
migrate_startup_vault(vault_path);
Vault::load(vault_path.to_path_buf())
.with_context(|| format!("load vault {}", vault_path.display()))
}
#[cfg(test)]
pub(crate) fn prepare_startup_vault(
vault_path: impl AsRef<Path>,
server_env_path: impl AsRef<Path>,
env_entries: &HashMap<String, String>,
) -> anyhow::Result<Vault> {
let mut vault = load_startup_vault(vault_path)?;
let report = migrations::migrate_optional_server_env_secrets_to_vault(
&mut vault,
server_env_path.as_ref(),
env_entries,
)
.context("migrate optional server env secrets into vault")?;
for warning in &report.warnings {
warn!(
warning = %warning,
removal_deadline = migrations::OPTIONAL_SERVER_ENV_SECRETS_REMOVAL_DEADLINE,
"Optional server env secrets migration warning"
);
}
if report.changed() {
let backup_path = report
.backup_path
.as_ref()
.map_or_else(|| "<none>".to_string(), |path| path.display().to_string());
warn!(
migrated_secrets = report.migrated_secrets,
removed_env_entries = report.removed_env_entries,
preserved_env_entries = report.preserved_env_entries,
backup_path = %backup_path,
removal_deadline = migrations::OPTIONAL_SERVER_ENV_SECRETS_REMOVAL_DEADLINE,
"Migrated optional server env secrets into vault"
);
}
Ok(vault)
}
pub fn validate_startup(
env_path: &Path,
env_entries: HashMap<String, String>,
@ -123,9 +73,10 @@ mod tests {
use fabro_config::{ServerSettingsBuilder, envfile};
use fabro_static::EnvVars;
use fabro_types::settings::ServerNamespace;
use fabro_vault::{SecretType, Vault};
use fabro_vault::{SecretStore, SecretType, Vault};
use super::{prepare_startup_vault, validate_startup};
use super::validate_startup;
use crate::migrations;
fn resolved_settings(auth_methods: &[&str]) -> ServerNamespace {
ServerSettingsBuilder::from_toml(&format!(
@ -159,8 +110,12 @@ client_id = "Iv1.test"
dir.path().join("server.env")
}
fn vault_path(dir: &tempfile::TempDir) -> PathBuf {
dir.path().join("secrets.json")
async fn test_secret_store(dir: &tempfile::TempDir) -> SecretStore {
let database = fabro_db::Database::connect(dir.path().join("fabro.db"))
.await
.unwrap();
database.migrate().await.unwrap();
SecretStore::new(database.clone_pool())
}
#[expect(
@ -280,8 +235,8 @@ client_id = "Iv1.test"
.expect("github client secret in vault should satisfy startup");
}
#[test]
fn prepare_startup_vault_migrates_server_env_optional_secrets_to_vault() {
#[tokio::test]
async fn migrate_optional_secrets_moves_server_env_secrets_to_store() {
let dir = tempfile::tempdir().unwrap();
let server_env_path = env_path(&dir);
envfile::write_env_file(
@ -303,33 +258,36 @@ client_id = "Iv1.test"
]),
)
.unwrap();
let store = test_secret_store(&dir).await;
let vault = prepare_startup_vault(vault_path(&dir), &server_env_path, &HashMap::new())
.expect("legacy optional secrets should migrate");
migrations::migrate_optional_server_env_secrets_to_store(
&store,
&server_env_path,
&HashMap::new(),
)
.await
.expect("legacy optional secrets should migrate");
assert_eq!(
vault.get(EnvVars::GITHUB_APP_CLIENT_SECRET),
Some("legacy-client-secret")
);
assert_eq!(
vault
.get_entry(EnvVars::GITHUB_APP_CLIENT_SECRET)
.unwrap()
.secret_type,
SecretType::Token
);
assert_eq!(
vault.get(EnvVars::GITHUB_APP_PRIVATE_KEY),
Some("legacy-private-key")
);
assert_eq!(
vault
.get_entry(EnvVars::GITHUB_APP_PRIVATE_KEY)
.unwrap()
.secret_type,
SecretType::File
);
assert_eq!(vault.get(EnvVars::OPENAI_API_KEY), Some("sk-legacy"));
let client_secret = store
.get(EnvVars::GITHUB_APP_CLIENT_SECRET)
.await
.unwrap()
.expect("client secret should be stored");
assert_eq!(client_secret.value, "legacy-client-secret");
assert_eq!(client_secret.secret_type, SecretType::Token);
let private_key = store
.get(EnvVars::GITHUB_APP_PRIVATE_KEY)
.await
.unwrap()
.expect("private key should be stored");
assert_eq!(private_key.value, "legacy-private-key");
assert_eq!(private_key.secret_type, SecretType::File);
let openai_key = store
.get(EnvVars::OPENAI_API_KEY)
.await
.unwrap()
.expect("openai key should be stored");
assert_eq!(openai_key.value, "sk-legacy");
let server_env = envfile::read_env_file(&server_env_path).unwrap();
assert!(server_env.contains_key(EnvVars::SESSION_SECRET));
@ -339,8 +297,8 @@ client_id = "Iv1.test"
assert_eq!(migration_backups(dir.path()).len(), 1);
}
#[test]
fn prepare_startup_vault_prefers_process_env_and_preserves_conflicting_server_env() {
#[tokio::test]
async fn migrate_optional_secrets_prefers_process_env_and_preserves_conflicting_server_env() {
let dir = tempfile::tempdir().unwrap();
let server_env_path = env_path(&dir);
envfile::write_env_file(
@ -355,14 +313,22 @@ client_id = "Iv1.test"
EnvVars::GITHUB_APP_CLIENT_SECRET.to_string(),
"process-client-secret".to_string(),
)]);
let store = test_secret_store(&dir).await;
let vault = prepare_startup_vault(vault_path(&dir), &server_env_path, &env_entries)
.expect("process env secret should migrate");
migrations::migrate_optional_server_env_secrets_to_store(
&store,
&server_env_path,
&env_entries,
)
.await
.expect("process env secret should migrate");
assert_eq!(
vault.get(EnvVars::GITHUB_APP_CLIENT_SECRET),
Some("process-client-secret")
);
let client_secret = store
.get(EnvVars::GITHUB_APP_CLIENT_SECRET)
.await
.unwrap()
.expect("client secret should be stored");
assert_eq!(client_secret.value, "process-client-secret");
let server_env = envfile::read_env_file(&server_env_path).unwrap();
assert_eq!(
server_env
@ -373,8 +339,9 @@ client_id = "Iv1.test"
assert!(migration_backups(dir.path()).is_empty());
}
#[test]
fn prepare_startup_vault_keeps_existing_vault_secret_and_removes_matching_server_env() {
#[tokio::test]
async fn migrate_optional_secrets_keeps_existing_stored_secret_and_removes_matching_server_env()
{
let dir = tempfile::tempdir().unwrap();
let server_env_path = env_path(&dir);
envfile::write_env_file(
@ -385,30 +352,38 @@ client_id = "Iv1.test"
)]),
)
.unwrap();
let mut vault = Vault::load(vault_path(&dir)).unwrap();
vault
let store = test_secret_store(&dir).await;
store
.set(
EnvVars::GITHUB_APP_CLIENT_SECRET,
"vault-client-secret",
SecretType::Token,
None,
)
.await
.unwrap();
let vault = prepare_startup_vault(vault_path(&dir), &server_env_path, &HashMap::new())
.expect("redundant server env secret should be cleaned up");
migrations::migrate_optional_server_env_secrets_to_store(
&store,
&server_env_path,
&HashMap::new(),
)
.await
.expect("redundant server env secret should be cleaned up");
assert_eq!(
vault.get(EnvVars::GITHUB_APP_CLIENT_SECRET),
Some("vault-client-secret")
);
let client_secret = store
.get(EnvVars::GITHUB_APP_CLIENT_SECRET)
.await
.unwrap()
.expect("client secret should be stored");
assert_eq!(client_secret.value, "vault-client-secret");
let server_env = envfile::read_env_file(&server_env_path).unwrap();
assert!(!server_env.contains_key(EnvVars::GITHUB_APP_CLIENT_SECRET));
assert_eq!(migration_backups(dir.path()).len(), 1);
}
#[test]
fn prepare_startup_vault_migrated_github_client_secret_satisfies_startup() {
#[tokio::test]
async fn migrate_optional_secrets_migrated_github_client_secret_satisfies_startup() {
let dir = tempfile::tempdir().unwrap();
let server_env_path = env_path(&dir);
envfile::write_env_file(
@ -426,10 +401,17 @@ client_id = "Iv1.test"
)
.unwrap();
let settings = resolved_settings(&["github"]);
let store = test_secret_store(&dir).await;
let vault = prepare_startup_vault(vault_path(&dir), &server_env_path, &HashMap::new())
.expect("legacy github client secret should migrate");
migrations::migrate_optional_server_env_secrets_to_store(
&store,
&server_env_path,
&HashMap::new(),
)
.await
.expect("legacy github client secret should migrate");
let vault = store.snapshot().await.unwrap().into_vault();
validate_startup(&server_env_path, HashMap::new(), &settings, &vault)
.expect("migrated github client secret should satisfy startup");
}

View file

@ -1,9 +1,8 @@
use std::collections::HashMap;
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use chrono::{DateTime, Utc};
use fabro_db::DbPool;
use fabro_db::{DbPool, legacy};
use fabro_types::{Variable, is_env_style_name};
use sqlx::Row as _;
use sqlx::sqlite::SqliteRow;
@ -313,24 +312,13 @@ fn parse_legacy_entries(
}
async fn rename_imported_legacy_file(source_path: &Path) -> Result<PathBuf, Error> {
let backup_path = legacy_backup_path(source_path, Utc::now());
fs::rename(source_path, &backup_path)
legacy::rename_to_legacy_backup(source_path, "variables.json")
.await
.map_err(|source| Error::LegacyBackup {
.map_err(|err| Error::LegacyBackup {
source_path: source_path.to_path_buf(),
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}
fn legacy_backup_path(source_path: &Path, imported_at: DateTime<Utc>) -> PathBuf {
let timestamp = imported_at.format("%Y%m%dT%H%M%S%fZ");
let mut file_name = source_path
.file_name()
.map_or_else(|| OsString::from("variables.json"), OsString::from);
file_name.push(format!(".imported-{timestamp}.bak"));
source_path.with_file_name(file_name)
backup_path: err.backup_path,
source: err.source,
})
}
fn row_count(count: usize) -> Result<i64, Error> {
@ -354,11 +342,9 @@ fn variable_from_row(row: &SqliteRow) -> Result<Variable, Error> {
}
fn parse_timestamp(name: &str, column: &'static str, value: &str) -> Result<DateTime<Utc>, Error> {
DateTime::parse_from_rfc3339(value)
.map(|timestamp| timestamp.with_timezone(&Utc))
.map_err(|source| Error::Timestamp {
name: name.to_string(),
column,
source,
})
fabro_db::parse_rfc3339_utc(value).map_err(|source| Error::Timestamp {
name: name.to_string(),
column,
source,
})
}

View file

@ -1,10 +1,9 @@
use std::collections::HashMap;
use std::ffi::OsString;
use std::path::{Path, PathBuf};
use std::str::FromStr as _;
use chrono::{DateTime, Utc};
use fabro_db::{Database, DbPool};
use fabro_db::{Database, DbPool, legacy};
use fabro_types::{OAuthCredential, SecretMetadata, SecretType};
use sqlx::sqlite::SqliteRow;
use sqlx::{Row as _, Sqlite, Transaction};
@ -177,10 +176,7 @@ impl SecretStore {
.bind(name)
.fetch_optional(&self.pool)
.await?;
row.as_ref()
.map(entry_from_row)
.transpose()
.map(|entry| entry.map(|(_, entry)| entry))
row.as_ref().map(entry_from_row).transpose()
}
pub async fn list(&self) -> Result<Vec<SecretMetadata>, SecretStoreError> {
@ -288,7 +284,7 @@ impl SecretStore {
.await?;
if let Some(row) = row {
return entry_from_row(&row).map(|(_, entry)| entry);
return entry_from_row(&row);
}
match self.get(name).await? {
Some(entry) => Err(SecretStoreError::StaleRevision {
@ -309,8 +305,8 @@ impl SecretStore {
.await?;
let entries = rows
.iter()
.map(entry_from_row)
.collect::<Result<HashMap<_, _>, _>>()?;
.map(|row| Ok((row.try_get::<String, _>("name")?, entry_from_row(row)?)))
.collect::<Result<HashMap<_, _>, SecretStoreError>>()?;
Ok(SecretSnapshot(Vault::from_entries(entries)))
}
}
@ -444,7 +440,7 @@ async fn insert_legacy_entry(
Ok(result.rows_affected() == 1)
}
fn entry_from_row(row: &SqliteRow) -> Result<(String, SecretEntry), SecretStoreError> {
fn entry_from_row(row: &SqliteRow) -> Result<SecretEntry, SecretStoreError> {
let metadata = metadata_from_row(row)?;
let name = metadata.name;
let value = row.try_get::<String, _>("value")?;
@ -458,15 +454,14 @@ fn entry_from_row(row: &SqliteRow) -> Result<(String, SecretEntry), SecretStoreE
if revision <= 0 {
return Err(SecretStoreError::StoredRevision { name, revision });
}
let entry = SecretEntry {
Ok(SecretEntry {
value,
secret_type: metadata.secret_type,
description: metadata.description,
created_at: metadata.created_at,
updated_at: metadata.updated_at,
revision,
};
Ok((name, entry))
})
}
fn metadata_from_row(row: &SqliteRow) -> Result<SecretMetadata, SecretStoreError> {
@ -500,13 +495,11 @@ fn parse_timestamp(
column: &'static str,
value: &str,
) -> Result<DateTime<Utc>, SecretStoreError> {
DateTime::parse_from_rfc3339(value)
.map(|timestamp| timestamp.with_timezone(&Utc))
.map_err(|source| SecretStoreError::Timestamp {
name: name.to_string(),
column,
source,
})
fabro_db::parse_rfc3339_utc(value).map_err(|source| SecretStoreError::Timestamp {
name: name.to_string(),
column,
source,
})
}
fn validate_oauth_json(value: &str) -> Result<(), serde_json::Error> {
@ -536,22 +529,11 @@ fn validate_stored_name(name: &str, secret_type: SecretType) -> Result<(), Secre
}
async fn rename_imported_legacy_file(source_path: &Path) -> Result<PathBuf, SecretStoreError> {
let backup_path = legacy_backup_path(source_path, Utc::now());
fs::rename(source_path, &backup_path)
legacy::rename_to_legacy_backup(source_path, "secrets.json")
.await
.map_err(|source| SecretStoreError::LegacyBackup {
.map_err(|err| SecretStoreError::LegacyBackup {
source_path: source_path.to_path_buf(),
backup_path: backup_path.clone(),
source,
})?;
Ok(backup_path)
}
fn legacy_backup_path(source_path: &Path, imported_at: DateTime<Utc>) -> PathBuf {
let timestamp = imported_at.format("%Y%m%dT%H%M%S%fZ");
let mut file_name = source_path
.file_name()
.map_or_else(|| OsString::from("secrets.json"), OsString::from);
file_name.push(format!(".imported-{timestamp}.bak"));
source_path.with_file_name(file_name)
backup_path: err.backup_path,
source: err.source,
})
}