Repository navigation
feat(events): give a session an event envelope carrying run lifecycle - #794
justintime4tea wants to merge 10 commits into
Conversation
|
e5ce350 to
8dbd0bb
Compare
aa3e3ba to
4f5f612
Compare
Adopting the run's id in persistence waits for #780 because that's where every run's id gets minted. Resume, in #785 and #760, however then has to bind a resumed run's persistence to the parked run's id. Just a point of clarification regarding the deferal. |
|
Thanks for the clarification. That distinction makes sense: this PR intentionally leaves orchestration persistence on its independently minted checkpoint/park-owner ID, while #780 will establish the runtime-owned run ID. Resume work in #785/#760 can then explicitly rebind resumed persistence state to the parked run’s ID, preserving identity across the park/resume boundary. I don’t see a review issue here. |
4f5f612 to
73c1fbb
Compare
cba59a3 to
245e703
Compare
dhable
left a comment
There was a problem hiding this comment.
Just had one thought on the use of the expanding enum type pattern in our code. It might be something to chat about at the dev sync meeting in the future.
| clippy::large_enum_variant, | ||
| reason = "agent events are nearly every event on a session's stream; \ | ||
| boxing them would add an allocation to each to save memory \ | ||
| only on the rare lifecycle ones" |
There was a problem hiding this comment.
LifecycleEvent seems like it's grown past what should be represented as an enum. Disabling the linting rule and the use of non_exhaustive suggests that we should probably be using a common trait, structs and type params to structure this code.
pub trait SessionEvent {}
pub struct Started {
/// The configured agent running, matching
/// [`AgentInfo::id`](crate::AgentInfo::id).
agent: String,
prompt: String,
#[serde(
rename = "timeout_ms",
default,
skip_serializing_if = "Option::is_none",
with = "crate::duration_ms::option"
)]
timeout: Option<Duration>,
#[serde(default)]
liveness: Liveness,
}
impl SessionEvent for Started {}
// .... additional lifecycle event structs ...
impl SessionEvent for Agent {}
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct SessionEvent<P: SessionEvent> {
pub session_id: SessionId,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub run_id: Option<RunId>,
pub seq: SequenceNumber,
pub at: Timestamp,
pub payload: P,
}
If you needed to also have a hierarchy of LifecycleEvents, you can add a second trait for them as well. Then use something like downcast_rs (or implement the pattern by hand) in order to safely cast the type into specific instances when you need them in code.
Just some thoughts on other ways to get strong typed objects without huge enum values everywhere.
There was a problem hiding this comment.
I like where your head is at @dhable! Good idea, I added something to the AURA Dev sync "chat" that should show up when we have our meeting to remind us to discuss.
https://chat.google.com/room/AAQAKVebBm4/hybZb20OIIw/hybZb20OIIw?cls=10
|
couple of comments |
| pub fn ends_run(&self) -> bool { | ||
| matches!( | ||
| self, | ||
| Self::Parked { .. } |
There was a problem hiding this comment.
ends_run() counts Parked as the end of the run, but a parked run resumes. Resume reloads the run id from the parked document (park/continuation.rs:133).
A consumer that closes the run on Parked will then get the resumed run's events for a run it thinks is over. Started has nothing linking a resumed run back to the parked one.
Whether a resumed run keeps its RunId is the #780 question.
There was a problem hiding this comment.
Parked ending the run is #778's call, and #785 settles what a resume is: a new run that continues the parked one (its POST /v1/runs/{id}/resume). The link you're pointing at was missing from Started, so fa9af6f adds continues: Option<RunId> there, in #785's shape; a start that continues nothing carries no such field, so an older payload still reads as a fresh run. Today's resume path reloads persistence's own id from the parked document (continuation.rs:133), not the envelope's; whether that id becomes the RunId is #780's question, as you say.
| /// name and attributes an event *within* a run. | ||
| #[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)] | ||
| #[serde(try_from = "uuid::Uuid", into = "uuid::Uuid")] | ||
| pub struct RunId(uuid::Uuid); |
There was a problem hiding this comment.
This adds a second RunId beside aura::orchestration::RunId, and the two disagree. This one rejects the nil UUID and that one accepts it. RunParked.run_id is also still a String.
So nothing enforces the module doc's "the envelope never carries a second id for the same run". An orchestrated run can end up with the envelope's v7 id and the persistence id in one stream.
Predates: the orchestration type accepting nil, and RunParked.run_id being a String.
There was a problem hiding this comment.
The disagreement goes away on #796 (f4e2779): rather than re-exporting the envelope's RunId over orchestration's, persistence's id gets its own type, PersistenceRunId, which also refuses nil. The sentence you quoted was overclaiming; 8c0a23f drops it. RunParked.run_id being a String predates this and waits on #780.
| #[derive(Clone, Debug, Serialize, Deserialize)] | ||
| #[serde(tag = "type", rename_all = "snake_case")] | ||
| #[non_exhaustive] | ||
| pub enum LifecycleEvent { |
There was a problem hiding this comment.
#[non_exhaustive] only helps Rust callers. On the wire, an unknown variant fails the whole SessionEvent.
A probe test confirms it. {"type":"suspended"}, "reason":"preempted", and "policy":"hibernate" all fail with "unknown variant". An older CLI then can't read that event's seq, so the next event looks like a gap.
A #[serde(other)] fallback variant would fix this, here and on RunCancelReason and LivenessPolicy.
Predates: no enum in aura-events has a fallback today, AgentEventPayload included.
There was a problem hiding this comment.
Fixed in 2a25276: every enum the envelope carries now reads an unknown tag as Unknown, including ObserverKind and DetachCause (which were not non_exhaustive before, and now are) and AgentEventPayload, since a new agent event would have failed the envelope the same way. Unknown refuses to serialize, so a relay has to forward the bytes it received. SessionEventPayload's kind stays closed: serde cannot fall back on an adjacent tag that carries content, and the two kinds are not expected to grow. ToolOutcome has the same gap and predates this PR; filed separately as #805.
| /// One event in a session's ordered stream. | ||
| #[derive(Clone, Debug, Serialize, Deserialize)] | ||
| pub struct SessionEvent { | ||
| pub session_id: SessionId, |
There was a problem hiding this comment.
session_id is required, but runs without a session still exist. AgentRuntimeConfig.session_id is Option<String>, and the removed aura::config::SessionId doc said "not every run has one". A producer for those runs would have to invent an id or share one.
#778 also disagrees with itself here. Its summary says session-keyed, but its data-structures section still shows session_id: Option<SessionId>.
There was a problem hiding this comment.
Required is the rebase's call (the "stream is the session's" bullet), and #780 answers who supplies it: the runtime finds or creates the session when a run starts. Every ingress already has one today — chat mints one if the client sent none (fa3861d), A2A uses the context id, Slack the thread — and there is no producer for a library caller yet. The contradiction you found is the issue's stale "as implemented in PR #730" block, which still showed the pre-rebase Option<SessionId>; I have replaced it with the current shape.
| string_newtype! { | ||
| /// Identifier for a session: an agent's identity over time, and the key | ||
| /// its [`SessionEvent`](run::SessionEvent) stream is ordered under. | ||
| SessionId |
There was a problem hiding this comment.
SessionId accepts "", and a probe parses "session_id": "". As the stream key, that puts every run with an empty id into one stream and one sequence. fa3861d just made a blank chat session id count as missing.
ObserverId and CheckpointRef use nonempty_string_newtype, but the stream key doesn't.
Predates: aura::config::SessionId already accepted "".
|
|
||
| /// Whether this is the position immediately after `previous`. `false` | ||
| /// means at least one position between them was missed. | ||
| pub fn follows(self, previous: Self) -> bool { |
There was a problem hiding this comment.
The doc says false means a position was missed. A duplicate (5 after 5) and a reordered event (4 after 5) also return false, so a consumer following the doc would resync on a redelivery.
The doc should cover those cases, or the check should tell them apart.
There was a problem hiding this comment.
| impl $name { | ||
| pub fn new(value: impl Into<String>) -> Result<Self, $crate::EmptyId> { | ||
| let value = value.into(); | ||
| if value.is_empty() { |
There was a problem hiding this comment.
nit: this only refuses "". A probe accepts " " as an ObserverId and " " as a CheckpointRef.
fa3861d treats a blank id as missing, so checking trim().is_empty() here would match it.
| )] | ||
| timeout: Option<Duration>, | ||
| #[serde(default)] | ||
| liveness: Liveness, |
There was a problem hiding this comment.
nit: "liveness": null fails with "invalid type: null, expected struct Liveness", while "timeout_ms": null parses.
A producer that writes absent optionals as null would break the older-payload fallback. A null-tolerant deserializer that falls back to the default would fix it.
There was a problem hiding this comment.
Fixed in 778eaa4: Started.liveness reads null as the default, as it already read a missing field. timeout_ms tolerated null already, being an Option.
| /// Generates a string newtype that cannot hold the empty string: a fallible | ||
| /// constructor, `TryFrom` the string types, and deserialization that refuses | ||
| /// `""`. | ||
| macro_rules! nonempty_string_newtype { |
There was a problem hiding this comment.
nit: this repeats most of string_newtype!, and the two have already drifted. This one lacks PartialEq<String>. A shared inner macro for the common impls would keep them in step.
The saturating millisecond conversion is also written three times (duration_ms.rs:11, duration_ms.rs:30, lib.rs:449).
| /// | ||
| /// This is the epic's "unique agent id for that run". It is not | ||
| /// [`AgentContext::agent_id`], which is [`CONVERSATION_AGENT_ID`] or a worker | ||
| /// name and attributes an event *within* a run. |
There was a problem hiding this comment.
nit: this describes what AgentContext::agent_id holds ("CONVERSATION_AGENT_ID or a worker name"), which goes stale if its producers change. CLAUDE.md's drift rule puts that on agent_id itself.
AgentEvent says what an agent did and nothing about the run it belongs to. Observers reconstruct the run's story from the transport — a stream closing, an A2A status update, a [DONE] — and no two agree. aura_events::run adds SessionEvent: the session id, the run id when the event belongs to a run, a sequence number dense per session across its runs, a timestamp, and either an AgentEvent or a LifecycleEvent (started, observer attached and detached, claims exhausted, liveness decided, parked, finished, cancelled, failed). One stream per session means an agent's history loads as one stream, and a run boundary is a position in it rather than a second stream. The payload is adjacently tagged, kind beside event, so no field of an event can collide with the envelope's tag. Lifecycle events stay internally tagged by type, like AgentEventPayload; the roundtrip test over every variant is what catches a field named after a tag. Internal tagging and flatten need a self-describing format, which the module doc records. RunId is a UUID, v7 when minted, so it agrees with the orchestration run id HITL parses and the park owner key; aura-events takes uuid for it. SessionId and RunCancelReason move into aura-events, re-exported from aura::config and aura::hooks. RunCancelReason is non_exhaustive and gains Unclaimed and Shutdown, so a liveness or shutdown cancel is not reported as External. ObserverDetached says why with a DetachCause, and Started carries the run's Liveness, read as cancel-at-once when an older payload omits it. No producer emits it yet; this is the type and its tests. Ref: GH-778 Ref: GH-578
An identifier that names nothing should not be constructible. The ones this PR introduces now refuse their empty value on every path they can be built by: - ObserverId and CheckpointRef come from nonempty_string_newtype!, whose constructor returns a Result, with TryFrom the string types and deserialization that refuses "". - RunId refuses the nil UUID through TryFrom<Uuid>, FromStr and deserialization, and its parse error is InvalidRunId. - SequenceNumber refuses 0, since a session's sequence starts at 1. The string ids that predate this PR (ToolCallId, ToolName, ToolNamespace, SessionId) keep string_newtype!; hardening them reaches well past the envelope. Ref: GH-778
Timestamp, SequenceNumber, and the duration_ms serde module are not specific to a session's envelope, so they move from the run module to the crate root, beside the other shared value types. duration_ms becomes a public module so other crates can serialize a Duration as whole milliseconds without keeping their own copy, and it gains tests of its own, including the Option form and saturation past u64. SequenceNumber's docs describe any dense sequence starting at 1; the per-session rules stay on SessionEvent and the run module doc. That doc now notes that `at` is Unix milliseconds while other instants on the wire are RFC 3339 strings. Ref: GH-778
The CLI kept its own duration_millis module, the same whole-millisecond encoding aura-events now exports as duration_ms. Saved REPL events use the shared module instead. The format is unchanged for any duration that fits a u64 of milliseconds; a longer one saturates instead of being written as a number the reader cannot parse back. Ref: GH-778
245e703 to
5db77be
Compare
A reader older than its producer failed the whole SessionEvent on a variant it did not know, losing the event's seq, so the next event looked like a gap. The lifecycle enums, the observer and detach enums, the cancel reason, the liveness policy, and AgentEventPayload now read any other tag as an Unknown variant that cannot be written back, so a relay forwards the bytes it received rather than a re-encoded event. ObserverKind and DetachCause become non_exhaustive to match. SessionEventPayload keeps a closed kind: serde cannot fall back on an adjacent tag that carries content, and its two kinds are not expected to grow. Ref: GH-778
Started.liveness defaulted only when the field was missing; an explicit null failed the event, so a producer writing absent optionals as null broke the older-payload fallback. A null now reads as the default, as a missing field does. Ref: GH-778
`follows` answers whether a position is exactly one past the previous one, so a repeat and an out-of-order arrival return false as a gap does. The doc named only the gap, which would have a consumer resync on a redelivery. Ref: GH-778
A park ends its run, and a resume is a new run that continues the parked one, in the shape GH-785's resume route gives it. Nothing on the stream linked the two, so a consumer that closed a run on Parked could not tell the resumed run's Started from a fresh one. Started now carries the run it continues, when it continues one; a start that does not carries no such field, so an older payload reads as a fresh run. Ref: GH-778 Ref: GH-785
| /// Any other value on the wire, read by a version of this crate that does | ||
| /// not know it. It cannot be written back. |
There was a problem hiding this comment.
Variant comments describe runtime behavior
The new Unknown comment explains how serde reads and writes values. CLAUDE.md requires type comments to describe only what the value means, not runtime behavior. Remove the read/write clauses. The same wording appears on ObserverKind::Unknown and LifecycleEvent::Unknown, with the pattern repeated elsewhere in run.rs.
This repository requirement must be satisfied before merging.
Context Used: CLAUDE.md (source)
Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!
Only Finished carried the run's cumulative provider-billed tokens, so a consumer reading cost from the terminal event got nothing for a run that was cancelled, failed, or parked after billing turns. Parked, Cancelled, and Failed now carry the same flattened usage as Finished. It is required rather than optional: a flattened Option reads a malformed usage as None and hides a producer bug, and a producer always has a cumulative figure, zero included. Ref: GH-778
RunParked, an agent event the envelope carries, names its run by orchestration persistence's own id until the runtime mints both, so the module doc overclaimed. A run still has one RunId on the envelope. Ref: GH-778
feat(events): give a session an event envelope carrying run lifecycleaura_events::runaddsSessionEvent: the session id, the run id when the event belongs to a run, a sequence number dense per session across its runs, a timestamp, and either anAgentEventor aLifecycleEvent(started, observer attached and detached, claims exhausted, liveness decided, parked, finished, cancelled, failed). No producer emits it yet; #779 is the first.The items #778 lists to settle before a producer exists:
session_idis required,run_idoptional,seqdense per session.ends_run()says whether an event ends its run; the stream goes on.RunIdis a UUID, v7 when minted withRunId::mint().aura-eventstakesuuidfor it (the CLI already depends on the same workspace crate).kindbesideevent, so no event field can collide with the envelope's tag. Lifecycle events stay internally tagged bytype; the roundtrip test over every variant catches a field named after a tag. The module doc records that the family needs a self-describing format.RunCancelReasonmoves here fromaura::hooks(re-exported there), is#[non_exhaustive], and gainsUnclaimedandShutdown.ObserverDetachedcarries aDetachCause(Released,Expired,Displaced).Startedcarries the run'sLiveness, read as cancel-at-once when an older payload omits it.feat(events): refuse empty and nil identifiersIdentifiers that name nothing can't be built:
ObserverIdandCheckpointRefrefuse""(fallible constructor,TryFrom, deserialization),RunIdrefuses the nil UUID (TryFrom<Uuid>,FromStr, deserialization; errors areInvalidRunId), andSequenceNumberrefuses0. The string ids that predate this PR (ToolCallId,ToolName,ToolNamespace,SessionId) are unchanged.refactor(events): share timestamp, sequence number, and duration_msTimestamp,SequenceNumber(withZeroSequenceNumber), and theduration_msserde module aren't specific to the envelope, so they live at theaura-eventscrate root beside the other shared value types.duration_msis public, with tests of its own (includingOption<Duration>and saturation past au64of milliseconds). The run module doc notes thatatis Unix milliseconds while other instants on the wire, such as an approval'sexpires_at, are RFC 3339 strings.refactor(cli): use the shared duration_ms serde moduleThe CLI's
duration_milliswas the same encoding, so saved REPL events useaura_events::duration_msinstead. The format is the same for any duration that fits au64of milliseconds; a longer one saturates instead of being written as a number the reader can't parse back.aura-config'screated_atstays a plainu64(that crate doesn't depend onaura-events), and the A2A bus frame'sseqis left for #783.Verification
cargo +nightly fmt --check,cargo clippy --workspace --all-targets --all-features -- -D warnings, andcargo test --workspace(2,572 passed) onnightlyat3dfd2bbd.Ref: GH-778
Ref: GH-578