Skip to content

Commit 6436088

Browse files
committed
add observability metadata to tracing
1 parent 4689e24 commit 6436088

7 files changed

Lines changed: 418 additions & 1 deletion

File tree

vigilo/src/cli/commands/coordinator.rs

Lines changed: 37 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -308,11 +308,23 @@ async fn drain_dispatch_batch(
308308
debug!(coordinator_id = %coordinator_id, "draining dispatchable run-shard windows");
309309

310310
let mut dispatched = 0usize;
311+
let mut dispatched_by_alias = std::collections::BTreeMap::<String, usize>::new();
311312
for _ in 0..config.max_dispatch_per_cycle {
312313
let control_db = database.control().await?;
313314
let Some(route) = run_dispatch::select_next_dispatch_route(control_db).await? else {
314315
break;
315316
};
317+
let database_alias = route.database_alias.clone();
318+
319+
debug!(
320+
coordinator_id = %coordinator_id,
321+
run_id = %route.run_id,
322+
run_shard = route.run_shard,
323+
database_alias = %database_alias,
324+
placement_status = %route.placement_status,
325+
routing_decision = "selected_dispatch_route",
326+
"selected dispatch route"
327+
);
316328

317329
let Some(snapshot) = run_dispatch::prepare_dispatch_run_snapshot(
318330
control_db,
@@ -339,16 +351,22 @@ async fn drain_dispatch_batch(
339351
};
340352

341353
dispatched += 1;
354+
*dispatched_by_alias
355+
.entry(database_alias.clone())
356+
.or_default() += 1;
342357
debug!(
343358
run_id = %run.id,
344359
run_key = %run.run_key,
345360
run_shard = run.run_shard,
361+
database_alias = %database_alias,
362+
routing_decision = "dispatch_window_claimed",
346363
"claimed dispatchable run shard window"
347364
);
348365
info!(
349366
run_id = %run.id,
350367
run_key = %run.run_key,
351368
run_shard = run.run_shard,
369+
database_alias = %database_alias,
352370
chunk_events_enqueued = run.chunk_events_enqueued,
353371
chunks_marked_dispatched = run.chunks_marked_dispatched,
354372
run_started_events_enqueued = run.run_started_events_enqueued,
@@ -364,6 +382,14 @@ async fn drain_dispatch_batch(
364382
dispatch_windows_prepared = dispatched,
365383
"completed coordinator dispatch drain pass"
366384
);
385+
for (alias, alias_dispatched) in dispatched_by_alias {
386+
info!(
387+
coordinator_id = %coordinator_id,
388+
database_alias = %alias,
389+
dispatch_windows_prepared = alias_dispatched,
390+
"completed coordinator dispatch drain pass for execution placement"
391+
);
392+
}
367393
}
368394

369395
Ok(dispatched)
@@ -491,6 +517,7 @@ async fn collect_run_shard_summaries(
491517
) -> anyhow::Result<Vec<run_shard_summary::RunShardSummary>> {
492518
let routes = database.execution_routes_for_run(run_id).await?;
493519
let mut summaries = Vec::with_capacity(routes.len());
520+
let mut summaries_by_alias = std::collections::BTreeMap::<String, usize>::new();
494521

495522
for (run_shard, alias, db) in routes {
496523
let Some(summary) =
@@ -514,8 +541,18 @@ async fn collect_run_shard_summaries(
514541
expected_execution_count = summary.expected_execution_count,
515542
"loaded run shard summary for finalization"
516543
);
544+
*summaries_by_alias.entry(alias).or_default() += 1;
517545
summaries.push(summary);
518546
}
519547

548+
for (alias, count) in summaries_by_alias {
549+
debug!(
550+
run_id = %run_id,
551+
database_alias = %alias,
552+
shard_summaries_loaded = count,
553+
"loaded routed shard summaries for finalization from execution placement"
554+
);
555+
}
556+
520557
Ok(summaries)
521558
}

vigilo/src/cli/commands/shard.rs

Lines changed: 105 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -40,6 +40,16 @@ pub(crate) enum SubCommand {
4040
command: PlacementSubCommand,
4141
},
4242

43+
/// Inspect the resolved route for one run shard
44+
Route {
45+
/// Run UUID
46+
run_id: String,
47+
48+
/// Logical run shard
49+
#[arg(value_parser = clap::value_parser!(i16).range(0..=127))]
50+
run_shard: i16,
51+
},
52+
4353
/// Move one run shard to another database placement
4454
Move {
4555
/// Run UUID
@@ -160,6 +170,9 @@ impl Executable for Command {
160170
.await
161171
}
162172
SubCommand::Placements { command } => exec_placement_command(context, command).await,
173+
SubCommand::Route { run_id, run_shard } => {
174+
exec_route_command(context, run_id, run_shard).await
175+
}
163176
}
164177
}
165178
}
@@ -239,6 +252,28 @@ async fn exec_move_command(
239252
Ok(())
240253
}
241254

255+
async fn exec_route_command(
256+
context: Context,
257+
run_id: String,
258+
run_shard: i16,
259+
) -> anyhow::Result<()> {
260+
let run_id = parse_run_id(&run_id)?;
261+
let database = context.db().await?;
262+
let out = context.out().await?;
263+
let route = shard_admin::inspect_shard_route(database, run_id, run_shard).await?;
264+
265+
info!(
266+
run_id = %run_id,
267+
run_shard,
268+
database_alias = %route.database_alias,
269+
placement_status = %route.shard_placement_status,
270+
routing_decision = route.routing_decision,
271+
"inspected shard route"
272+
);
273+
out.write_value(&route_payload(&route))?;
274+
Ok(())
275+
}
276+
242277
async fn exec_placement_command(
243278
context: Context,
244279
command: PlacementSubCommand,
@@ -378,6 +413,24 @@ fn move_payload(outcome: &shard_admin::ShardMoveOutcome) -> Value {
378413
})
379414
}
380415

416+
fn route_payload(route: &shard_admin::ShardRouteInspection) -> Value {
417+
json!({
418+
"data": {
419+
"run_id": route.run_id,
420+
"run_shard": route.run_shard,
421+
"database_alias": route.database_alias,
422+
"shard_placement_status": route.shard_placement_status,
423+
"database_role": route.database_role,
424+
"database_status": route.database_status,
425+
"database_url_env": route.database_url_env,
426+
"database_url_env_resolved": route.database_url_env_resolved,
427+
"dispatchable": route.dispatchable,
428+
"readable": route.readable,
429+
"routing_decision": route.routing_decision,
430+
}
431+
})
432+
}
433+
381434
#[cfg(test)]
382435
mod tests {
383436
use chrono::Utc;
@@ -442,6 +495,23 @@ mod tests {
442495
assert_eq!(alias, "shard_001");
443496
}
444497

498+
#[test]
499+
fn route_command_matches_argument_shape() {
500+
let run_id = Uuid::now_v7().to_string();
501+
let cli = TestCli::try_parse_from(["vigilo", "route", &run_id, "4"]).unwrap();
502+
503+
let SubCommand::Route {
504+
run_id: parsed_run_id,
505+
run_shard,
506+
} = cli.command
507+
else {
508+
panic!("expected shard route command");
509+
};
510+
511+
assert_eq!(parsed_run_id, run_id);
512+
assert_eq!(run_shard, 4);
513+
}
514+
445515
#[test]
446516
fn database_list_payload_reports_count() {
447517
let placement = DatabasePlacement {
@@ -531,4 +601,39 @@ mod tests {
531601
assert_eq!(payload["meta"]["verified"], json!(true));
532602
assert_eq!(payload["data"]["tables"][0]["table"], json!("run_chunks"));
533603
}
604+
605+
#[test]
606+
fn route_payload_excludes_secret_url_value() {
607+
let run_id = Uuid::now_v7();
608+
let route = shard_admin::ShardRouteInspection {
609+
run_id,
610+
run_shard: 4,
611+
database_alias: "shard_001".to_string(),
612+
shard_placement_status: "active".to_string(),
613+
database_role: "shard".to_string(),
614+
database_status: "active".to_string(),
615+
database_url_env: "VIGILO_SHARD_001_DATABASE_URL".to_string(),
616+
database_url_env_resolved: true,
617+
dispatchable: true,
618+
readable: true,
619+
routing_decision: "dispatchable",
620+
};
621+
622+
let payload = route_payload(&route);
623+
624+
assert_eq!(payload["data"]["run_id"], json!(run_id));
625+
assert_eq!(payload["data"]["run_shard"], json!(4));
626+
assert_eq!(payload["data"]["database_alias"], json!("shard_001"));
627+
assert_eq!(
628+
payload["data"]["database_url_env"],
629+
json!("VIGILO_SHARD_001_DATABASE_URL")
630+
);
631+
assert!(
632+
payload
633+
.to_string()
634+
.contains("VIGILO_SHARD_001_DATABASE_URL")
635+
);
636+
assert!(!payload.to_string().contains("postgres://"));
637+
assert_eq!(payload["data"]["routing_decision"], json!("dispatchable"));
638+
}
534639
}

vigilo/src/context/database.rs

Lines changed: 60 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -202,6 +202,14 @@ impl Db {
202202
let key = ShardPlacementKey { run_id, run_shard };
203203
if let Some(placement) = self.shard_placement_cache.get(&key).await {
204204
self.validate_shard_placement_alias(&placement).await?;
205+
debug!(
206+
run_id = %run_id,
207+
run_shard,
208+
database_alias = %placement.database_alias,
209+
placement_status = %placement.status,
210+
routing_decision = "cached_execution_placement",
211+
"resolved execution placement from cache"
212+
);
205213
return Ok(placement);
206214
}
207215

@@ -213,6 +221,14 @@ impl Db {
213221
};
214222

215223
self.validate_shard_placement_alias(&placement).await?;
224+
debug!(
225+
run_id = %run_id,
226+
run_shard,
227+
database_alias = %placement.database_alias,
228+
placement_status = %placement.status,
229+
routing_decision = "control_lookup_execution_placement",
230+
"resolved execution placement from control metadata"
231+
);
216232
self.shard_placement_cache
217233
.insert(key, placement.clone())
218234
.await;
@@ -229,6 +245,14 @@ impl Db {
229245
let placement = self.execution_placement(run_id, run_shard).await?;
230246

231247
if !placement.is_dispatchable() {
248+
debug!(
249+
run_id = %run_id,
250+
run_shard,
251+
database_alias = %placement.database_alias,
252+
placement_status = %placement.status,
253+
routing_decision = "blocked_non_dispatchable",
254+
"execution placement is not dispatchable"
255+
);
232256
return Err(ExecutionRouteError::NonDispatchableShardPlacement {
233257
run_id,
234258
run_shard,
@@ -237,6 +261,14 @@ impl Db {
237261
.into());
238262
}
239263

264+
debug!(
265+
run_id = %run_id,
266+
run_shard,
267+
database_alias = %placement.database_alias,
268+
placement_status = %placement.status,
269+
routing_decision = "dispatchable_execution_pool",
270+
"resolved dispatchable execution pool"
271+
);
240272
self.placement(&placement.database_alias).await
241273
}
242274

@@ -253,6 +285,14 @@ impl Db {
253285
for placement in placements {
254286
self.validate_shard_placement_alias(&placement).await?;
255287
if !placement.is_dispatchable() {
288+
debug!(
289+
run_id = %run_id,
290+
run_shard = placement.run_shard,
291+
database_alias = %placement.database_alias,
292+
placement_status = %placement.status,
293+
routing_decision = "blocked_non_dispatchable",
294+
"execution route is not dispatchable"
295+
);
256296
anyhow::bail!(
257297
"shard placement for run {} shard {} has status {}, which is not dispatchable",
258298
run_id,
@@ -262,6 +302,14 @@ impl Db {
262302
}
263303

264304
let pool = self.placement(&placement.database_alias).await?.clone();
305+
debug!(
306+
run_id = %run_id,
307+
run_shard = placement.run_shard,
308+
database_alias = %placement.database_alias,
309+
placement_status = %placement.status,
310+
routing_decision = "dispatchable_execution_route",
311+
"resolved dispatchable execution route"
312+
);
265313
routed.push((placement.run_shard, placement.database_alias, pool));
266314
}
267315

@@ -284,6 +332,14 @@ impl Db {
284332
for placement in placements {
285333
self.validate_shard_placement_alias(&placement).await?;
286334
let pool = self.placement(&placement.database_alias).await?.clone();
335+
debug!(
336+
run_id = %run_id,
337+
run_shard = placement.run_shard,
338+
database_alias = %placement.database_alias,
339+
placement_status = %placement.status,
340+
routing_decision = "readable_execution_route",
341+
"resolved readable execution route"
342+
);
287343
routed.push((placement.run_shard, placement.database_alias, pool));
288344
}
289345

@@ -364,6 +420,10 @@ impl Db {
364420
})
365421
}
366422

423+
pub(crate) fn database_url_env_is_resolved(&self, database_url_env: &str) -> bool {
424+
database_url_env == DEFAULT_DATABASE_URL_ENV || std::env::var_os(database_url_env).is_some()
425+
}
426+
367427
#[allow(dead_code)]
368428
async fn placement_catalog(&self) -> anyhow::Result<&PlacementCatalog> {
369429
self.placement_catalog

vigilo/src/db/workflows/run_dispatch.rs

Lines changed: 9 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,8 @@ pub(crate) struct DispatchRunSnapshot {
5454
pub(crate) struct DispatchRoute {
5555
pub(crate) run_id: Uuid,
5656
pub(crate) run_shard: i16,
57+
pub(crate) database_alias: String,
58+
pub(crate) placement_status: String,
5759
}
5860

5961
#[derive(Debug, Clone, sqlx::FromRow)]
@@ -279,7 +281,9 @@ pub(crate) async fn select_next_dispatch_route(
279281
r#"
280282
SELECT
281283
c.run_id,
282-
c.run_shard
284+
c.run_shard,
285+
sp.database_alias,
286+
sp.status AS placement_status
283287
FROM run_shard_dispatch_cursors c
284288
JOIN runs r
285289
ON r.id = c.run_id
@@ -1143,6 +1147,8 @@ mod tests {
11431147

11441148
assert_eq!(route.run_id, run_id);
11451149
assert_eq!(route.run_shard, 1);
1150+
assert_eq!(route.database_alias, "primary");
1151+
assert_eq!(route.placement_status, "active");
11461152
}
11471153

11481154
#[sqlx::test(migrations = "../migrations")]
@@ -1152,6 +1158,8 @@ mod tests {
11521158
let route = DispatchRoute {
11531159
run_id,
11541160
run_shard: 1,
1161+
database_alias: "primary".to_string(),
1162+
placement_status: "active".to_string(),
11551163
};
11561164
let snapshot = prepare_dispatch_run_snapshot(&pool, &route, Uuid::now_v7(), 60)
11571165
.await

0 commit comments

Comments
 (0)