feature/durable turns and heartbeats - #2
Conversation
jvjvjv
commented
Aug 31, 2026
- fix: never discard an interrupted turn
- docs: spec, plan and openspec change for heartbeats and durable turns
- feat: emit heartbeats during idle provider reads
- feat: forward stream heartbeats to the browser as comment frames
- feat: add turn run and turn event models
- feat: add TurnRunStore for detached turn events
- feat: run a conversation turn as a queued job
- feat: add resumable turn event reader
- fix: page the turn reader's terminal drain instead of reading once
- feat: add dispatchTurn, resumeTurn and cancelTurn
- feat: add ai:prune-turn-events command
- fix: page the turn-event prune sweep and make run terminality explicit
- fix: dispatch turn jobs after commit and correct the resumeTurn docblock
- docs: document heartbeats and dispatched turns for 0.15.0
- chore: check off completed durable-turns-and-heartbeats tasks
- fix: keep heartbeats out of stored turn events and fail runs on provider errors
- docs: record deferred review findings for the durable-turns change
Persist the assistant message when a turn produced text, reasoning, or tool calls, and record the turn even when it was cut short before producing any of them. Report a cut-short turn as stop_reason 'incomplete' and log it as AiInteractionStatus::Aborted rather than a clean success, while still counting the tokens it burned towards conversation usage. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
The turn runner intercepts Heartbeat stream events and yields a
`{type: heartbeat}` turn event: never logged, never appended to the
stored event array, never resetting the step clock — but checked by the
max-duration guard, so a stalled stream can now trip it. SseFrameEncoder
renders it as an SSE comment (`: ping`) that every consumer ignores
without handling, while putting a byte on the wire so a dead connection
is detected.
Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Page the ai:prune-turn-events deletion so a first sweep over a large backlog never binds an unbounded id list into one delete, make AiTurnRunStatus::isTerminal() enumerate terminal cases explicitly so a future non-terminal case fails loudly instead of being pruned, and pin that retention_days of 0 disables pruning entirely. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Dispatch RunConversationTurnJob with afterCommit() so a host wrapping dispatchTurn() in its own transaction with a Redis/SQS queue cannot have a worker pick the job up before the run row commits, and rewrite the resumeTurn() docblock: synthesized heartbeat and max_stream_duration events carry no _seq, so a host encoder must read that key conditionally. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
The published client records the SSE id of each frame through a new onSequence callback so a reconnecting caller can resume a dispatched turn from where it left off; the README gains the dispatched-turn walkthrough, the turns.* config keys, and the fifth scheduled job; the 0.15.0 changelog entry gains the heartbeat and dispatched-turn features and drops the abandoned-connection known issue heartbeats fixed. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_015NAV1EaRmCUBd1Tfc7HuhZ
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
…der errors A detached turn's job stored every event continueConversation() yielded, including the heartbeats the runner emits per gateway tick. SseFrameEncoder never transmits a heartbeat (it becomes an SSE comment with no id), so each stored beat burned a sequence number a reader never saw — a long silent gap left a reconnecting client's after cursor behind the real sequence, replaying events it already had. The job now consumes heartbeats instead of storing them; TurnEventStream synthesizes its own beat whenever the store is quiet, which is what a reader actually needs. The job also finished a run Completed when the turn failed: provider failures are caught inside continueConversation() and yielded as a terminal error event, so the generator ends normally. The run's status then contradicted the interaction log and the wire. The job now tracks whether an error event was appended and finishes the run Failed with the provider's message, with an explicit precedence: error event, then stop request, then Completed. README: note that turns.max_stream_seconds bounds only a resumeTurn() read — generation in the worker is capped per provider request by conversations.max_stream_seconds, which large-context local models may need raised — and fix the dispatchTurn() snippet passing a Stringable where a string is declared, a TypeError under strict_types. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
There was a problem hiding this comment.
🟡 Changes recommended
There are a few confirmed correctness edge cases (notably _seq construction and heartbeat parser detach fallback) that can silently break resumable replay or drop provider output.
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Pull request overview
This PR introduces “durable turns” (turn execution detached into a queued job with resumable replay) and “heartbeats” (keep-alive ticks during provider/store silence) to prevent interrupted turns from being discarded, improve abandonment detection, and support browser reconnect/resume semantics.
Changes:
- Add heartbeat emission during idle provider SSE reads and forward heartbeats to browsers as SSE comment frames.
- Add durable turn infrastructure:
AiTurnRun/AiTurnEventmodels + store/reader + queued job execution + retention pruning. - Update frontend contract/docs to support resumable SSE
id:sequencing and new stop/incomplete semantics.
File summaries
| File | Description |
|---|---|
| tests/Feature/TurnRunStoreTest.php | Coverage for sequencing, cancellation/abandonment, throttled stop checks, and finish recording. |
| tests/Feature/TurnEventStreamTest.php | Coverage for replay/resume behavior, terminal drain, SSE id framing, and polling touch behavior. |
| tests/Feature/StreamTranslatorTest.php | Ensures stop_reason becomes incomplete when a stream never cleanly ended. |
| tests/Feature/RunConversationTurnJobTest.php | Verifies job-driven turn execution records events and sets run terminal state correctly. |
| tests/Feature/ReasoningOpenAiCompatibleGatewayTest.php | Exercises idle-gap heartbeat behavior and parser correctness across split frames. |
| tests/Feature/PruneTurnEventsCommandTest.php | Validates retention pruning deletes terminal runs/events while preserving live runs and supports “disabled” retention. |
| tests/Feature/ConversationUsageServiceTest.php | Ensures aborted turns still count toward token usage totals. |
| tests/Feature/ChatTurnLibraryTest.php | Verifies heartbeat is encoded as an SSE comment frame and stream still terminates normally. |
| tests/Feature/AiTurnRunModelTest.php | Covers run public id generation, status terminality, and event uniqueness/payload casting. |
| tests/Feature/AiPersonaConversationServiceTest.php | Adds tests for dispatch/resume/cancel APIs, heartbeat handling, and interrupted-turn persistence semantics. |
| src/Services/RawExchange/TeeingStream.php | Adds record() to capture bytes read directly from detached resources. |
| src/Services/LaravelAi/StreamTranslator.php | Tracks open-stream state to return stop_reason: incomplete when appropriate. |
| src/Services/LaravelAi/Streaming/Heartbeat.php | New stream-event type used to surface heartbeat ticks through laravel/ai streaming. |
| src/Services/LaravelAi/ReasoningOpenAiCompatibleGateway.php | Uses idle-read heartbeat parsing and forwards Heartbeat events through processing. |
| src/Services/LaravelAi/Concerns/HeartbeatsIdleSseReads.php | Implements timeout-bounded SSE reading with a partial-line buffer and heartbeat yields. |
| src/Services/ConversationUsageService.php | Counts both success and aborted interaction logs when summarizing usage. |
| src/Services/Conversation/TurnRunStore.php | Implements single-writer event sequencing, cancellation/abandonment logic, and throttled stop checks. |
| src/Services/Conversation/TurnEventStream.php | Streams stored events with replay/resume cursoring, terminal drain, heartbeat synthesis, and read ceiling. |
| src/Services/ChatBot/SseFrameEncoder.php | Encodes heartbeats as SSE comments and _seq as SSE id: while stripping from payload. |
| src/Services/ChatBot/Conversation/TurnRecorder.php | Persists interrupted turns (even empty), adds incomplete metadata, and logs aborted status appropriately. |
| src/Services/ChatBot/Conversation/ConversationTurnRunner.php | Forwards heartbeat events without logging/storing them and ensures guards can evaluate during silence. |
| src/Services/ChatBot/ChatBotPresenter.php | Surfaces incomplete on transcript rows based on stored message metadata. |
| src/Services/AiPersonaConversationService.php | Adds dispatchTurn(), resumeTurn(), and cancelTurn() durable-turn APIs. |
| src/Models/AiTurnRun.php | New model for durable turn runs with status/timestamp casts and ordered events relation. |
| src/Models/AiTurnEvent.php | New model for durable turn events with payload casting and created-at defaulting. |
| src/Jobs/RunConversationTurnJob.php | Queued execution of turns that records events, skips heartbeat storage, and finalizes run state. |
| src/Enums/AiTurnRunStatus.php | New status enum with isTerminal() used by reader/job/store semantics. |
| src/Enums/AiInteractionStatus.php | Adds Aborted status to represent caller-hung-up turns. |
| src/Console/Commands/PruneTurnEventsCommand.php | Adds ai:prune-turn-events retention command with paging. |
| src/CodeTalkerServiceProvider.php | Registers and schedules turn-event pruning alongside existing scheduled jobs. |
| resources/js/types/code-talker.d.ts | Updates client contract: incomplete transcript field, StopReason includes incomplete, documents heartbeat omission from wire union. |
| resources/js/code-talker-stream.ts | Adds onSequence callback by parsing SSE id: for resume support. |
| README.md | Documents detached turns, heartbeat behavior/encoding, and interrupted turn semantics. |
| openspec/changes/durable-turns-and-heartbeats/tasks.md | Implementation task checklist for the change set. |
| openspec/changes/durable-turns-and-heartbeats/specs/chat-turn-library/spec.md | Spec updates for heartbeat + resumable SSE id framing + durable turns lifecycle. |
| openspec/changes/durable-turns-and-heartbeats/proposal.md | Proposal describing motivation, constraints, and scope. |
| openspec/changes/durable-turns-and-heartbeats/design.md | Design summary referencing deeper design/plan docs. |
| openspec/changes/durable-turns-and-heartbeats/deferred-findings.md | Deferred review findings tracked for later follow-up. |
| docs/superpowers/specs/2026-08-31-durable-turns-and-heartbeats-design.md | Detailed design document for heartbeats + durable turns. |
| docs/superpowers/plans/2026-08-31-durable-turns-and-heartbeats.md | Detailed implementation plan and constraints checklist. |
| database/migrations/2026_08_31_000002_create_ai_turn_events_table.php | Adds ai_turn_events schema with unique (run, sequence). |
| database/migrations/2026_08_31_000001_create_ai_turn_runs_table.php | Adds ai_turn_runs schema for durable turn lifecycle. |
| config/code-talker.php | Adds conversations.heartbeat_seconds and new turns.* config block. |
| CHANGELOG.md | Documents 0.15.0 breaking changes, new features, and bug fixes for durable turns/heartbeats. |
Review details
Suppressed comments (1)
src/Services/Conversation/TurnEventStream.php:67
- Same
_seqoverride issue as above: using array union with$event->payloadon the left allows an_seqkey in the stored payload to mask the actual sequence, breaking resumable replay cursors.
/** @var AiTurnEvent $event */
$after = $event->sequence;
yield $event->payload + ['_seq' => $event->sequence];
}
- Files reviewed: 43/44 changed files
- Comments generated: 3
- Review effort level: Lite
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
| // Raw-exchange capture tees every byte the parser reads. Detaching | ||
| // takes the resource out from under the tee, so the reader below | ||
| // feeds it the bytes itself — otherwise enabling heartbeats would | ||
| // silently blank ai_provider_exchanges.raw_response. | ||
| $tee = $streamBody instanceof TeeingStream ? $streamBody : null; | ||
|
|
||
| $resource = $streamBody->detach(); | ||
|
|
||
| if (! is_resource($resource)) { | ||
| return; | ||
| } |
| $after = $event->sequence; | ||
| $lastEmittedAt = microtime(true); | ||
|
|
||
| yield $event->payload + ['_seq' => $event->sequence]; |
| RunConversationTurnJob::dispatch($run->id)->afterCommit(); | ||
|
|
||
| return $run; |