-- SPDX-FileCopyrightText: © 2026 Vladimir Zorin -- SPDX-License-Identifier: LicenseRef-OWL-1.0-or-later -- Licensed under OWL v1.0+. See LICENSE. --- Core turn loop: orchestrates LLM calls, tool execution, event emission, --- cancellation, transcript updates, and normalized turn_result production. --- --- Exports: --- run(runtime, request, options) → turn_result | nil, agent_error local event = require("agent_smith.core.event") local cancel = require("agent_smith.cancel") local agent_error = require("agent_smith.error") local role = require("agent_smith.role") local llm = require("agent_smith.llm") local tool_call_mod = require("agent_smith.llm.tool.call") local tool_result_mod = require("agent_smith.llm.tool.result") local lev = require("lev") -- A safety bound for the tool loops local MAX_LOOP_ITERATIONS = 65535 -- How many times a delegated child (a role that must terminate via `report`) is -- re-prompted when it answers in prose instead of calling the terminal tool, -- before the dispatch worker falls back to salvaging that prose. local MAX_REPORT_NUDGES = 2 -- Corrective message appended when such a child skips the terminal `report` tool. local REPORT_NUDGE_TEXT = "You ended your turn without submitting a result. " .. "A plain-text reply is discarded and never reaches the caller — you MUST call the `report` tool now " .. "(with `status`, `summary`, and your `findings`/`changes`) as your final action." -- Retry fallbacks, used only when the effective config omits these settings. -- The documented defaults live in config/loader.lua; these mirror them so a -- standalone turn (e.g. a test driving turn.run directly) still retries sensibly. local DEFAULT_TRANSPORT_MAX_RETRIES = 2 local DEFAULT_EMPTY_RESPONSE_MAX_RETRIES = 2 local DEFAULT_EMPTY_RESPONSE_BACKOFF_MS = { 500, 1000 } local M = {} --- Aggregate usage delta into the running total. --- @param delta table — usage from one LLM response --- @param total table — accumulated usage local function aggregate_usage(delta, total) if type(delta) ~= "table" then return end total.input_tokens = (total.input_tokens or 0) + (delta.input_tokens or 0) total.output_tokens = (total.output_tokens or 0) + (delta.output_tokens or 0) total.total_tokens = (total.total_tokens or 0) + (delta.total_tokens or 0) end --- Load project memories for prompt construction. --- Returns {} when project scope or persistence repositories are unavailable, --- or when repository lookup/listing fails. --- @param services table --- @param active_turn table --- @return table local function load_project_memories(services, active_turn) if type(active_turn) ~= "table" or type(active_turn.project_id) ~= "string" or active_turn.project_id == "" then return {} end if type(services) ~= "table" then return {} end if type(services.persistence) ~= "table" or type(services.persistence.repositories) ~= "function" then return {} end local repos_ok, repos = pcall(services.persistence.repositories) if not repos_ok or type(repos) ~= "table" then return {} end if type(repos.project_memory) ~= "table" or type(repos.project_memory.list) ~= "function" then return {} end local list_ok, memories = pcall(repos.project_memory.list, active_turn.project_id) if not list_ok or type(memories) ~= "table" then return {} end return memories end --- Build the LLM request payload from resolved model, trimmed messages, and tool specs. --- @param selected_model table — { model = string, ... } --- @param model_messages table — array of message records --- @param tool_specs table or nil — LLM tool specs from tool service --- @param request table — original request (may carry stream flag) --- @param cancel_token table — cancel token to attach --- @return table llm_request local function build_llm_request(selected_model, model_messages, tool_specs, request, cancel_token) local req = { model = selected_model.model, messages = model_messages, stream = request.stream == true, } if tool_specs and #tool_specs > 0 then req.tools = tool_specs end if cancel_token then req.cancel = cancel_token.lev or cancel_token end return req end local function is_retryable_transport_error(err) local normalized = agent_error.normalize(err, { origin = "turn" }) if normalized.kind == agent_error.kind.transport then return true, normalized end if normalized.code == "http_error" then local status = type(normalized.data) == "table" and normalized.data.status or nil if type(status) == "number" and (status >= 500 or status == 429) then return true, normalized end end return false, normalized end local function retry_backoff_ms(backoff_ms, retry_index) if type(backoff_ms) ~= "table" or #backoff_ms == 0 then return 0 end local value = backoff_ms[retry_index] or backoff_ms[#backoff_ms] or 0 if type(value) ~= "number" or value < 0 then return 0 end return value end local function sleep_with_backoff(backoff_ms, cancel_token) if type(backoff_ms) ~= "number" or backoff_ms <= 0 then return true end lev.sleep(backoff_ms / 1000) if cancel.is_cancelled(cancel_token) then return nil, agent_error.cancelled("turn") end return true end local function has_assistant_payload(response) local message = response and response.message or nil if type(message) ~= "table" then return false end local content = message.content local has_content = type(content) == "string" and content ~= "" local has_reasoning = type(message.reasoning_text) == "string" and message.reasoning_text ~= "" local tool_calls = tool_call_mod.extract_tool_calls(message) local has_tool_calls = type(tool_calls) == "table" and #tool_calls > 0 return has_content or has_reasoning or has_tool_calls end --- Execute a single non-streaming LLM call and return the normalized response. --- @param llm_config table — LLM client configuration --- @param llm_request table — built request payload --- @param event_bus table — per-turn event bus --- @param turn_id string — current turn id --- @return table|nil response, table|nil error local function call_llm_non_streaming(llm_config, llm_request, event_bus, turn_id) local response, err = llm.chat(llm_config, llm_request) if err then if agent_error.is_cancelled(err) then event_bus:emit(event.event("llm_request_cancelled", turn_id, { model = llm_request.model, })) end return nil, err end return response end --- Execute a streaming LLM call and return the assembled final response. --- @param llm_config table — LLM client configuration --- @param llm_request table — built request payload --- @param event_bus table — per-turn event bus --- @param turn_id string — current turn id --- @return table|nil response, table|nil error local function call_llm_streaming(llm_config, llm_request, event_bus, turn_id, stream_state) -- The client's stream assembler accumulates text/reasoning/tool deltas and -- always returns the fully assembled response; the handlers here only drive -- the UI event stream and the retry-policy output_started flag. local mark_output_started = function() if type(stream_state) == "table" then stream_state.output_started = true end end local handlers = { on_text = function(delta) mark_output_started() event_bus:emit(event.event("llm_text_delta", turn_id, { text = delta })) end, on_reasoning = mark_output_started, on_tool_delta = mark_output_started, } local response, err = llm.chat_stream(llm_config, llm_request, handlers) if err then if agent_error.is_cancelled(err) then event_bus:emit(event.event("llm_request_cancelled", turn_id, { model = llm_request.model, })) end return nil, err end if not response or type(response.message) ~= "table" then return nil, agent_error.new({ kind = agent_error.kind.protocol, code = "invalid_stream_response", message = "stream completed without an assembled response message", origin = "turn", }) end return response end local function call_llm_with_transport_retries(llm_config, llm_request, event_bus, turn_id, cancel_token) local max_retries = type(llm_config.retry_max_retries) == "number" and llm_config.retry_max_retries or DEFAULT_TRANSPORT_MAX_RETRIES local backoff_ms = llm_config.retry_backoff_ms local stream_retry_after_output = llm_config.stream_retry_after_output == true local max_attempts = max_retries + 1 local attempt = 1 while attempt <= max_attempts do local response, llm_err local stream_state = { output_started = false } if llm_request.stream then response, llm_err = call_llm_streaming(llm_config, llm_request, event_bus, turn_id, stream_state) else response, llm_err = call_llm_non_streaming(llm_config, llm_request, event_bus, turn_id) end if not llm_err then return response end if agent_error.is_cancelled(llm_err) then return nil, llm_err end local retryable, normalized_err = is_retryable_transport_error(llm_err) if not retryable then return nil, normalized_err end local retry_index = attempt local next_attempt = attempt + 1 local phase = "request" if llm_request.stream and stream_state.output_started then phase = "post_output" if not stream_retry_after_output then event_bus:emit(event.event("llm_request_retry", turn_id, { retry_index = retry_index, max_retries = max_retries, attempt = next_attempt, max_attempts = max_attempts, reason = normalized_err.message, error_code = normalized_err.code, stream = true, phase = phase, })) return nil, normalized_err end elseif llm_request.stream then phase = "pre_output" end if retry_index > max_retries then return nil, normalized_err end event_bus:emit(event.event("llm_request_retry", turn_id, { retry_index = retry_index, max_retries = max_retries, attempt = next_attempt, max_attempts = max_attempts, reason = normalized_err.message, error_code = normalized_err.code, stream = llm_request.stream == true, phase = phase, })) local wait_ms = retry_backoff_ms(backoff_ms, retry_index) local _, wait_err = sleep_with_backoff(wait_ms, cancel_token) if wait_err then return nil, wait_err end attempt = next_attempt end return nil, agent_error.new({ kind = agent_error.kind.transport, code = "request_failed", message = "request failed", origin = "turn", }) end --- Execute tool calls deterministically with budget tracking. --- @param tool_calls table — array of extracted tool calls --- @param tool_service table — tool service handle --- @param tool_context table — execution context --- @param remaining_budget number|nil — remaining tool step budget (nil = unlimited) --- @param event_bus table — per-turn event bus --- @param turn_id string — current turn id --- @param conversation_svc table — conversation service --- @param conversation_state table — conversation state --- @return table results — array of { call, result, tool_result_message } --- @return number|nil remaining_budget — updated budget --- @return boolean budget_exhausted — true if budget was exhausted --- @return boolean cancelled — true if cancellation was detected --- @return table|nil append_err — transcript append failure (fails the turn) local function execute_tool_calls( tool_calls, tool_service, tool_context, remaining_budget, event_bus, turn_id, conversation_svc, conversation_state ) local results = {} local budget_exhausted = false local was_cancelled = false -- A tool result that cannot be appended leaves the assistant's tool_calls -- dangling, which fails wire validation one LLM call later with a confusing -- schema error — so surface the append failure as the turn error instead. local append_result = function(msg) local ok, err = conversation_svc.append_tool_result(conversation_state, msg) if not ok then return agent_error.normalize(err, { origin = "turn" }) end return nil end for _, call in ipairs(tool_calls) do -- Check cancellation before each tool if cancel.is_cancelled(tool_context.cancel_token) then -- Produce synthetic cancelled result for this and remaining calls local synth_msg = tool_result_mod.build_tool_result_message(call.id, '{"error":"cancelled","kind":"cancelled"}') local append_err = append_result(synth_msg) if append_err then return results, remaining_budget, budget_exhausted, was_cancelled, append_err end results[#results + 1] = { call = call, result = { call_id = call.id, ok = false, cancelled = true }, tool_result_message = synth_msg, } was_cancelled = true elseif type(remaining_budget) == "number" and remaining_budget <= 0 then -- Produce synthetic budget-exhausted result local synth_msg = tool_result_mod.build_tool_result_message( call.id, '{"error":"budget_exhausted","kind":"budget_exhausted"}' ) local append_err = append_result(synth_msg) if append_err then return results, remaining_budget, budget_exhausted, was_cancelled, append_err end results[#results + 1] = { call = call, result = { call_id = call.id, ok = false, cancelled = false, error = agent_error.new({ kind = agent_error.kind.budget_exhausted, code = "budget_exhausted", message = "tool step budget exhausted", origin = "turn", }), }, tool_result_message = synth_msg, } budget_exhausted = true else -- Real tool execution local exec_options = { event_bus = event_bus, } local tool_result = tool_service:execute(call, tool_context, exec_options) if type(remaining_budget) == "number" then remaining_budget = remaining_budget - 1 end -- Build and append tool result message local output = tool_result.content or "" if tool_result.ok == false and tool_result.error then output = tool_result.error.message or "tool error" end local result_msg = tool_result_mod.build_tool_result_message(call.id, output) -- A tool may declare that its (often large) result is recoverable from -- a durable store under some key (e.g. dispatch persists to the -- scratchpad). Carry the key onto the message so context-window elision -- can point the model at it instead of dropping the content silently. if type(tool_result.metadata) == "table" and tool_result.metadata.recall_key then result_msg.recall_key = tool_result.metadata.recall_key end local append_err = append_result(result_msg) if append_err then return results, remaining_budget, budget_exhausted, was_cancelled, append_err end results[#results + 1] = { call = call, result = tool_result, tool_result_message = result_msg, } if tool_result.cancelled then was_cancelled = true end end end return results, remaining_budget, budget_exhausted, was_cancelled end --- Main turn entry point. --- @param runtime table — runtime instance --- @param request table — { user_message, role, model_override, stream, conversation_state, tool_step_budget } --- @param options table or nil — { llm_config } --- @return table turn_result on success, or nil, agent_error on failure function M.run(runtime, request, options) options = options or {} return runtime:with_turn(request, function(active_turn) -- ── Step 9: Initialize per-turn state ── local event_bus = event.new() if type(options.on_event) == "function" then event_bus:subscribe(options.on_event) end active_turn.event_bus = event_bus local services = runtime:service_handles() local ctx = runtime:execution_context() local conversation_svc = services.conversation local conversation_state = request.conversation_state or conversation_svc.new_state() local usage_delta = { input_tokens = 0, output_tokens = 0, total_tokens = 0, } local dispatch_usage = { input_tokens = 0, output_tokens = 0, total_tokens = 0 } -- Usage from the most recent LLM response; the prompt shows context -- occupancy of the last response, not the per-turn running sum. local last_response_usage = nil local tool_calls_executed = {} local remaining_budget = request.tool_step_budget if remaining_budget == nil and runtime.cfg then remaining_budget = runtime.cfg.tool_step_budget end local stop_reason = nil local assistant_message = nil local turn_err = nil local detected_tool_calls = {} local detected_tool_call_ids = {} -- ── Step 10: Emit turn_started ── event_bus:emit(event.event("turn_started", active_turn.id, { role = request.role or ctx.active_role, conversation_id = ctx.conversation_id, cwd = active_turn.cwd, project_id = active_turn.project_id, })) -- ── Step 11: Resolve role and model ── local role_name = request.role or ctx.active_role local role_spec = role.resolve_role(role_name) if not role_spec then turn_err = agent_error.new({ kind = agent_error.kind.policy, code = "unknown_role", message = "unknown role: " .. tostring(role_name), origin = "turn", }) event_bus:emit(event.event("turn_failed", active_turn.id, { error = turn_err, })) event_bus:emit(event.event("turn_completed", active_turn.id, { stop_reason = "error", })) return { assistant_message = nil, tool_calls_executed = tool_calls_executed, usage_delta = usage_delta, events = event_bus:events(), stop_reason = "error", error = turn_err, cancelled = false, terminal_result = nil, } end local selected_model = services.model.select_model(request, role_spec, runtime.cfg) if not selected_model then turn_err = agent_error.new({ kind = agent_error.kind.policy, code = "no_model", message = "no model available for role: " .. role_name, origin = "turn", }) event_bus:emit(event.event("turn_failed", active_turn.id, { error = turn_err, })) event_bus:emit(event.event("turn_completed", active_turn.id, { stop_reason = "error", })) return { assistant_message = nil, tool_calls_executed = tool_calls_executed, usage_delta = usage_delta, events = event_bus:events(), stop_reason = "error", error = turn_err, cancelled = false, terminal_result = nil, } end -- ── Step 12: Build initial prompt and transcript state ── local project_memories = load_project_memories(services, active_turn) local system_prompt = services.prompt.build_system_prompt(active_turn, role_spec, project_memories) -- Everything appended from here on is this turn's transcript slice; the -- persistence layer writes it verbatim so the stored conversation keeps -- the true per-round message sequence. local turn_messages_start = #conversation_state.messages + 1 local append_ok, append_err = conversation_svc.append_user(conversation_state, request.user_message) if not append_ok then turn_err = agent_error.normalize(append_err, { origin = "turn" }) event_bus:emit(event.event("turn_failed", active_turn.id, { error = turn_err, })) event_bus:emit(event.event("turn_completed", active_turn.id, { stop_reason = "error", })) return { assistant_message = nil, tool_calls_executed = tool_calls_executed, usage_delta = usage_delta, events = event_bus:events(), stop_reason = "error", error = turn_err, cancelled = false, terminal_result = nil, } end event_bus:emit(event.event("prompt_built", active_turn.id, { model = selected_model.model, role = role_name, message_count = #conversation_state.messages, context_window = selected_model.context_window, })) -- ── LLM config from options or runtime ── local llm_config = options.llm_config or runtime.cfg.llm_config or {} -- ── Tool specs ── local tool_specs = services.tool:llm_specs(role_spec.name) -- ── Main LLM/tool loop ── local iteration = 0 -- Delegated child roles (worker/explorer) must terminate via the `report` -- terminal tool; user-facing roles (direct/orchestrator) legitimately end in -- prose. When a child ends in prose we re-prompt it up to MAX_REPORT_NUDGES -- times before letting the turn complete (the dispatch worker then salvages -- the prose rather than discarding the whole turn). local requires_report = not role_spec.user_facing local report_nudges = 0 while iteration < MAX_LOOP_ITERATIONS do iteration = iteration + 1 -- Check cancellation at loop top if cancel.is_cancelled(active_turn.cancel_token) then stop_reason = "cancelled" break end -- ── Step 13: Build and trim messages ── local model_messages = conversation_svc.build_model_messages(conversation_state, system_prompt) local cm = runtime.cfg.context_management or {} local trim_options = {} if cm.enabled ~= false then trim_options.keep_rounds = cm.tool_output_keep_rounds trim_options.elide_over_chars = cm.tool_output_elide_over_chars trim_options.trim_threshold_pct = cm.trim_threshold_pct trim_options.context_window = selected_model.context_window end local trimmed_messages, trim_meta = conversation_svc.trim_for_model(model_messages, trim_options) -- ── Step 14: Construct LLM request ── local llm_request = build_llm_request(selected_model, trimmed_messages, tool_specs, request, active_turn.cancel_token) local empty_response_cfg = runtime.cfg or {} local empty_response_max_retries = type(empty_response_cfg.empty_response_retry_max_retries) == "number" and empty_response_cfg.empty_response_retry_max_retries or DEFAULT_EMPTY_RESPONSE_MAX_RETRIES local empty_response_backoff_ms = empty_response_cfg.empty_response_retry_backoff_ms or DEFAULT_EMPTY_RESPONSE_BACKOFF_MS local empty_response_max_attempts = empty_response_max_retries + 1 -- ── Steps 15-17: LLM call ── local response, llm_err local empty_attempt = 1 while empty_attempt <= empty_response_max_attempts do -- ── Step 15: Emit llm_request_started ── event_bus:emit(event.event("llm_request_started", active_turn.id, { model = selected_model.model, stream = llm_request.stream, message_count = #trimmed_messages, tool_count = tool_specs and #tool_specs or 0, remaining_budget = remaining_budget, estimated_tokens = trim_meta and trim_meta.estimated_tokens or nil, dropped_messages = trim_meta and trim_meta.dropped_count or nil, elided_results = trim_meta and trim_meta.elided_count or nil, })) response, llm_err = call_llm_with_transport_retries( llm_config, llm_request, event_bus, active_turn.id, active_turn.cancel_token ) if llm_err then break end -- ── Aggregate usage ── if response and response.usage then aggregate_usage(response.usage, usage_delta) last_response_usage = response.usage end if has_assistant_payload(response) then break end local retry_index = empty_attempt if retry_index > empty_response_max_retries then llm_err = agent_error.new({ kind = agent_error.kind.protocol, code = "empty_response", message = "assistant returned empty response", origin = "turn", }) break end event_bus:emit(event.event("llm_empty_response_retry", active_turn.id, { retry_index = retry_index, max_retries = empty_response_max_retries, attempt = empty_attempt + 1, max_attempts = empty_response_max_attempts, })) local wait_ms = retry_backoff_ms(empty_response_backoff_ms, retry_index) local _, wait_err = sleep_with_backoff(wait_ms, active_turn.cancel_token) if wait_err then llm_err = wait_err break end empty_attempt = empty_attempt + 1 end -- ── Handle LLM errors ── if llm_err then if agent_error.is_cancelled(llm_err) then stop_reason = "cancelled" else turn_err = agent_error.normalize(llm_err, { origin = "turn" }) event_bus:emit(event.event("turn_failed", active_turn.id, { error = turn_err, })) stop_reason = "error" end break end -- ── Step 18: Extract tool calls and handle response ── local extracted_calls = response and response.message and tool_call_mod.extract_tool_calls(response.message) or nil if extracted_calls and #extracted_calls > 0 then -- Emit llm_tool_call_detected for each call for _, call in ipairs(extracted_calls) do if type(call.id) == "string" and call.id ~= "" then if not detected_tool_call_ids[call.id] then detected_tool_calls[#detected_tool_calls + 1] = call detected_tool_call_ids[call.id] = #detected_tool_calls event_bus:emit(event.event("llm_tool_call_detected", active_turn.id, { call_id = call.id, tool_name = call.name, arguments = call.arguments, })) else detected_tool_calls[detected_tool_call_ids[call.id]] = call end end end local assistant_msg = { role = "assistant", content = response.message.content or "", tool_calls = extracted_calls, } conversation_svc.append_assistant(conversation_state, assistant_msg) assistant_message = assistant_msg -- ── Step 19: Execute tool calls ── local tool_context = { role = role_spec.name, active_role = role_spec.name, turn_id = active_turn.id, conversation_id = ctx.conversation_id, cancel_token = active_turn.cancel_token, event_bus = event_bus, cwd = active_turn.cwd, project = active_turn.project, project_id = active_turn.project_id, runtime = runtime, services = services, service_handles = services, on_event = options.on_event, } local exec_results, new_budget, budget_exhausted, was_cancelled, exec_append_err = execute_tool_calls( extracted_calls, services.tool, tool_context, remaining_budget, event_bus, active_turn.id, conversation_svc, conversation_state ) remaining_budget = new_budget if exec_append_err then turn_err = exec_append_err event_bus:emit(event.event("turn_failed", active_turn.id, { error = turn_err, })) stop_reason = "error" for _, r in ipairs(exec_results) do tool_calls_executed[#tool_calls_executed + 1] = r.result end break end -- Collect executed tool calls for _, r in ipairs(exec_results) do tool_calls_executed[#tool_calls_executed + 1] = r.result end -- Compute child dispatch usage (separate from parent usage_delta) for _, r in ipairs(exec_results) do if type(r) == "table" and type(r.result) == "table" and type(r.result.content) == "table" and type(r.result.content.usage) == "table" then aggregate_usage(r.result.content.usage, dispatch_usage) end end -- ── Emit llm_request_completed after tool execution ── event_bus:emit(event.event("llm_request_completed", active_turn.id, { model = selected_model.model, finish_reason = response and response.finish_reason, })) -- ── Step 20: Check post-tool conditions ── if was_cancelled or cancel.is_cancelled(active_turn.cancel_token) then stop_reason = "cancelled" break end if budget_exhausted then stop_reason = "budget_exhausted" break end -- ── Check for terminal tool execution ── local terminal_result = nil for _, r in ipairs(exec_results) do local spec = services.tool:resolve(r.call.name, role_spec.name) if spec and spec.is_terminal and r.result.ok then terminal_result = r.result break end end if terminal_result then stop_reason = "completed" active_turn.terminal_result = terminal_result break end -- Continue loop: next LLM request with tool results else -- No tool calls: plain assistant answer, append and complete -- ── Emit llm_request_completed for completed no-tool response ── event_bus:emit(event.event("llm_request_completed", active_turn.id, { model = selected_model.model, finish_reason = response and response.finish_reason, })) local content = response and response.message and response.message.content or "" assistant_message = { role = "assistant", content = content } conversation_svc.append_assistant(conversation_state, assistant_message) -- A delegated child answered in prose instead of calling `report`. -- Re-prompt it to submit properly rather than completing the turn -- with no terminal result (which the dispatch worker would otherwise -- have to salvage or block). Bounded by MAX_REPORT_NUDGES. if requires_report and report_nudges < MAX_REPORT_NUDGES then report_nudges = report_nudges + 1 conversation_svc.append_user(conversation_state, REPORT_NUDGE_TEXT) -- Continue loop: re-ask the model, now with the corrective instruction. else stop_reason = "completed" break end end end -- ── Safety: if loop exhausted without stop_reason ── if not stop_reason then stop_reason = "completed" end -- ── Determine cancelled state ── local is_cancelled = cancel.is_cancelled(active_turn.cancel_token) if is_cancelled then stop_reason = "cancelled" end -- ── Build normalized turn result ── local result = { assistant_message = assistant_message, tool_calls_executed = tool_calls_executed, usage_delta = usage_delta, last_usage = last_response_usage, dispatch_usage = dispatch_usage, events = nil, -- set after all terminal events stop_reason = stop_reason, error = nil, cancelled = false, terminal_result = active_turn.terminal_result or nil, } if stop_reason == "cancelled" then result.cancelled = true result.error = agent_error.cancelled("turn", { turn_id = active_turn.id, reason = active_turn.cancel_reason, }) elseif stop_reason == "error" then result.error = turn_err elseif stop_reason == "budget_exhausted" then result.error = agent_error.new({ kind = agent_error.kind.budget_exhausted, code = "budget_exhausted", message = "tool step budget exhausted", origin = "turn", }) end -- ── Emit turn_completed event ── event_bus:emit(event.event("turn_completed", active_turn.id, { stop_reason = stop_reason, })) -- ── Step 24: Persist turn record ── -- The turn's transcript slice, in true per-round order (user message, -- then assistant/tool messages exactly as appended by the loop). local turn_messages = {} for i = turn_messages_start, #conversation_state.messages do turn_messages[#turn_messages + 1] = conversation_state.messages[i] end local persist_record = { turn_id = active_turn.id, conversation_id = ctx.conversation_id, project_id = active_turn.project_id, project_origin = active_turn.project and active_turn.project.normalized or nil, user_message = request.user_message, messages = turn_messages, role = role_spec.name, stop_reason = result.stop_reason, cancelled = result.cancelled, error = result.error, assistant_message = result.assistant_message, tool_calls_executed = result.tool_calls_executed, usage_delta = result.usage_delta, detected_tool_calls = detected_tool_calls, events = event_bus:events(), } local persist_ok, persist_result, persist_err = pcall(function() return services.persistence.persist_turn(persist_record) end) if not persist_ok then -- pcall caught an error: normalize as persistence error local persist_error = agent_error.normalize(persist_result, { kind = agent_error.kind.persistence, origin = "turn", }) event_bus:emit(event.event("turn_persist_failed", active_turn.id, { error = persist_error, })) elseif persist_err then local persist_error = agent_error.normalize(persist_err, { kind = agent_error.kind.persistence, origin = "turn", }) event_bus:emit(event.event("turn_persist_failed", active_turn.id, { error = persist_error, })) elseif persist_result and persist_result.ok then event_bus:emit(event.event("turn_persisted", active_turn.id, { turn_id = active_turn.id, })) else -- persistence returned but without ok=true local persist_error = agent_error.new({ kind = agent_error.kind.persistence, code = "persist_failed", message = "persistence returned non-ok", origin = "turn", }) event_bus:emit(event.event("turn_persist_failed", active_turn.id, { error = persist_error, })) end -- ── Capture final events including turn_completed and turn_persisted ── result.events = event_bus:events() return result end) end return M