MCPcopy Create free account
hub / github.com/atomicdotdev/atomic / handle_turn_end

Method handle_turn_end

atomic-agent/src/turn/orchestrator/turn.rs:128–448  ·  view source on GitHub ↗

Handle a TurnEnd event (Stop). Records an Atomic change for the turn (status → add → record), then transitions the session back to Idle. The recording workflow lets the repository figure out what changed: 1. `repo.status()` — find modified, deleted, and untracked files 2. `repo.add()` — track any new files the agent created 3. `repo.record(all: true)` — record everything that's dirty This avoid

(
        &mut self,
        mut event: TurnEvent,
    )

Source from the content-addressed store, hash-verified

126 /// snapshots between TurnStart and TurnEnd. Instead, we ask the
127 /// repository what changed since the last recorded state.
128 pub(super) async fn handle_turn_end(
129 &mut self,
130 mut event: TurnEvent,
131 ) -> AgentResult<DispatchResult> {
132 // Clone so `event` stays mutable for enrichment below while the
133 // id is borrowed throughout this function.
134 let session_id_owned = event.session_id.clone();
135 let session_id = session_id_owned.as_str();
136 let _turn_end_lock = match self.try_turn_end_lock(session_id) {
137 TurnEndLock::Acquired(guard) => Some(guard),
138 TurnEndLock::Busy => {
139 log::warn!(
140 "Turn end for session {} is already being recorded; skipping duplicate Stop hook",
141 session_id
142 );
143 return Ok(DispatchResult::new(session_id, phase::Phase::Idle)
144 .with_warning("duplicate Stop hook skipped: turn already recording"));
145 }
146 TurnEndLock::Unavailable => None,
147 };
148
149 // The session lock deduplicates one session's Stops. Independent
150 // sessions (including sandboxes) still share pristine and must not
151 // race status/add/record or checkpoint publication. Acquire before
152 // reading session/status, and retain through publication and save.
153 let _publication_lock = self.wait_turn_publication_lock(session_id)?;
154
155 let mut session = self.load_or_create_session(session_id, &event)?;
156
157 // Tool hooks can reserve and populate the next journal turn without
158 // a TurnStart (e.g. programmatic OpenCode/subagent turns). Consult that
159 // durable state before classifying an idle session's Stop as a retry.
160 if session.phase == phase::Phase::Idle {
161 if let Some(sink) = &self.journal_sink {
162 let pending = sink
163 .turn_status(session_id, session.turn_count.saturating_add(1))
164 .map_err(|reason| AgentError::ProvenanceJournalFailed {
165 session_id: session_id.to_owned(),
166 reason,
167 })?;
168 if let Some(pending) = pending {
169 match pending.lifecycle {
170 JournalTurnLifecycle::Running | JournalTurnLifecycle::Checkpointing => {
171 session.begin_turn();
172 let transition = phase::transition(
173 session.phase,
174 Event::TurnStart,
175 TransitionContext::default(),
176 );
177 phase::apply_common_actions(&mut session, &transition);
178 // Persist activation before recording or publication so
179 // a failed checkpoint retains the normal retry path.
180 self.session_store.save(&session)?;
181 }
182 lifecycle => {
183 return Err(AgentError::ProvenanceJournalFailed {
184 session_id: session_id.to_owned(),
185 reason: format!(

Callers 1

dispatchMethod · 0.80

Calls 15

transitionFunction · 0.85
apply_common_actionsFunction · 0.85
has_pending_changesFunction · 0.85
record_turnFunction · 0.85
try_turn_end_lockMethod · 0.80
is_turn_activeMethod · 0.80
set_transcript_pathMethod · 0.80
enrich_opencode_turnMethod · 0.80

Tested by

no test coverage detected