Module: TurnKit::Reconciliation
- Defined in:
- lib/turnkit/reconciliation.rb
Overview
Reconciles turns abandoned by a dead worker: atomically marks them stale, marks their unfinished tool executions interrupted, and appends synthetic error tool results so the persisted transcript stays structurally complete for continuation. TurnKit never reruns an interrupted tool; the synthetic result tells the continued model the outcome is unknown.
Constant Summary collapse
- INTERRUPTED_MESSAGE =
"Tool execution was interrupted before a result was recorded. " \ "It is unknown whether the operation ran; do not assume it did or did not."
Class Method Summary collapse
- .emit(type, turn, payload = {}) ⇒ Object
- .interrupt_tool_executions(turn) ⇒ Object
- .reconcile!(before:) ⇒ Object
- .repair_transcript(turn, executions) ⇒ Object
Class Method Details
.emit(type, turn, payload = {}) ⇒ Object
68 69 70 |
# File 'lib/turnkit/reconciliation.rb', line 68 def emit(type, turn, payload = {}) TurnKit.on_event&.call(Event.new(type: type, turn_id: turn.fetch("id"), conversation_id: turn.fetch("conversation_id"), payload: payload)) end |
.interrupt_tool_executions(turn) ⇒ Object
23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 |
# File 'lib/turnkit/reconciliation.rb', line 23 def interrupt_tool_executions(turn) store = TurnKit.store store.list_tool_executions(turn_id: turn.fetch("id")).map do |execution| next execution unless %w[pending running].include?(execution.fetch("status")) interrupted = store.claim_tool_execution( execution.fetch("id"), from: execution.fetch("status"), to: "interrupted", error: { "message" => "interrupted: worker terminated while the tool was executing" }, completed_at: Clock.now ) next execution unless interrupted emit("tool_call.interrupted", turn, id: interrupted.fetch("tool_call_id"), name: interrupted.fetch("tool_name"), tool_execution_id: interrupted.fetch("id")) interrupted end end |
.reconcile!(before:) ⇒ Object
15 16 17 18 19 20 21 |
# File 'lib/turnkit/reconciliation.rb', line 15 def reconcile!(before:) TurnKit.store.reconcile_stale_turns(before: before).each do |turn| emit("turn.stale", turn) executions = interrupt_tool_executions(turn) repair_transcript(turn, executions) end end |
.repair_transcript(turn, executions) ⇒ Object
42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 |
# File 'lib/turnkit/reconciliation.rb', line 42 def repair_transcript(turn, executions) store = TurnKit.store = store.(turn.fetch("conversation_id")) resolved = .select { || ["kind"] == "tool_result" } .flat_map { || ["content"].map { |part| part["tool_call_id"] } } .select { || ["turn_id"] == turn.fetch("id") && ["kind"] == "tool_call" } .flat_map { || ["content"].select { |part| part["type"] == "tool_call" } } .reject { |part| resolved.include?(part["id"]) } .each do |part| execution = executions.find { |candidate| candidate["tool_call_id"] == part["id"] } = store.( "conversation_id" => turn.fetch("conversation_id"), "turn_id" => turn.fetch("id"), "role" => "tool", "kind" => "tool_result", "content" => [ { "type" => "tool_result", "tool_call_id" => part["id"], "text" => { "error" => true, "message" => INTERRUPTED_MESSAGE }.to_json, "error" => true } ], "tool_execution_id" => execution&.fetch("id"), "metadata" => { "tool_name" => part["name"], "interrupted" => true } ) emit("message.created", turn, message_id: .fetch("id"), role: "tool", kind: "tool_result") end end |