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
2 changes: 1 addition & 1 deletion .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -180,7 +180,7 @@ jobs:
with:
# Pinned — matches `packageManager` in web/package.json.
# Bump together with that field, never alone.
version: 11.28.4
version: 12.9.1
- uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0
with:
node-version: 24
Expand Down
12 changes: 5 additions & 7 deletions crates/common/src/clickhouse_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -46,13 +46,11 @@ pub fn create_client(config: &AuditConfig) -> Option<clickhouse::Client> {
.with_url(url)
.with_database(&config.clickhouse_db)
.with_product_info("think-watch", env!("CARGO_PKG_VERSION"))
.with_setting(NETWORK_COMPRESSION_METHOD.0, NETWORK_COMPRESSION_METHOD.1)
// Plain `RowBinary`, as before clickhouse 0.14. Validation reads
// results as `RowBinaryWithNamesAndTypes` and fails a query whose
// row struct differs from the column types at all, i64 against
// UInt64 included. Not every query has been checked against
// that yet; until they have, it stays off.
.with_validation(false);
.with_setting(NETWORK_COMPRESSION_METHOD.0, NETWORK_COMPRESSION_METHOD.1);
// Row validation stays at the crate's default, on: each result is
// checked against the Rust row it lands in, so a type that drifts
// fails the query instead of decoding into garbage. Every read is
// exercised with complete rows by tests/clickhouse_row_types.rs.

if let Some(ref user) = config.clickhouse_user {
client = client.with_user(user);
Expand Down
2 changes: 1 addition & 1 deletion crates/server/src/handlers/gateway_logs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -402,7 +402,7 @@ pub async fn get_gateway_log_body(
let row: Option<GatewayLogBodyRow> = ch
.query(
"SELECT id, trace_id, user_id, model_id, \
formatDateTime(created_at, '%Y-%m-%dT%H:%M:%S.%fZ', 'UTC') AS created_at, \
formatDateTime(created_at, '%Y-%m-%dT%H:%i:%S.%fZ', 'UTC') AS created_at, \
request_body, response_body, request_body_bytes, \
response_body_bytes, body_capture_status \
FROM gateway_logs WHERE id = ? LIMIT 1",
Expand Down
2 changes: 1 addition & 1 deletion crates/server/src/handlers/mcp_logs.rs
Original file line number Diff line number Diff line change
Expand Up @@ -272,7 +272,7 @@ pub async fn get_mcp_log_body(
let row: Option<McpLogBodyRow> = ch
.query(
"SELECT id, trace_id, user_id, server_id, server_name, tool_name, \
formatDateTime(created_at, '%Y-%m-%dT%H:%M:%S.%fZ', 'UTC') AS created_at, \
formatDateTime(created_at, '%Y-%m-%dT%H:%i:%S.%fZ', 'UTC') AS created_at, \
tool_arguments, tool_result, arguments_bytes, result_bytes, \
body_capture_status \
FROM mcp_logs WHERE id = ? LIMIT 1",
Expand Down
12 changes: 7 additions & 5 deletions crates/server/src/handlers/models.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1304,7 +1304,7 @@ pub async fn get_route_history(
// provisioned yet, we fall back to an empty response — the
// sparkline is a hint, not load-bearing.
let sql = format!(
"SELECT toUnixTimestamp(toStartOfMinute(created_at)) AS bucket_ts, \
"SELECT toInt64(toUnixTimestamp(toStartOfMinute(created_at))) AS bucket_ts, \
quantile(0.50)(latency_ms) AS p50, \
quantile(0.95)(latency_ms) AS p95, \
count() AS requests, \
Expand All @@ -1318,11 +1318,13 @@ pub async fn get_route_history(
ORDER BY bucket_ts"
);

// `latency_ms` is Nullable, so its quantiles are too: NULL for a
// minute whose requests all failed before a latency was recorded.
#[derive(Debug, clickhouse::Row, serde::Deserialize)]
struct Row {
bucket_ts: i64,
p50: f64,
p95: f64,
p50: Option<f64>,
p95: Option<f64>,
requests: u64,
errors: u64,
}
Expand All @@ -1348,8 +1350,8 @@ pub async fn get_route_history(
.into_iter()
.map(|r| RouteHistoryBucket {
ts: r.bucket_ts,
p50_ms: if r.p50.is_finite() { Some(r.p50) } else { None },
p95_ms: if r.p95.is_finite() { Some(r.p95) } else { None },
p50_ms: r.p50.filter(|v| v.is_finite()),
p95_ms: r.p95.filter(|v| v.is_finite()),
requests: r.requests,
errors: r.errors,
})
Expand Down
14 changes: 9 additions & 5 deletions crates/server/src/handlers/trace.rs
Original file line number Diff line number Diff line change
Expand Up @@ -108,7 +108,7 @@ pub async fn get_trace(
let rows: Vec<GatewayRow> = ch
.query(
"SELECT id, \
formatDateTime(created_at, '%Y-%m-%dT%H:%M:%S.%fZ', 'UTC') AS created_at, \
formatDateTime(created_at, '%Y-%m-%dT%H:%i:%S.%fZ', 'UTC') AS created_at, \
model_id, status_code, latency_ms, user_id \
FROM gateway_logs \
WHERE trace_id = ?",
Expand Down Expand Up @@ -141,7 +141,7 @@ pub async fn get_trace(
let rows: Vec<McpRow> = ch
.query(
"SELECT id, \
formatDateTime(created_at, '%Y-%m-%dT%H:%M:%S.%fZ', 'UTC') AS created_at, \
formatDateTime(created_at, '%Y-%m-%dT%H:%i:%S.%fZ', 'UTC') AS created_at, \
tool_name, status, duration_ms, user_id \
FROM mcp_logs \
WHERE trace_id = ?",
Expand Down Expand Up @@ -172,7 +172,7 @@ pub async fn get_trace(
let rows: Vec<AuditRow> = ch
.query(
"SELECT id, \
formatDateTime(created_at, '%Y-%m-%dT%H:%M:%S.%fZ', 'UTC') AS created_at, \
formatDateTime(created_at, '%Y-%m-%dT%H:%i:%S.%fZ', 'UTC') AS created_at, \
action, user_id \
FROM audit_logs \
WHERE trace_id = ?",
Expand Down Expand Up @@ -210,6 +210,10 @@ pub async fn get_trace(
level: String,
message: String,
}
// `app_logs.created_at`, qualified: unqualified, the name means
// the String alias above, the comparison with a DateTime fails,
// and `unwrap_or_default` turned that into "no app events".
//
// Escape LIKE metacharacters in the user-supplied trace_id so
// `trace_id=%` doesn't trigger a 1h-window app_logs scan
// bypassing the substring-search intent. LIMIT 200 + PREWHERE
Expand All @@ -222,10 +226,10 @@ pub async fn get_trace(
let app_rows: Vec<AppLogRow> = ch
.query(
"SELECT id, \
formatDateTime(created_at, '%Y-%m-%dT%H:%M:%S.%fZ', 'UTC') AS created_at, \
formatDateTime(created_at, '%Y-%m-%dT%H:%i:%S.%fZ', 'UTC') AS created_at, \
level, message \
FROM app_logs \
PREWHERE created_at >= now() - INTERVAL 1 HOUR \
PREWHERE app_logs.created_at >= now() - INTERVAL 1 HOUR \
AND (fields LIKE ? OR span LIKE ?) \
LIMIT 200",
)
Expand Down
20 changes: 12 additions & 8 deletions crates/test-support/src/ch.rs
Original file line number Diff line number Diff line change
Expand Up @@ -43,14 +43,18 @@ impl IsolatedClickHouseDatabase {
.await
.context("CREATE DATABASE on test ClickHouse")?;

// Test client targets the new DB.
let mut client = client_for(url).with_database(&name);
if let Some(u) = user {
client = client.with_user(u);
}
if let Some(p) = password {
client = client.with_password(p);
}
// The application's own client, pointed at the new DB: the same
// settings, row validation and connector as production, so a
// query that would fail there fails here.
let client = think_watch_common::clickhouse_client::create_client(
&think_watch_common::audit::AuditConfig {
clickhouse_url: Some(url.to_string()),
clickhouse_db: name.clone(),
clickhouse_user: user.map(String::from),
clickhouse_password: password.map(String::from),
},
)
.context("build the ClickHouse client")?;

// Load the production schema into the per-test DB. The
// bundled init SQL contains a `CREATE DATABASE IF NOT EXISTS
Expand Down
Loading
Loading