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
17 changes: 16 additions & 1 deletion benchmarks/harness/graphforge_bench/lifecycle_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,16 +20,31 @@ def summarize(documents: dict[str, Any]) -> dict[str, Any]:
for operation, timing in receipt.get("operation_timings", {}).items():
total = calls.setdefault(operation, dict.fromkeys(("calls", "errors", "elapsed_ns"), 0))
for key, value in timing.items():
total[key] += value
# `.get` rather than `+=`: the receipt gained `cpu_ns` and
# `cpu_unmeasured_calls` (#1462) and a fixed seed would raise
# KeyError on any field added after this was written.
total[key] = total.get(key, 0) + value
if not calls:
raise ValueError("runtime diagnosis requires operation timing receipts")
call_ns = sum(timing["elapsed_ns"] for timing in calls.values())
# Process CPU over elapsed wall, per construction operation: how many cores'
# worth each used (#1462). This is the figure #1387's serialized-fraction
# budget is read from -- `seal` is 68-80% of ingest, so its value is the one
# that matters. Omitted per operation when CPU was unavailable, rather than
# reported as zero, and absent entirely for evidence recorded before the
# receipt carried CPU, so historical bundles stay comparable.
effective_cores = {
operation: timing["cpu_ns"] / timing["elapsed_ns"]
for operation, timing in calls.items()
if timing.get("elapsed_ns") and not timing.get("cpu_unmeasured_calls", 1)
}
return {
"identities": documents["result"]["identities"],
"whole_lifecycle_benchexec": documents["benchexec"]["authority"],
"phase_wall_ms": {phase["phase"]: phase["duration_ms"] for phase in phases},
"process_peak_rss_bytes": documents["rung"]["metrics"]["peak_rss_bytes"],
"construction_calls": calls,
**({"construction_call_effective_cores": effective_cores} if effective_cores else {}),
"ingest_outside_construction_calls_ns": ingest["duration_ms"] * 1_000_000 - call_ns,
"lifecycle_outside_phase_wall_seconds": documents["benchexec"]["authority"]["wall_seconds"]
- sum(phase["duration_ms"] for phase in phases) / 1000,
Expand Down
21 changes: 20 additions & 1 deletion benchmarks/runners/certify/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1283,7 +1283,13 @@ fn sanitized_import_operation_timings(value: &serde_json::Value) -> bool {
let Some(timing) = phases.get(*phase).and_then(serde_json::Value::as_object) else {
return false;
};
if timing.len() != 3 {
// Three fields before #1462, five after: `cpu_ns` and
// `cpu_unmeasured_calls` carry process CPU so the serialized
// fraction #1387 budgets can be read off a run. Both shapes are
// accepted, because rejecting the old one would invalidate every
// bundle recorded before the change, and rejecting the new one
// fails every rung after it.
if timing.len() != 3 && timing.len() != 5 {
return false;
}
let (Some(calls), Some(errors), Some(elapsed)) = (
Expand All @@ -1293,6 +1299,19 @@ fn sanitized_import_operation_timings(value: &serde_json::Value) -> bool {
) else {
return false;
};
if timing.len() == 5 {
let (Some(cpu_ns), Some(unmeasured)) = (
timing.get("cpu_ns").and_then(serde_json::Value::as_u64),
timing
.get("cpu_unmeasured_calls")
.and_then(serde_json::Value::as_u64),
) else {
return false;
};
if unmeasured > calls || (calls == 0 && (cpu_ns != 0 || unmeasured != 0)) {
return false;
}
}
errors <= calls && (calls != 0 || elapsed == 0)
})
}
Expand Down
13 changes: 11 additions & 2 deletions benchmarks/schemas/certification-evidence.json
Original file line number Diff line number Diff line change
Expand Up @@ -68,10 +68,19 @@
"properties": {
"calls": { "type": "integer", "minimum": 0, "maximum": 18446744073709551615 },
"errors": { "type": "integer", "minimum": 0, "maximum": 18446744073709551615 },
"elapsed_ns": { "type": "integer", "minimum": 0, "maximum": 18446744073709551615 }
"elapsed_ns": { "type": "integer", "minimum": 0, "maximum": 18446744073709551615 },
"cpu_ns": { "type": "integer", "minimum": 0, "maximum": 18446744073709551615 },
"cpu_unmeasured_calls": { "type": "integer", "minimum": 0, "maximum": 18446744073709551615 }
},
"if": { "properties": { "calls": { "const": 0 } } },
"then": { "properties": { "errors": { "const": 0 }, "elapsed_ns": { "const": 0 } } }
"then": {
"properties": {
"errors": { "const": 0 },
"elapsed_ns": { "const": 0 },
"cpu_ns": { "const": 0 },
"cpu_unmeasured_calls": { "const": 0 }
}
}
},
"importOperationTimings": {
"type": "object",
Expand Down
127 changes: 119 additions & 8 deletions crates/graphforge-api/src/import_session.rs
Original file line number Diff line number Diff line change
Expand Up @@ -328,8 +328,13 @@ struct SessionManifest {
updated_unix_millis: u64,
}

/// Monotonic wall time for attempted calls, including returned errors.
/// These observations are not durable progress or performance limits.
/// Monotonic wall time and process CPU for attempted calls, including returned
/// errors. These observations are not durable progress or performance limits.
///
/// CPU is here so that [`ImportCallTiming::effective_cores`] can answer how much
/// of the machine an operation used. #1387 budgets a serialized fraction of the
/// ingest path, and `seal` is 68-80% of it, so `seal.effective_cores()` is the
/// figure that budget is read from.
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq, Serialize)]
pub struct ImportCallTiming {
/// Number of attempted calls.
Expand All @@ -338,16 +343,69 @@ pub struct ImportCallTiming {
pub errors: u64,
/// Sum of elapsed nanoseconds around these calls.
pub elapsed_ns: u64,
/// Sum of process CPU nanoseconds consumed across these calls.
pub cpu_ns: u64,
/// Calls for which process CPU could not be read. Non-zero makes
/// [`ImportCallTiming::effective_cores`] refuse rather than under-report.
pub cpu_unmeasured_calls: u64,
}

impl ImportCallTiming {
fn record(&mut self, started: Instant, failed: bool) {
fn record(&mut self, started: CallStart, failed: bool) {
self.elapsed_ns = self
.elapsed_ns
.saturating_add(u64::try_from(started.elapsed().as_nanos()).unwrap_or(u64::MAX));
.saturating_add(u64::try_from(started.wall.elapsed().as_nanos()).unwrap_or(u64::MAX));
match (
started.cpu,
graphforge_storage::concurrency_attribution::process_cpu_time(),
) {
(Some(before), Some(after)) => {
self.cpu_ns = self.cpu_ns.saturating_add(
u64::try_from(after.saturating_sub(before).as_nanos()).unwrap_or(u64::MAX),
);
}
_ => self.cpu_unmeasured_calls = self.cpu_unmeasured_calls.saturating_add(1),
}
self.calls = self.calls.saturating_add(1);
self.errors = self.errors.saturating_add(u64::from(failed));
}

/// Process CPU divided by elapsed wall across these calls: how many cores'
/// worth the operation used. `1.0` means it ran on one core.
///
/// `None` when no wall time elapsed or any call's CPU could not be read,
/// rather than a fabricated zero — a reader must be able to tell "not
/// measured" from "measured as idle".
///
/// **This is a process-level ratio.** The CPU term counts every thread in
/// the process, so it answers "how much of this machine did the process use
/// during this operation", which is only the operation's own figure when the
/// process is doing one thing. A ladder rung is; a busy embedding host is
/// not. See `graphforge_storage::concurrency_attribution`.
#[must_use]
pub fn effective_cores(&self) -> Option<f64> {
if self.elapsed_ns == 0 || self.cpu_unmeasured_calls > 0 {
return None;
}
#[allow(clippy::cast_precision_loss, reason = "reporting-only ratio")]
Some(self.cpu_ns as f64 / self.elapsed_ns as f64)
}
}

/// Wall and process CPU captured at the start of a timed call.
#[derive(Clone, Copy, Debug)]
struct CallStart {
wall: Instant,
cpu: Option<std::time::Duration>,
}

impl CallStart {
fn now() -> Self {
Self {
wall: Instant::now(),
cpu: graphforge_storage::concurrency_attribution::process_cpu_time(),
}
}
}

/// Disjoint call timings from the latest import validation or commit invocation.
Expand Down Expand Up @@ -745,7 +803,7 @@ impl GraphImportSession {
self.persist_manifest()?;
}
let chunk_id = format!("import-{:020}-{:020}", source.sequence, batch_index);
let started = Instant::now();
let started = CallStart::now();
let staged = match (input_kind, cancellation) {
(BulkInputKind::Node, Some(token)) => {
construction.append_nodes_with_cancellation(&chunk_id, &batch, token)
Expand Down Expand Up @@ -805,7 +863,7 @@ impl GraphImportSession {
construction: &mut crate::GraphConstructionSession<'_>,
cancellation: Option<&CancellationToken>,
) -> Result<ImportProgress, GfError> {
let started = Instant::now();
let started = CallStart::now();
let sealed = construction.validate_and_seal(cancellation);
self.operation_timings.seal.record(started, sealed.is_err());
sealed?;
Expand Down Expand Up @@ -859,7 +917,7 @@ impl GraphImportSession {
}
self.ensure_base(graph)?;
let mut construction = self.open_construction(graph)?;
let started = Instant::now();
let started = CallStart::now();
let publication = match cancellation {
Some(token) => construction.seal_and_publish_with_cancellation(token),
None => construction.seal_and_publish(),
Expand Down Expand Up @@ -910,7 +968,7 @@ impl GraphImportSession {
graph: &'a GraphForge,
) -> Result<crate::GraphConstructionSession<'a>, GfError> {
let budgets = self.construction_budgets();
let started = Instant::now();
let started = CallStart::now();
if let Some(session_uuid) = self.manifest.construction_session_uuid {
let resumed = graph.resume_graph_construction(session_uuid, budgets);
self.operation_timings
Expand Down Expand Up @@ -1441,6 +1499,59 @@ mod tests {
.join(session_uuid.simple().to_string())
}

#[test]
fn operation_timings_carry_process_cpu_for_every_measured_call() {
// #1462: #1387 budgets a serialized fraction of ingest and nothing on
// the path computed one. This is the end-to-end check that a real
// import now carries CPU, so `seal.effective_cores()` -- the figure the
// budget is read from, seal being 68-80% of ingest -- is available off
// an ordinary run rather than inferred from phase totals.
let (_directory, _project, graph) = fixture();
let mut session = graph
.begin_import_session(
OperationId(Uuid::now_v7()),
ImportSessionLimits {
batch_rows: 1,
..ImportSessionLimits::default()
},
)
.unwrap();
let ids = [Uuid::now_v7(), Uuid::now_v7()];
session
.append_arrow(BulkInputKind::Node, &[nodes(&ids[..1]), nodes(&ids[1..])])
.unwrap();
session
.append_arrow(
BulkInputKind::Edge,
&[edges(Uuid::now_v7(), ids[0], ids[1])],
)
.unwrap();
session.validate(&graph).unwrap();
let timing = session.operation_timings();

for (name, call) in [
("begin", timing.begin),
("append", timing.append),
("seal", timing.seal),
] {
assert!(call.calls > 0, "{name} should have been called");
assert_eq!(
call.cpu_unmeasured_calls, 0,
"{name}: process CPU unavailable on this platform"
);
assert!(
call.effective_cores().is_some(),
"{name}: effective cores must be reported once CPU is measured"
);
}

// An operation never invoked reports nothing rather than zero cores,
// so a reader cannot mistake "not run" for "ran on no CPU".
assert_eq!(timing.publish.calls, 0);
assert_eq!(timing.publish.cpu_ns, 0);
assert!(timing.publish.effective_cores().is_none());
}

#[test]
fn operation_timings_are_scoped_non_durable_and_preserve_cancelled_commit() {
let (_directory, _project, graph) = fixture();
Expand Down