Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1,355 changes: 1,355 additions & 0 deletions cli/src/services/agent_trace_db/lock_contention_tests.rs

Large diffs are not rendered by default.

2 changes: 2 additions & 0 deletions cli/src/services/agent_trace_db/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,8 @@ use crate::services::{
use serde_json::Value;

pub mod lifecycle;
#[cfg(test)]
mod lock_contention_tests;
pub mod repository;

/// Payload type discriminator for diff trace source payloads.
Expand Down
122 changes: 114 additions & 8 deletions cli/src/services/agent_trace_db/repository.rs
Original file line number Diff line number Diff line change
Expand Up @@ -76,6 +76,16 @@ pub fn is_valid_source_instance_id(value: &str) -> bool {
!value.trim().is_empty()
}

fn ensure_repository_id_matches(stored_repository_id: &str, repository_id: &str) -> Result<()> {
if stored_repository_id != repository_id {
anyhow::bail!(
"repository Agent Trace DB metadata mismatch: stored repository ID \
{stored_repository_id} does not match resolved repository ID {repository_id}"
);
}
Ok(())
}

/// Repository-scoped Agent Trace database configuration.
pub struct RepositoryAgentTraceDbSpec;

Expand Down Expand Up @@ -157,7 +167,19 @@ impl RepositoryAgentTraceDb {
&self,
repository_id: &str,
) -> Result<RepositoryMetadata> {
self.execute(INSERT_REPOSITORY_METADATA_SQL, (repository_id,))?;
if let Some((stored_repository_id, source_instance_id)) =
self.select_repository_metadata_row()?
{
ensure_repository_id_matches(&stored_repository_id, repository_id)?;
if is_valid_source_instance_id(&source_instance_id) {
return Ok(RepositoryMetadata {
repository_id: stored_repository_id,
source_instance_id,
});
}
}

self.execute_idempotent_write(INSERT_REPOSITORY_METADATA_SQL, (repository_id,))?;

let Some((stored_repository_id, source_instance_id)) =
self.select_repository_metadata_row()?
Expand All @@ -168,12 +190,7 @@ impl RepositoryAgentTraceDb {
);
};

if stored_repository_id != repository_id {
anyhow::bail!(
"repository Agent Trace DB metadata mismatch: stored repository ID \
{stored_repository_id} does not match resolved repository ID {repository_id}"
);
}
ensure_repository_id_matches(&stored_repository_id, repository_id)?;

if is_valid_source_instance_id(&source_instance_id) {
return Ok(RepositoryMetadata {
Expand All @@ -183,7 +200,7 @@ impl RepositoryAgentTraceDb {
}

let candidate = generate_source_instance_id();
self.execute(CLAIM_SOURCE_INSTANCE_ID_SQL, (candidate.as_str(),))?;
self.execute_idempotent_write(CLAIM_SOURCE_INSTANCE_ID_SQL, (candidate.as_str(),))?;

let (final_repository_id, final_source_instance_id) =
self.select_repository_metadata_row()?.ok_or_else(|| {
Expand Down Expand Up @@ -906,6 +923,95 @@ mod tests {
remove_test_db(&db_path);
}

#[test]
fn initialized_repository_metadata_hook_runtime_open_issues_no_writes() {
let db_path = unique_test_db_path("metadata-hook-open-no-writes");
let repository_id = "a".repeat(64);

let setup = RepositoryAgentTraceDb::new_at(&db_path).expect("repository DB should open");
let (first, first_writes) = crate::services::db::count_write_statements(|| {
setup.verify_or_initialize_repository_metadata(&repository_id)
});
let first = first.expect("first metadata initialization should succeed");
assert!(
first_writes > 0,
"first initialization must seed and claim metadata"
);
drop(setup);

let (reopened, writes) = crate::services::db::count_write_statements(|| {
let hook = RepositoryAgentTraceDb::open_for_hooks_without_migrations_at(&db_path)?;
hook.ensure_schema_ready_for_hooks()?;
hook.verify_or_initialize_repository_metadata(&repository_id)
});
let reopened = reopened.expect("initialized hook-runtime open should succeed");
assert_eq!(
writes, 0,
"initialized hook-runtime open must issue no writes"
);
assert_eq!(reopened, first);

remove_test_db(&db_path);
}

#[test]
fn mismatched_repository_metadata_errors_without_writes() {
let db_path = unique_test_db_path("metadata-mismatch-no-writes");
let stored_repository_id = "a".repeat(64);
let other_repository_id = "b".repeat(64);

let db = RepositoryAgentTraceDb::new_at(&db_path).expect("repository DB should open");
db.verify_or_initialize_repository_metadata(&stored_repository_id)
.expect("first metadata initialization should succeed");

let (result, writes) = crate::services::db::count_write_statements(|| {
db.verify_or_initialize_repository_metadata(&other_repository_id)
});
let message = result
.expect_err("mismatched repository ID should fail validation")
.to_string();
assert!(
message.contains("metadata mismatch"),
"unexpected error: {message}"
);
assert_eq!(writes, 0, "a rejected mismatch must issue no writes");

remove_test_db(&db_path);
}

#[test]
fn repository_metadata_with_empty_source_instance_id_is_claimed_once() {
let db_path = unique_test_db_path("metadata-empty-source-instance");
let repository_id = "a".repeat(64);

let db = RepositoryAgentTraceDb::new_at(&db_path).expect("repository DB should open");
db.verify_or_initialize_repository_metadata(&repository_id)
.expect("first metadata initialization should succeed");
db.execute(
"UPDATE repository_metadata SET source_instance_id = '' WHERE id = 1",
(),
)
.expect("source-instance ID should reset to the empty placeholder");

let claimed = db
.verify_or_initialize_repository_metadata(&repository_id)
.expect("an empty source-instance ID should be claimed");
assert_eq!(claimed.repository_id, repository_id);
assert!(is_valid_source_instance_id(&claimed.source_instance_id));

let (reopened, writes) = crate::services::db::count_write_statements(|| {
db.verify_or_initialize_repository_metadata(&repository_id)
});
assert_eq!(
reopened.expect("claimed metadata should validate"),
claimed,
"a claimed source-instance ID must not be overwritten"
);
assert_eq!(writes, 0);

remove_test_db(&db_path);
}

#[test]
fn mismatched_repository_metadata_errors_on_open() {
let db_path = unique_test_db_path("mismatch");
Expand Down
159 changes: 154 additions & 5 deletions cli/src/services/config/render.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ use super::resolver::{
AuthConfigKeySpec, RuntimeConfig, CONTROL_PLANE_BASE_URL_KEY, PRECEDENCE_DESCRIPTION,
WORKOS_CLIENT_ID_KEY,
};
use super::types::DatabaseRetryConfig;
use super::types::{AgentTraceDbRetryConfig, DatabaseRetryConfig};
use super::{ConfigPathSource, ReportFormat, ResolvedOptionalValue, ValueSource};

#[allow(clippy::too_many_lines)]
Expand Down Expand Up @@ -429,15 +429,35 @@ fn format_per_db_retry_text(
lines
}

fn format_agent_trace_db_retry_text(config: &AgentTraceDbRetryConfig) -> Vec<String> {
let db_label = "agent_trace_db";
let mut lines = format_per_db_retry_text(&config.retry, db_label);
if let Some(busy_timeout_ms) = config.busy_timeout_ms {
lines.push(format!(
" {}: {} (busy_timeout_ms)",
style::label(db_label),
style::value(&format!("{busy_timeout_ms}ms"))
));
}
if let Some(contention_deadline_ms) = config.contention_deadline_ms {
lines.push(format!(
" {}: {} (contention_deadline_ms)",
style::label(db_label),
style::value(&format!("{contention_deadline_ms}ms"))
));
}
lines
}

fn format_database_retry_text(value: &ResolvedOptionalValue<DatabaseRetryConfig>) -> String {
match (value.value.as_ref(), value.source) {
(Some(config), Some(source)) => {
let mut lines = vec![format!(" {}:", style::label("policies.database_retry"))];
if let Some(ref per_db) = config.local_db {
lines.extend(format_per_db_retry_text(per_db, "local_db"));
}
if let Some(ref per_db) = config.agent_trace_db {
lines.extend(format_per_db_retry_text(per_db, "agent_trace_db"));
if let Some(ref agent_trace_db) = config.agent_trace_db {
lines.extend(format_agent_trace_db_retry_text(agent_trace_db));
}
if let Some(ref per_db) = config.auth_db {
lines.extend(format_per_db_retry_text(per_db, "auth_db"));
Expand Down Expand Up @@ -492,17 +512,33 @@ fn format_per_db_retry_json(config: &super::types::PerDbRetryConfig) -> Value {
Value::Object(obj)
}

fn format_agent_trace_db_retry_json(config: &AgentTraceDbRetryConfig) -> Value {
let mut value = format_per_db_retry_json(&config.retry);
if let Value::Object(ref mut obj) = value {
if let Some(busy_timeout_ms) = config.busy_timeout_ms {
obj.insert("busy_timeout_ms".to_string(), json!(busy_timeout_ms));
}
if let Some(contention_deadline_ms) = config.contention_deadline_ms {
obj.insert(
"contention_deadline_ms".to_string(),
json!(contention_deadline_ms),
);
}
}
value
}

fn format_database_retry_json(value: &ResolvedOptionalValue<DatabaseRetryConfig>) -> Value {
let config = value.value.as_ref();
let mut resolved = serde_json::Map::new();
if let Some(c) = config {
if let Some(ref per_db) = c.local_db {
resolved.insert("local_db".to_string(), format_per_db_retry_json(per_db));
}
if let Some(ref per_db) = c.agent_trace_db {
if let Some(ref agent_trace_db) = c.agent_trace_db {
resolved.insert(
"agent_trace_db".to_string(),
format_per_db_retry_json(per_db),
format_agent_trace_db_retry_json(agent_trace_db),
);
}
if let Some(ref per_db) = c.auth_db {
Expand All @@ -515,3 +551,116 @@ fn format_database_retry_json(value: &ResolvedOptionalValue<DatabaseRetryConfig>
"config_source": value.source.and_then(ValueSource::config_source).map(ConfigPathSource::as_str),
})
}

#[cfg(test)]
mod database_retry_render_tests {
use serde_json::json;

use super::{format_database_retry_json, format_database_retry_text};
use crate::services::config::{
AgentTraceDbRetryConfig, ConfigPathSource, DatabaseRetryConfig, PerDbRetryConfig,
ResolvedOptionalValue, ValueSource,
};
use crate::services::resilience::RetryPolicy;

fn query_policy() -> RetryPolicy {
RetryPolicy {
max_attempts: 3,
timeout_ms: 150,
initial_backoff_ms: 10,
max_backoff_ms: 50,
}
}

fn resolved(
busy_timeout_ms: Option<u64>,
contention_deadline_ms: Option<u64>,
) -> ResolvedOptionalValue<DatabaseRetryConfig> {
ResolvedOptionalValue {
value: Some(DatabaseRetryConfig {
local_db: Some(PerDbRetryConfig {
connection_open: None,
query: Some(query_policy()),
}),
agent_trace_db: Some(AgentTraceDbRetryConfig {
retry: PerDbRetryConfig {
connection_open: None,
query: Some(query_policy()),
},
busy_timeout_ms,
contention_deadline_ms,
}),
auth_db: None,
}),
source: Some(ValueSource::ConfigFile(ConfigPathSource::Flag)),
}
}

#[test]
fn database_retry_json_renders_contention_keys_for_agent_trace_db_only() {
let rendered = format_database_retry_json(&resolved(Some(750), Some(2_000)));

assert_eq!(
rendered["resolved"],
json!({
"local_db": {
"query": {
"max_attempts": 3,
"timeout_ms": 150,
"initial_backoff_ms": 10,
"max_backoff_ms": 50,
},
},
"agent_trace_db": {
"query": {
"max_attempts": 3,
"timeout_ms": 150,
"initial_backoff_ms": 10,
"max_backoff_ms": 50,
},
"busy_timeout_ms": 750,
"contention_deadline_ms": 2000,
},
})
);
}

#[test]
fn database_retry_json_omits_unset_contention_keys() {
let rendered = format_database_retry_json(&resolved(None, None));

let agent_trace_db = rendered["resolved"]["agent_trace_db"].as_object().unwrap();
assert!(!agent_trace_db.contains_key("busy_timeout_ms"));
assert!(!agent_trace_db.contains_key("contention_deadline_ms"));
assert_eq!(
rendered["resolved"]["agent_trace_db"]["query"]["timeout_ms"],
150
);
}

#[test]
fn database_retry_text_renders_contention_keys_for_agent_trace_db() {
let rendered = format_database_retry_text(&resolved(Some(750), Some(0)));

let busy_line = rendered
.lines()
.find(|line| line.contains("(busy_timeout_ms)"))
.unwrap();
assert!(busy_line.contains("agent_trace_db"), "{rendered}");
assert!(busy_line.contains("750ms"), "{rendered}");
let deadline_line = rendered
.lines()
.find(|line| line.contains("(contention_deadline_ms)"))
.unwrap();
assert!(deadline_line.contains("agent_trace_db"), "{rendered}");
assert!(deadline_line.contains("0ms"), "{rendered}");
assert_eq!(
rendered
.lines()
.filter(|line| line.contains("3 attempts, 150ms timeout, 10..50ms backoff (query)"))
.count(),
2,
"{rendered}"
);
}
}
Loading
Loading