Compare commits

..
4 changed files with 26 additions and 34 deletions

View file

@ -217,7 +217,15 @@ impl JobQueue {
) -> anyhow::Result<u64> {
let mut inner = self.lock();
let container = inner
.append(NodeKind::Dag { source, reason }, Vec::new(), None)
.append(
NodeKind::Dag {
source,
reason,
created_at: Utc::now(),
},
Vec::new(),
None,
)
.map_err(|e| anyhow::anyhow!("job_queue: container insert failed: {e}"))?;
insert_group(&mut inner, declare, Some(container))?;
drop(inner);
@ -399,15 +407,15 @@ fn container(sched: &Sched, dag_id: u64) -> Option<NodeId> {
/// DAG drops out of the snapshot entirely (a `Failed` one lingers until
/// aged out).
fn dag_view(sched: &Sched, container: NodeId) -> Option<DagView> {
// Read straight off the container's payload: the domain fields have a
// Read straight off the container's payload: the three fields have a
// single home there, so an intermediate owned copy of them was a second
// type describing the same data rather than a grouping side-table.
//
// `created_at` is NOT among them — it comes off the container node itself
// below, where the graph stamps it for every node. Keeping a payload copy
// would record one instant in two places.
let node = sched.graph().node(container)?;
let NodeKind::Dag { source, reason } = &node.payload else {
let NodeKind::Dag {
source,
reason,
created_at,
} = &sched.graph().node(container)?.payload
else {
return None;
};
let all: Vec<_> = sched.graph().descendants(container).collect();
@ -486,7 +494,7 @@ fn dag_view(sched: &Sched, container: NodeId) -> Option<DagView> {
id: container.get(),
source: *source,
reason: reason.clone(),
created_at: node.created_at,
created_at: *created_at,
started_at: started.into_iter().min(),
finished_at: is_terminal.then(|| finished.into_iter().max()).flatten(),
nodes,

View file

@ -12,6 +12,7 @@
//! DAG can span agents). See `docs/coordinator.md::Job queue` for the
//! full design.
use chrono::{DateTime, Utc};
pub use hive_host_sock::jobs::{DagView, PermPayload, Source, State};
use serde::Serialize;
@ -275,12 +276,11 @@ pub enum NodeKind {
/// the DAG state. Pure grouping — lease- and
/// build-slot-exempt; the executor instant-completes it (`Done`) so it
/// reaches `Finishing` and its children start.
///
/// No `created_at` here: the graph stamps [`hive_jobq::Node::created_at`] on
/// every node at insert, so the container already has one. A second copy in
/// the payload would be the same instant recorded twice, with only this
/// variant's version reachable to a generic viewer.
Dag { source: Source, reason: String },
Dag {
source: Source,
reason: String,
created_at: DateTime<Utc>,
},
}
/// How a hive-c0re node describes itself to a generic graph viewer.

View file

@ -152,9 +152,6 @@ pub struct GraphNode {
/// What must hold before this node runs. Omitted when empty.
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub deps: Vec<GraphDep>,
/// When the node was inserted into the graph. Always present — a node that
/// exists was created, so unlike the two below this is not an `Option`.
pub created_at: DateTime<Utc>,
/// When the node entered `Running`. `None` until it starts; a node that
/// never ran keeps `None`.
#[serde(default, skip_serializing_if = "Option::is_none")]
@ -257,7 +254,6 @@ fn wire_node<N: WireNode, R: WireResource>(node: &hive_jobq::Node<N, R>) -> Grap
parent: node.parent.map(NodeId::get),
state: node.state,
deps: node.deps.iter().map(wire_dep).collect(),
created_at: node.created_at,
started_at: node.started_at,
finished_at: node.finished_at,
error: node.error.clone(),
@ -447,7 +443,6 @@ mod tests {
parent,
state,
deps,
created_at: chrono::DateTime::<chrono::Utc>::default(),
started_at: None,
finished_at: None,
error: None,

View file

@ -255,9 +255,9 @@ impl State {
}
/// Wall-clock UTC now — the source for node lifecycle timestamps
/// ([`Node::created_at`] / [`Node::started_at`] / [`Node::finished_at`]). The
/// graph stamps its own timestamps rather than threading a clock through every
/// call, so a node's timing is self-contained. Derived from `SystemTime` (the workspace `chrono`
/// ([`Node::started_at`] / [`Node::finished_at`]). The graph stamps its own
/// timestamps rather than threading a clock through every call, so a node's
/// timing is self-contained. Derived from `SystemTime` (the workspace `chrono`
/// carries no `clock` feature, matching `hive_sh4re::wire_time`), truncated to
/// whole seconds; a pre-epoch or out-of-range clock clamps to the epoch.
fn now_utc() -> DateTime<Utc> {
@ -291,14 +291,6 @@ pub struct Node<N, R> {
pub deps: Vec<Dep<R>>,
/// Lifecycle state.
pub state: State,
/// UTC instant the node was inserted into the graph. Stamped by the graph.
///
/// Not an `Option`, unlike its two siblings below: starting and finishing
/// are events that may never happen, but a node that exists was created.
/// Together the three are the node's whole lifecycle — a caller asking
/// "how long did this sit before it ran" needs this end of the interval,
/// and previously had to get it from a wrapper the graph knows nothing about.
pub created_at: DateTime<Utc>,
/// UTC instant the node entered [`State::Running`] (`None` until it starts;
/// a cancelled node never ran, so it stays `None`). Stamped by the graph.
pub started_at: Option<DateTime<Utc>>,
@ -489,7 +481,6 @@ impl<N, R> Graph<N, R> {
payload,
deps,
state: State::Pending,
created_at: now_utc(),
started_at: None,
finished_at: None,
error: None,
@ -943,7 +934,6 @@ mod tests {
when: DepWhen::AFTER_OK,
}],
state: State::Pending,
created_at: now_utc(),
started_at: None,
finished_at: None,
error: None,
@ -964,7 +954,6 @@ mod tests {
payload: "x",
deps: vec![],
state: State::Pending,
created_at: now_utc(),
started_at: None,
finished_at: None,
error: None,