16#include <nlohmann/json.hpp>
40 if (hooks.fire_info !=
nullptr) {
41 hooks.fire_info(hooks.registry, point, json);
54 std::string json =
"{\"message_count\":"
55 + std::to_string(ctx.
messages.size())
56 +
",\"delegation_depth\":"
70 std::string json =
"{\"final_state\":\""
72 +
"\",\"iterations\":"
86 std::string json =
"{\"iteration\":"
89 +
"\",\"consecutive_errors\":"
103 std::string json =
"{\"message_count\":"
104 + std::to_string(ctx.
messages.size()) +
"}";
111 const std::string& key);
120 auto now = std::chrono::steady_clock::now();
121 auto dur = now.time_since_epoch();
122 auto ms = std::chrono::duration_cast<
123 std::chrono::milliseconds>(dur);
124 return static_cast<double>(ms.count()) / 1000.0;
136 const InferenceInterface& inference,
139 : inference_(inference),
140 loop_config_(loop_config),
141 token_counter_(loop_config.context_length),
142 compaction_manager_(compaction_config, token_counter_),
144 compaction_manager_, callbacks_,
146 [this](
LoopContext& ctx) { reinject_context_anchors(ctx); }
149 inference, loop_config, callbacks_,
150 GenerationEvents{&interrupt_flag_, &pause_flag_}) {
151 register_directive_handlers();
161 callbacks_ = callbacks;
170void AgentEngine::set_tool_executor(
172 tool_exec_ = tool_exec;
181void AgentEngine::set_tier_resolution(
183 tier_res_ = tier_res;
194 compaction_manager_.set_storage(&storage_);
203void AgentEngine::set_hooks(
const HookInterface& hooks) {
205 response_generator_.set_hooks(hooks);
206 context_manager_.set_hooks(hooks);
207 directive_processor_.set_hooks(hooks);
221void AgentEngine::set_stream_observer(
223 response_generator_.set_stream_observer(observer, user_data);
234void AgentEngine::set_delegation_callbacks(
238 std::lock_guard<std::mutex> lock(delegation_cb_mutex_);
239 delegation_cb_.start = on_start;
240 delegation_cb_.complete = on_complete;
241 delegation_cb_.user_data = user_data;
251AgentEngine::delegation_callbacks_snapshot()
const {
252 std::lock_guard<std::mutex> lock(delegation_cb_mutex_);
253 return delegation_cb_;
263void AgentEngine::set_validation_provider(
264 char* (*provider)(
void*),
void* user_data) {
265 validation_provider_ = provider;
266 validation_provider_data_ = user_data;
292void AgentEngine::apply_identity_overrides(
LoopContext& ctx) {
293 if (ctx.
locked_tier.empty() || tier_res_.get_tier_param ==
nullptr) {
297 ctx.
locked_tier,
"max_iterations", tier_res_.user_data);
300 logger->info(
"[override] tier={} max_iterations={}",
303 auto mt = tier_res_.get_tier_param(
304 ctx.
locked_tier,
"max_tool_calls_per_turn", tier_res_.user_data);
307 logger->info(
"[override] tier={} max_tool_calls_per_turn={}",
320int AgentEngine::resolve_max_iterations(
const LoopContext& ctx)
const {
321 return ctx.effective_max_iterations >= 0
322 ? ctx.effective_max_iterations
323 : loop_config_.max_iterations;
333int AgentEngine::resolve_max_tool_calls(
const LoopContext& ctx)
const {
334 return ctx.effective_max_tool_calls_per_turn >= 0
335 ? ctx.effective_max_tool_calls_per_turn
336 : loop_config_.max_tool_calls_per_turn;
348void AgentEngine::run_loop(
LoopContext& ctx,
bool inherit_interrupt) {
353 if (!inherit_interrupt) {
355 pause_flag_.store(
false);
357 apply_identity_overrides(ctx);
358 reinject_context_anchors(ctx);
362 set_state(ctx, AgentState::PLANNING);
366 auto& tm = per_tier_metrics_[
386std::vector<Message> AgentEngine::run(std::vector<Message> messages,
387 const std::string& tier_override) {
391 last_activity_epoch_s_.store(
392 std::chrono::duration_cast<std::chrono::seconds>(
393 std::chrono::system_clock::now().time_since_epoch()).
count());
400 init_session_conversation(ctx);
403 pause_flag_.store(
false);
405 reinject_context_anchors(ctx);
406 set_state(ctx, AgentState::PLANNING);
411 accumulate_run_metrics(ctx);
412 logger->info(
"Loop complete: {} iterations, {}ms",
423void AgentEngine::init_session_conversation(
LoopContext& ctx) {
431 if (storage_.create_conversation ==
nullptr) {
return; }
433 if (storage_.create_conversation(
"session", conv_id, storage_.user_data)
434 && !conv_id.empty()) {
437 logger->warn(
"Storage create_conversation failed at run() init; "
438 "delegations will not persist this session (gh#48)");
448void AgentEngine::accumulate_run_metrics(LoopContext& ctx) {
449 last_metrics_ = ctx.metrics;
451 auto& tm = per_tier_metrics_[
452 ctx.locked_tier.empty() ?
"lead" : ctx.locked_tier];
453 tm.iterations += ctx.metrics.iterations;
454 tm.tool_calls += ctx.metrics.tool_calls;
455 tm.tokens_used += ctx.metrics.tokens_used;
456 tm.errors += ctx.metrics.errors;
457 tm.end_time += (ctx.metrics.end_time - ctx.metrics.start_time);
467void AgentEngine::loop(LoopContext& ctx) {
470 while (!should_stop(ctx)) {
471 ctx.metrics.iterations++;
473 if (interrupt_flag_.load()) {
474 set_state(ctx, AgentState::INTERRUPTED);
478 execute_iteration(ctx);
480 if (ctx.state == AgentState::ERROR) {
485 if (ctx.metrics.iterations >= resolve_max_iterations(ctx)
486 && !is_terminal_state(ctx)) {
492 logger->warn(
"Loop ended due to max iterations ({}/{}) — "
493 "forcing synthetic entropic.complete",
494 ctx.metrics.iterations, resolve_max_iterations(ctx));
496 forced.role =
"assistant";
497 forced.content =
"[iteration cap reached after "
498 + std::to_string(ctx.metrics.iterations)
499 +
" iterations — returning current state]";
500 ctx.messages.push_back(std::move(forced));
501 ctx.metadata[
"terminal_reason"] =
"budget_exhausted";
502 set_state(ctx, AgentState::COMPLETE);
515void AgentEngine::execute_iteration(LoopContext& ctx) {
516 logger->info(
"[LOOP] iter {}/{} state={} msgs={}",
517 ctx.metrics.iterations,
518 resolve_max_iterations(ctx),
520 ctx.messages.size());
524 context_manager_.refresh_context_limit(ctx, 0);
525 context_manager_.prune_old_tool_results(ctx);
526 context_manager_.check_compaction(ctx);
530 set_state(ctx, AgentState::EXECUTING);
534 set_state(ctx, AgentState::COMPLETE);
538 auto result = response_generator_.generate_response(ctx);
539 dispatch_post_generate(ctx, result);
540 process_generation_result(ctx, result);
550void AgentEngine::process_generation_result(LoopContext& ctx,
551 GenerateResult& result) {
552 auto [cleaned, tool_calls] = parse_tool_calls(result.content);
553 logger->info(
"[ITER] finish={}, {} tool call(s), {} chars",
554 result.finish_reason, tool_calls.size(),
556 Message assistant_msg{
"assistant", cleaned};
557 ctx.messages.push_back(std::move(assistant_msg));
559 bool made_tool_call =
560 !tool_calls.empty() && tool_exec_.process_tool_calls !=
nullptr;
564 bool made_progress =
false;
565 if (made_tool_call) {
566 made_progress = process_tool_results(ctx, tool_calls);
570 ctx.metadata[
"zero_tool_call_retries"] =
"0";
573 evaluate_no_tool_decision(ctx, cleaned, result.finish_reason);
578 charge_thinking_budget(ctx, result.content.size(), made_progress);
580 dispatch_pending_or_halt(ctx);
588void AgentEngine::charge_thinking_budget(
589 LoopContext& ctx,
size_t content_len,
bool made_tool_call) {
590 if (loop_config_.budget_mode == BudgetMode::off) {
return; }
594 if (made_tool_call) {
595 ctx.budget_tokens_since_tool = 0;
596 ctx.budget_window_start_s = 0.0;
597 ctx.budget_completion_nudged =
false;
601 int units = budget_units_consumed(ctx, content_len);
602 if (units < loop_config_.budget_limit) {
return; }
604 if (!ctx.budget_completion_nudged) {
605 nudge_budget_completion(ctx);
607 hard_cut_budget(ctx);
618int AgentEngine::budget_units_consumed(LoopContext& ctx,
size_t content_len) {
619 if (loop_config_.budget_mode == BudgetMode::tokens) {
624 constexpr size_t kCharsPerToken = 4;
625 ctx.budget_tokens_since_tool +=
626 static_cast<int>(content_len / kCharsPerToken);
627 return ctx.budget_tokens_since_tool;
630 if (ctx.budget_window_start_s == 0.0) {
633 return static_cast<int>(
now_seconds() - ctx.budget_window_start_s);
641void AgentEngine::nudge_budget_completion(LoopContext& ctx) {
642 ctx.budget_completion_nudged =
true;
643 ctx.budget_tokens_since_tool = 0;
644 ctx.budget_window_start_s =
645 (loop_config_.budget_mode == BudgetMode::wall_clock)
647 logger->warn(
"[BUDGET] thinking budget reached ({} {}); nudging completion",
648 loop_config_.budget_limit,
649 loop_config_.budget_mode == BudgetMode::tokens
650 ?
"tokens" :
"seconds");
651 Message nudge{
"user",
652 "[engine] thinking budget reached — emit entropic.complete now "
653 "with your best current answer. Further deliberation without a "
654 "tool call will be cut off."};
655 ctx.messages.push_back(std::move(nudge));
663void AgentEngine::hard_cut_budget(LoopContext& ctx) {
664 logger->warn(
"[BUDGET] thinking budget exhausted after nudge — "
665 "hard-cutting turn");
666 Message cut{
"assistant",
667 "[thinking budget exhausted — the turn was hard-cut after the "
668 "completion nudge went unheeded; no tool call was emitted]"};
669 ctx.messages.push_back(std::move(cut));
670 ctx.metadata[
"terminal_reason"] =
"budget_exhausted_thinking";
671 set_state(ctx, AgentState::COMPLETE);
692void AgentEngine::dispatch_pending_or_halt(LoopContext& ctx) {
693 if (interrupt_flag_.load() && !is_terminal_state(ctx)) {
694 logger->info(
"[ITER] interrupt observed during tool processing — "
695 "halting before pending dispatch");
696 set_state(ctx, AgentState::INTERRUPTED);
700 if (is_terminal_state(ctx)) {
701 if (ctx.pending_delegation.has_value()
702 || ctx.pending_pipeline.has_value()) {
703 logger->info(
"[ITER] terminal state reached during batch "
704 "tool processing — dropping pending {}",
705 ctx.pending_delegation.has_value()
706 ?
"delegation" :
"pipeline");
711 if (ctx.pending_delegation.has_value()) {
712 ctx.metrics.iterations--;
713 execute_pending_delegation(ctx);
714 }
else if (ctx.pending_pipeline.has_value()) {
715 ctx.metrics.iterations--;
716 execute_pending_pipeline(ctx);
728void AgentEngine::evaluate_no_tool_decision(
730 const std::string& content,
731 const std::string& finish_reason) {
732 if (handle_terminal_finish_reasons(ctx, finish_reason)) {
return; }
733 if (try_auto_chain(ctx, finish_reason, content)) {
734 logger->info(
"[DECISION] auto-chain triggered");
737 if (record_explicit_completion_failure(ctx, finish_reason)) {
return; }
738 if (response_generator_.is_response_complete(content,
"[]")) {
739 logger->info(
"[DECISION] response complete");
740 set_state(ctx, AgentState::COMPLETE);
757bool AgentEngine::handle_terminal_finish_reasons(
758 LoopContext& ctx,
const std::string& finish_reason) {
759 if (finish_reason ==
"interrupted") {
760 logger->info(
"[DECISION] interrupted");
761 set_state(ctx, AgentState::INTERRUPTED);
764 if (finish_reason ==
"length") {
765 logger->info(
"[DECISION] length, continuing");
781bool AgentEngine::record_explicit_completion_failure(
782 LoopContext& ctx,
const std::string& finish_reason) {
783 if (!tier_requires_explicit_completion(ctx.locked_tier)) {
786 auto& retries = ctx.metadata[
"zero_tool_call_retries"];
787 int n = retries.empty() ? 0 : std::atoi(retries.c_str());
790 if (tier_res_.get_tier_param !=
nullptr) {
791 auto cv = tier_res_.get_tier_param(
792 ctx.locked_tier,
"max_consecutive_empty_turns",
793 tier_res_.user_data);
794 if (!cv.empty()) { ceiling = std::atoi(cv.c_str()); }
798 ctx.metadata[
"failure_reason"] =
799 "zero_tool_calls_with_explicit_completion";
800 ctx.metadata[
"failure_tier"] = ctx.locked_tier;
801 logger->error(
"[DECISION] empty-turn allowance exhausted "
802 "(tier={}, n={}/{}), failing turn",
803 ctx.locked_tier, n, ceiling);
804 set_state(ctx, AgentState::ERROR);
808 logger->info(
"[DECISION] tier '{}' empty turn {}/{} "
809 "(finish_reason={}), nudging toward action",
810 ctx.locked_tier, n + 1, ceiling, finish_reason);
811 std::string tools_hint =
812 "(entropic.complete, entropic.delegate, or entropic.inspect)";
813 if (tier_res_.get_tier_param !=
nullptr) {
814 auto allowed = tier_res_.get_tier_param(
815 ctx.locked_tier,
"allowed_tools", tier_res_.user_data);
816 if (!allowed.empty()) { tools_hint =
"(" + allowed +
")"; }
819 correction.role =
"user";
821 "[SYSTEM] Your previous response contained no tool call. "
822 "You must end every turn with exactly one tool call "
823 + tools_hint +
". Retry.";
824 ctx.messages.push_back(std::move(correction));
825 retries = std::to_string(n + 1);
826 set_state(ctx, AgentState::EXECUTING);
842bool AgentEngine::tier_requires_explicit_completion(
843 const std::string& tier)
const {
844 if (tier_res_.get_tier_param ==
nullptr) {
return false; }
845 auto val = tier_res_.get_tier_param(
846 tier,
"explicit_completion", tier_res_.user_data);
847 return val ==
"true" || val ==
"1";
863bool AgentEngine::is_delegation_cycle(
864 const LoopContext& ctx,
const std::string& target)
const {
867 if (anc == target) {
return true; }
879bool AgentEngine::is_delegation_repeat_blocked(
880 const LoopContext& ctx,
const std::string& target)
const {
883 >= loop_config_.max_consecutive_failed_delegations;
894bool AgentEngine::fold_complete_into_assistant(
896 auto tn_it = tool_result_msg.
metadata.find(
"tool_name");
897 bool is_complete = tn_it != tool_result_msg.
metadata.end()
898 && tn_it->second ==
"entropic.complete";
899 if (!is_complete || ctx.
messages.empty()) {
return false; }
901 auto sum_it = ctx.
metadata.find(
"explicit_completion_summary");
902 bool foldable = last.role ==
"assistant" && last.content.empty()
904 if (!foldable) {
return false; }
905 last.content = sum_it->second;
914int64_t AgentEngine::seconds_since_last_activity()
const {
915 auto last = last_activity_epoch_s_.load();
916 if (last == 0) {
return 0; }
917 auto now = std::chrono::duration_cast<std::chrono::seconds>(
918 std::chrono::system_clock::now().time_since_epoch()).count();
919 return now > last ? (now - last) : 0;
929bool AgentEngine::should_stop(
const LoopContext& ctx)
const {
932 return is_terminal_state(ctx) || at_limit;
949bool AgentEngine::is_terminal_state(
const LoopContext& ctx) {
950 return ctx.state == AgentState::COMPLETE
951 || ctx.state == AgentState::ERROR
952 || ctx.state == AgentState::INTERRUPTED;
963void AgentEngine::set_state(LoopContext& ctx, AgentState state) {
964 auto prev = ctx.state;
967 if (callbacks_.on_state_change !=
nullptr) {
968 callbacks_.on_state_change(
969 static_cast<int>(state), callbacks_.user_data);
975 if (state_observer_ !=
nullptr) {
976 state_observer_(
static_cast<int>(state), state_observer_data_);
981 std::string json =
"{\"previous\":\""
983 +
"\",\"current\":\""
989 if (state == AgentState::ERROR) {
990 std::string json =
"{\"error_code\":\"STATE_ERROR\""
992 + std::to_string(ctx.metrics.iterations)
993 +
",\"consecutive\":" + std::to_string(ctx.consecutive_errors)
1010void AgentEngine::interrupt() {
1011 if (!interrupt_flag_.exchange(
true)) {
1012 logger->info(
"Engine interrupted");
1015 if (external_interrupt_cb_ !=
nullptr) {
1016 external_interrupt_cb_(external_interrupt_data_);
1019 pause_flag_.store(
false);
1034void AgentEngine::set_external_interrupt(
void (*cb)(
void*),
1036 external_interrupt_cb_ = cb;
1037 external_interrupt_data_ = user_data;
1045void AgentEngine::reset_interrupt() {
1046 interrupt_flag_.store(
false);
1054void AgentEngine::pause() {
1055 logger->info(
"Engine paused");
1056 pause_flag_.store(
true);
1064void AgentEngine::cancel_pause() {
1065 logger->info(
"Pause cancelled, interrupting");
1066 pause_flag_.store(
false);
1067 interrupt_flag_.store(
true);
1077std::pair<int, int> AgentEngine::context_usage(
1078 const std::vector<Message>& messages)
const {
1079 int used = token_counter_.count_messages(messages);
1080 return {used, token_counter_.max_tokens};
1089void AgentEngine::reinject_context_anchors(
LoopContext& ctx) {
1090 for (
const auto& [key, content] : context_anchors_) {
1093 dir_anchor(ctx, d, r);
1104void AgentEngine::register_directive_handlers() {
1106 directive_processor_.register_handler(t,
1107 [
this, fn](LoopContext& c,
const Directive& d,
1108 DirectiveResult& r) {
1109 (this->*fn)(c, d, r);
1132void AgentEngine::dir_stop(
1133 LoopContext&,
const Directive&, DirectiveResult& r) {
1134 logger->info(
"[DIRECTIVE] stop_processing");
1135 r.stop_processing =
true;
1143void AgentEngine::dir_tier_change(
1144 LoopContext& ctx,
const Directive& d, DirectiveResult& r) {
1145 const auto& tc =
static_cast<const TierChangeDirective&
>(d);
1146 ctx.locked_tier = tc.tier;
1147 r.tier_changed =
true;
1148 logger->info(
"[DIRECTIVE] tier_change: {}", tc.tier);
1156void AgentEngine::dir_delegate(
1157 LoopContext& ctx,
const Directive& d, DirectiveResult& r) {
1158 const auto& dl =
static_cast<const DelegateDirective&
>(d);
1159 ctx.pending_delegation = PendingDelegation{
1160 dl.target, dl.task, dl.max_turns,
1161 dl.resume_from_delegation_id};
1162 r.stop_processing =
true;
1163 if (dl.resume_from_delegation_id.empty()) {
1164 logger->info(
"[DIRECTIVE] delegate: target={} task='{}'",
1165 dl.target, dl.task);
1168 logger->info(
"[DIRECTIVE] resume_delegation: id={} task='{}'",
1169 dl.resume_from_delegation_id, dl.task);
1178void AgentEngine::dir_pipeline(
1179 LoopContext& ctx,
const Directive& d, DirectiveResult& r) {
1180 const auto& pl =
static_cast<const PipelineDirective&
>(d);
1181 ctx.pending_pipeline = PendingPipeline{pl.stages, pl.task};
1182 r.stop_processing =
true;
1183 logger->info(
"[DIRECTIVE] pipeline: {} stages", pl.stages.size());
1197void AgentEngine::dir_complete(
1198 LoopContext& ctx,
const Directive& d, DirectiveResult& r) {
1199 const auto& cd =
static_cast<const CompleteDirective&
>(d);
1203 auto prior_state = ctx.state;
1204 set_state(ctx, AgentState::VERIFYING);
1207 if (fire_complete_hook(cd.summary, ctx)) {
1208 logger->info(
"[DIRECTIVE] complete REJECTED by hook → revising");
1210 ctx.metadata[
"validator_phase"] =
"revising";
1211 set_state(ctx, prior_state);
1212 r.stop_processing =
false;
1216 ctx.metadata[
"explicit_completion_summary"] = cd.summary;
1223 if (cd.coverage_gap) {
1224 ctx.metadata[
"coverage_gap"] =
"true";
1225 ctx.metadata[
"gap_description"] = cd.gap_description;
1226 nlohmann::json files_arr = cd.suggested_files;
1227 ctx.metadata[
"suggested_files_json"] = files_arr.dump();
1229 set_state(ctx, AgentState::COMPLETE);
1230 r.stop_processing =
true;
1231 logger->info(
"[DIRECTIVE] complete coverage_gap={}",
1240void AgentEngine::dir_clear_todos(
1241 LoopContext&,
const Directive&, DirectiveResult&) {
1242 logger->debug(
"[DIRECTIVE] clear_self_todos (no-op)");
1250void AgentEngine::dir_inject(
1251 LoopContext&,
const Directive& d, DirectiveResult& r) {
1252 const auto& ic =
static_cast<const InjectContextDirective&
>(d);
1253 if (!ic.content.empty()) {
1256 msg.content = ic.content;
1257 r.injected_messages.push_back(std::move(msg));
1258 logger->info(
"[DIRECTIVE] inject_context");
1267void AgentEngine::dir_prune(
1268 LoopContext& ctx,
const Directive& d, DirectiveResult&) {
1269 const auto& pm =
static_cast<const PruneMessagesDirective&
>(d);
1270 auto [pruned, freed] = context_manager_.prune_tool_results(
1271 ctx, pm.keep_recent);
1272 logger->info(
"[DIRECTIVE] prune: {} results, {} chars", pruned, freed);
1280void AgentEngine::dir_anchor(
1281 LoopContext& ctx,
const Directive& d, DirectiveResult&) {
1282 const auto& ca =
static_cast<const ContextAnchorDirective&
>(d);
1283 if (ca.content.empty()) {
1284 context_anchors_.erase(ca.key);
1286 logger->info(
"Removed anchor: {}", ca.key);
1289 context_anchors_[ca.key] = ca.content;
1292 anchor.role =
"user";
1293 anchor.content = ca.content;
1294 anchor.metadata[
"is_context_anchor"] =
"true";
1295 anchor.metadata[
"anchor_key"] = ca.key;
1296 ctx.messages.push_back(std::move(anchor));
1297 logger->info(
"Updated anchor: {}", ca.key);
1305void AgentEngine::dir_phase(
1306 LoopContext& ctx,
const Directive& d, DirectiveResult&) {
1307 const auto& pc =
static_cast<const PhaseChangeDirective&
>(d);
1308 ctx.active_phase = pc.phase;
1309 logger->info(
"[DIRECTIVE] phase_change: {}", pc.phase);
1317void AgentEngine::dir_notify(
1318 LoopContext&,
const Directive& d, DirectiveResult&) {
1319 const auto& np =
static_cast<const NotifyPresenterDirective&
>(d);
1320 if (callbacks_.on_presenter_notify !=
nullptr) {
1321 callbacks_.on_presenter_notify(
1322 np.key.c_str(), np.data_json.c_str(),
1323 callbacks_.user_data);
1345 const nlohmann::json& obj,
const std::string& tc_str) {
1347 tc.
name = obj.value(
"name",
"");
1348 tc.
id =
"tc-" + std::to_string(
1349 std::hash<std::string>{}(tc.
name + tc_str) & 0xFFFF);
1350 if (obj.contains(
"arguments") && obj[
"arguments"].is_object()) {
1352 for (
auto& [k, v] : obj[
"arguments"].items()) {
1354 ? v.get<std::string>() : v.dump();
1368 const std::string& tc_str) {
1369 std::vector<ToolCall> calls;
1370 if (tc_str ==
"[]" || tc_str.empty()) {
return calls; }
1371 auto arr = nlohmann::json::parse(tc_str,
nullptr,
false);
1372 if (!arr.is_array()) {
return calls; }
1373 for (
const auto& obj : arr) {
1392std::pair<std::string, std::vector<ToolCall>>
1393AgentEngine::parse_tool_calls(
const std::string& raw_content) {
1394 if (inference_.parse_tool_calls ==
nullptr) {
1395 return {raw_content, {}};
1398 char* cleaned =
nullptr;
1399 char* tc_json =
nullptr;
1400 int rc = inference_.parse_tool_calls(
1401 raw_content.c_str(), &cleaned, &tc_json,
1402 inference_.adapter_data);
1422 std::string cleaned_str =
1423 mcp::sanitize_utf8(cleaned ? cleaned : raw_content);
1424 std::string tc_str = mcp::sanitize_utf8(tc_json ? tc_json :
"[]");
1426 if (inference_.free_fn !=
nullptr) {
1427 if (cleaned !=
nullptr) { inference_.free_fn(cleaned); }
1428 if (tc_json !=
nullptr) { inference_.free_fn(tc_json); }
1431 if (rc != 0) {
return {cleaned_str, {}}; }
1433 if (!calls.empty()) {
1434 logger->info(
"Parsed tool calls from model output");
1436 return {cleaned_str, std::move(calls)};
1444void AgentEngine::defang_meta_action_envelope(
Message& msg)
const {
1445 auto tn = msg.
metadata.find(
"tool_name");
1446 bool is_meta = tn != msg.
metadata.end()
1447 && tn->second.rfind(
"entropic.", 0) == 0;
1448 if (!is_meta) {
return; }
1449 auto j = nlohmann::json::parse(msg.
content,
nullptr,
false);
1450 bool is_envelope = j.is_object() && j.contains(
"action")
1451 && j[
"action"].is_string();
1452 if (!is_envelope) {
return; }
1453 msg.
content =
"(engine: " + j[
"action"].get<std::string>()
1454 +
" directive accepted)";
1466bool AgentEngine::process_tool_results(
1468 const std::vector<ToolCall>& tool_calls) {
1469 set_state(ctx, AgentState::WAITING_TOOL);
1471 auto results = tool_exec_.process_tool_calls(
1472 ctx, tool_calls, tool_exec_.user_data);
1474 logger->info(
"[TOOLS] {} call(s) -> {} result message(s)",
1475 tool_calls.size(), results.size());
1482 bool made_progress =
false;
1483 for (
auto& msg : results) {
1484 auto it = msg.metadata.find(
"result_kind");
1485 if (it != msg.metadata.end()
1486 && (it->second ==
"ok" || it->second ==
"ok_empty")) {
1487 made_progress =
true;
1489 if (fold_complete_into_assistant(ctx, msg)) {
1492 defang_meta_action_envelope(msg);
1493 ctx.
messages.push_back(std::move(msg));
1498 if (!is_terminal_state(ctx)) {
1499 set_state(ctx, AgentState::EXECUTING);
1501 return made_progress;
1514 const std::string& key) {
1517 std::remove_if(msgs.begin(), msgs.end(),
1519 auto it = m.metadata.find(
"anchor_key");
1520 return it != m.metadata.end() && it->second == key;
1536bool AgentEngine::fire_pre_hook(
1538 if (hooks_.fire_pre ==
nullptr) {
1541 std::string json =
"{\"iteration\":"
1542 + std::to_string(iteration) +
"}";
1543 char* modified =
nullptr;
1544 int rc = hooks_.fire_pre(hooks_.registry,
1545 point, json.c_str(), &modified);
1561 static const auto tbl = []() {
1562 std::array<const char*, 256> t{};
1563 t[
static_cast<unsigned char>(
'"')] =
"\\\"";
1564 t[
static_cast<unsigned char>(
'\\')] =
"\\\\";
1565 t[
static_cast<unsigned char>(
'\n')] =
"\\n";
1566 t[
static_cast<unsigned char>(
'\r')] =
"\\r";
1567 t[
static_cast<unsigned char>(
'\t')] =
"\\t";
1583 out.reserve(s.size() + 16);
1585 const char* esc = tbl[
static_cast<unsigned char>(c)];
1586 if (esc) { out += esc; }
1604 const std::vector<Message>& messages) {
1605 std::string manifest;
1606 for (
const auto& msg : messages) {
1607 auto it = msg.metadata.find(
"tool_name");
1608 if (it == msg.metadata.end() || it->second.empty()) {
1616 auto orig = msg.metadata.find(
"original_content");
1617 size_t real_size = (orig != msg.metadata.end())
1618 ? orig->second.size()
1619 : msg.content.size();
1620 manifest +=
"- " + it->second
1621 +
" \xe2\x86\x92 " + std::to_string(real_size)
1636 size_t chars_per_result) {
1637 auto orig = msg.
metadata.find(
"original_content");
1638 const std::string& src = (orig != msg.
metadata.end())
1641 auto iter = msg.
metadata.find(
"added_at_iteration");
1642 std::string out =
"## ";
1643 out += msg.
metadata.at(
"tool_name");
1645 out +=
" [iter " + iter->second +
"]";
1648 if (src.size() <= chars_per_result) {
1651 out += src.substr(0, chars_per_result);
1652 out +=
"\n[... truncated, "
1653 + std::to_string(src.size() - chars_per_result)
1688 const std::vector<Message>& messages,
1690 size_t chars_per_result) {
1691 std::vector<const Message*> tool_msgs;
1692 for (
const auto& msg : messages) {
1693 auto it = msg.metadata.find(
"tool_name");
1694 if (it == msg.metadata.end() || it->second.empty()) {
1697 tool_msgs.push_back(&msg);
1699 if (tool_msgs.empty()) {
return {}; }
1701 size_t elided = (tool_msgs.size() > max_results)
1702 ? (tool_msgs.size() - max_results) : 0;
1703 size_t start = elided;
1706 out +=
"(" + std::to_string(elided)
1707 +
" earlier tool results elided for length)\n";
1709 for (
size_t i = start; i < tool_msgs.size(); ++i) {
1730 const std::string& base,
1732 if (tool_exec.
history_json ==
nullptr) {
return base; }
1734 if (js ==
nullptr) {
return base; }
1735 std::string out = base +
"prior-iteration history: " + js +
"\n";
1748 const std::vector<Message>& messages) {
1749 for (
const auto& msg : messages) {
1750 if (msg.role ==
"system") {
return msg.content; }
1770void AgentEngine::fire_post_generate_hook(
1771 GenerateResult& result,
1772 const std::string& tier,
1773 const std::vector<Message>& messages) {
1774 if (hooks_.fire_post ==
nullptr) {
1789 "{\"finish_reason\":\"" + result.finish_reason
1791 +
"\",\"tier\":\"" + tier
1795 char* out =
nullptr;
1796 hooks_.fire_post(hooks_.registry,
1798 if (out !=
nullptr) {
1804 result.content = mcp::sanitize_utf8(out);
1805 logger->info(
"POST_GENERATE hook revised content");
1831void AgentEngine::dispatch_post_generate(
1832 LoopContext& ctx, GenerateResult& result) {
1837 ctx.pending_validation_feedback.clear();
1838 ctx.pending_anti_spiral_warning.clear();
1839 fire_post_generate_hook(result, ctx.locked_tier, ctx.messages);
1842 capture_validation_feedback(ctx);
1850void AgentEngine::capture_validation_feedback(LoopContext& ctx) {
1851 if (validation_provider_ ==
nullptr) {
return; }
1852 char* v = validation_provider_(validation_provider_data_);
1853 if (v ==
nullptr) {
return; }
1857 auto j = nlohmann::json::parse(raw);
1858 auto verdict = j.value(
"verdict",
"");
1859 if (verdict.rfind(
"rejected", 0) != 0) {
return; }
1861 if (j.contains(
"violations") && j[
"violations"].is_array()) {
1862 for (
const auto& vio : j[
"violations"]) {
1863 if (!joined.empty()) { joined +=
"; "; }
1864 joined += vio.is_string()
1865 ? vio.get<std::string>()
1869 if (joined.empty()) { joined = verdict; }
1870 ctx.pending_validation_feedback = joined;
1871 }
catch (
const nlohmann::json::exception&) {
1888 const std::vector<Message>& messages) {
1889 std::string arr =
"[";
1891 for (
const auto& msg : messages) {
1892 auto it = msg.metadata.find(
"tool_name");
1893 if (it == msg.metadata.end() || it->second.empty()) {
1896 if (!first) { arr +=
","; }
1926bool AgentEngine::fire_complete_hook(
1927 const std::string& summary,
1928 const LoopContext& ctx) {
1929 if (hooks_.fire_pre ==
nullptr) {
return false; }
1935 std::string validation_block =
",\"validation\":null";
1936 if (validation_provider_ !=
nullptr) {
1937 char* v = validation_provider_(validation_provider_data_);
1939 validation_block =
",\"validation\":" + std::string(v);
1945 +
"\",\"tier\":\"" + ctx.locked_tier
1946 +
"\",\"tool_results\":" + tool_results
1947 +
",\"iteration\":" + std::to_string(ctx.metrics.iterations)
1950 char* modified =
nullptr;
1951 int rc = hooks_.fire_pre(hooks_.registry,
1953 if (modified !=
nullptr) {
1959 feedback.role =
"user";
1961 std::string(
"[CITATION VALIDATION] ") + mcp::sanitize_utf8(modified);
1962 const_cast<LoopContext&
>(ctx).messages.push_back(
1963 std::move(feedback));
1978bool AgentEngine::fire_delegate_pre_hook(
1979 const PendingDelegation& pending,
int depth) {
1980 if (hooks_.fire_pre ==
nullptr) {
1983 std::string json =
"{\"target_tier\":\""
1984 + pending.target +
"\",\"task\":\""
1985 + pending.task +
"\",\"depth\":"
1986 + std::to_string(depth) +
"}";
1987 char* modified =
nullptr;
1988 int rc = hooks_.fire_pre(hooks_.registry,
2013 const std::string& target,
2015 const std::string& summary) {
2017 j[
"target_tier"] = target;
2018 j[
"success"] = success;
2020 ? ToolResultKind::ok
2022 j[
"summary"] = mcp::sanitize_utf8(summary);
2041void AgentEngine::fire_delegate_complete_hook(
2042 const std::string& target,
bool success,
2043 const std::string& summary) {
2044 if (hooks_.fire_post ==
nullptr) {
2047 std::string json = detail::build_delegate_complete_json(
2048 target, success, summary);
2049 char* out =
nullptr;
2050 hooks_.fire_post(hooks_.registry,
2065 auto* engine =
static_cast<AgentEngine*
>(user_data);
2079 reject.
role =
"user";
2080 reject.
content =
"[DELEGATION REJECTED] Maximum delegation "
2081 "depth (" + std::to_string(
2082 AgentEngine::MAX_DELEGATION_DEPTH) +
2084 ctx.
messages.push_back(std::move(reject));
2101 chain += anc +
" -> ";
2106 reject.
role =
"user";
2107 reject.
content =
"[DELEGATION REJECTED] Circular delegation "
2108 "detected: tier '" + target +
2109 "' already in the active delegation chain "
2110 "(" + chain +
"). Choose a different target.";
2111 ctx.
metadata[
"failure_reason"] =
"delegation_cycle";
2112 ctx.
metadata[
"failure_target"] = target;
2113 ctx.
messages.push_back(std::move(reject));
2126 std::string tag = result.
success ?
"COMPLETE" :
"FAILED";
2129 msg.
content =
"[DELEGATION " + tag +
": " + target +
"] " + result.
summary;
2130 ctx.
messages.push_back(std::move(msg));
2145 LoopContext& ctx,
const std::string& target,
int n) {
2147 reject.
role =
"user";
2148 reject.
content =
"[DELEGATION REJECTED] '" + target
2149 +
"' has just failed " + std::to_string(n)
2150 +
" times in a row. Stop retrying this target. Either "
2151 "respond to the user with what you have, or delegate to "
2152 "a different tier.";
2153 ctx.
metadata[
"failure_reason"] =
"delegation_repeat_blocked";
2154 ctx.
metadata[
"failure_target"] = target;
2155 ctx.
messages.push_back(std::move(reject));
2166bool AgentEngine::reject_delegation_if_guarded(
2167 LoopContext& ctx,
const PendingDelegation& pending) {
2168 bool rejected =
true;
2169 if (ctx.delegation_depth >= MAX_DELEGATION_DEPTH) {
2170 logger->warn(
"Delegation rejected: depth {} >= max {}",
2171 ctx.delegation_depth, MAX_DELEGATION_DEPTH);
2172 push_delegation_rejected(ctx);
2173 }
else if (is_delegation_cycle(ctx, pending.target)) {
2175 logger->warn(
"Delegation rejected: cycle on target tier '{}'",
2178 }
else if (is_delegation_repeat_blocked(ctx, pending.target)) {
2183 logger->warn(
"Delegation rejected: '{}' failed {}x in a row "
2184 "(>= max_consecutive_failed_delegations={})",
2186 ctx.consecutive_failed_delegations,
2187 loop_config_.max_consecutive_failed_delegations);
2189 ctx, pending.target, ctx.consecutive_failed_delegations);
2209void AgentEngine::execute_pending_delegation(LoopContext& ctx) {
2210 auto pending = std::move(*ctx.pending_delegation);
2211 ctx.pending_delegation.reset();
2213 if (reject_delegation_if_guarded(ctx, pending)) {
return; }
2215 set_state(ctx, AgentState::DELEGATING);
2224 std::vector<Message> resume_history;
2225 bool blocked =
false;
2226 if (!pending.resume_from_delegation_id.empty()
2227 && !resolve_resume_delegation(ctx, pending, resume_history)) {
2229 }
else if (fire_delegate_pre_hook(pending, ctx.delegation_depth)) {
2230 logger->info(
"ON_DELEGATE hook cancelled delegation");
2234 set_state(ctx, AgentState::EXECUTING);
2238 fire_delegation_start(ctx, pending.target, pending.task);
2239 auto result = run_pending_delegation(
2240 ctx, pending, std::move(resume_history));
2244 if (result.success) {
2245 ctx.last_failed_delegation_target.clear();
2246 ctx.consecutive_failed_delegations = 0;
2247 }
else if (pending.target == ctx.last_failed_delegation_target) {
2248 ++ctx.consecutive_failed_delegations;
2250 ctx.last_failed_delegation_target = pending.target;
2251 ctx.consecutive_failed_delegations = 1;
2254 fire_delegation_complete(ctx, pending.target, result);
2255 fire_delegate_complete_hook(pending.target, result.success,
2258 finalize_delegation_result(ctx, result);
2273void AgentEngine::relay_partial_result(
2274 LoopContext& ctx,
const std::string& summary) {
2275 GenerateResult relay_result;
2276 relay_result.content = summary;
2277 relay_result.finish_reason =
"stop";
2278 relay_result.tool_calls_json =
"[]";
2279 fire_post_generate_hook(
2280 relay_result, ctx.locked_tier, ctx.messages);
2281 ctx.metadata[
"explicit_completion_summary"] = relay_result.content;
2282 set_state(ctx, AgentState::COMPLETE);
2295 "[COVERAGE GAP from " + tier +
"]\n"
2296 "Summary so far: " + result.
summary +
"\n"
2299 body +=
"\nSuggested files to inspect:";
2301 body +=
"\n - " + f;
2324 for (
auto rit = ctx.
messages.rbegin();
2325 rit != ctx.
messages.rend(); ++rit) {
2326 if (rit->role ==
"assistant" && rit->content.empty()) {
2327 rit->content = summary;
2353void AgentEngine::finalize_delegation_result(
2354 LoopContext& ctx,
const DelegationResult& result) {
2357 const bool in_relay_tier =
2358 relay_single_delegate_tiers_.count(ctx.locked_tier) > 0;
2359 if (in_relay_tier && result.coverage_gap) {
2362 gap.content = build_coverage_gap_message(
2363 result.target_tier, result);
2364 ctx.messages.push_back(std::move(gap));
2365 ctx.metadata[
"relay_status"] =
"coverage_gap_suppressed";
2367 "[COVERAGE GAP] suppressing auto-relay for tier={} "
2368 "({} suggested files)",
2369 result.target_tier, result.suggested_files.size());
2370 set_state(ctx, AgentState::EXECUTING);
2373 if (result.success && in_relay_tier) {
2374 relay_partial_result(ctx, result.summary);
2375 log_relay_status(ctx);
2378 if (!result.terminal_reason.empty() && in_relay_tier) {
2379 relay_partial_result(ctx,
2380 "[partial — budget_exhausted] " + result.summary);
2381 log_relay_status(ctx, result.terminal_reason);
2384 bool needs_explicit = tier_requires_explicit_completion(
2386 if (result.success && !needs_explicit) {
2389 set_state(ctx, (result.success && !needs_explicit)
2390 ? AgentState::COMPLETE :
AgentState::EXECUTING);
2407void AgentEngine::log_relay_status(LoopContext& ctx,
2408 const std::string& terminal_reason) {
2409 std::string verdict;
2410 if (validation_provider_ !=
nullptr) {
2411 char* v = validation_provider_(validation_provider_data_);
2415 auto pos = jv.find(
"\"verdict\":\"");
2416 if (pos != std::string::npos) {
2418 auto end = jv.find(
'"', pos);
2419 if (end != std::string::npos) {
2420 verdict = jv.substr(pos, end - pos);
2425 if (!terminal_reason.empty()) {
2426 ctx.metadata[
"relay_status"] =
"budget_exhausted_relayed";
2428 "Relay: single-delegate result used "
2429 "(partial — terminal_reason={}, verdict={})",
2430 terminal_reason, verdict.empty() ?
"none" : verdict);
2433 if (verdict ==
"skipped" || verdict.empty()) {
2434 ctx.metadata[
"relay_status"] =
"validation_skipped";
2436 "Relay: single-delegate result used "
2437 "(lead validation skipped per config)");
2439 ctx.metadata[
"relay_status"] =
"validated";
2441 "Relay: single-delegate result used "
2442 "(passed lead validation)");
2453void AgentEngine::execute_pending_pipeline(LoopContext& ctx) {
2454 auto pending = std::move(*ctx.pending_pipeline);
2455 ctx.pending_pipeline.reset();
2457 if (ctx.delegation_depth >= MAX_DELEGATION_DEPTH) {
2458 logger->warn(
"Pipeline rejected: depth {} >= max {}",
2459 ctx.delegation_depth, MAX_DELEGATION_DEPTH);
2461 reject.role =
"user";
2462 reject.content =
"[PIPELINE REJECTED] Maximum delegation "
2464 ctx.messages.push_back(std::move(reject));
2468 set_state(ctx, AgentState::DELEGATING);
2471 auto repo_dir = get_repo_dir();
2472 DelegationManager mgr(run_child_loop_trampoline,
this,
2473 tier_res_, repo_dir,
2474 ensure_sandbox_manager());
2475 if (storage_.create_delegation !=
nullptr) {
2476 mgr.set_storage(&storage_);
2481 auto cb_snap = delegation_callbacks_snapshot();
2482 mgr.set_delegation_callbacks(
2483 cb_snap.start, cb_snap.complete, cb_snap.user_data);
2484 std::vector<DelegationResult> stage_log;
2485 auto result = mgr.execute_pipeline(
2486 ctx, pending.stages, pending.task, stage_log);
2488 std::string tag = result.success ?
"COMPLETE" :
"FAILED";
2489 std::string content =
"[PIPELINE " + tag +
"]";
2490 for (
size_t i = 0; i < stage_log.size(); ++i) {
2491 content +=
"\nStage " + std::to_string(i + 1)
2492 +
" (" + stage_log[i].target_tier +
"): "
2493 + stage_log[i].summary;
2495 content +=
"\nFinal: " + result.summary;
2497 result_msg.role =
"user";
2498 result_msg.content = std::move(content);
2499 ctx.messages.push_back(std::move(result_msg));
2501 set_state(ctx, AgentState::EXECUTING);
2512void AgentEngine::fire_delegation_start(
2513 const LoopContext& ,
2514 const std::string& tier,
2515 const std::string& task) {
2516 if (callbacks_.on_delegation_start !=
nullptr) {
2517 callbacks_.on_delegation_start(
2518 "", tier.c_str(), task.c_str(),
2519 callbacks_.user_data);
2531void AgentEngine::fire_delegation_complete(
2532 const LoopContext& ,
2533 const std::string& tier,
2534 const DelegationResult& result) {
2535 if (callbacks_.on_delegation_complete !=
nullptr) {
2536 callbacks_.on_delegation_complete(
2537 "", tier.c_str(), result.summary.c_str(),
2538 result.success ? 1 : 0,
2539 callbacks_.user_data);
2554bool AgentEngine::should_auto_chain(
2555 const LoopContext& ctx,
2556 const std::string& finish_reason,
2557 const std::string& content) {
2558 if (ctx.locked_tier.empty() || tier_res_.get_tier_param ==
nullptr) {
2562 std::string auto_chain = tier_res_.get_tier_param(
2563 ctx.locked_tier,
"auto_chain", tier_res_.user_data);
2564 if (auto_chain.empty()) {
2568 bool triggered = (finish_reason ==
"length") ||
2569 (finish_reason ==
"stop" &&
2570 response_generator_.is_response_complete(content,
"[]"));
2583bool AgentEngine::try_auto_chain(
2585 const std::string& finish_reason,
2586 const std::string& content) {
2587 if (!should_auto_chain(ctx, finish_reason, content)) {
2591 if (ctx.delegation_depth > 0) {
2592 logger->info(
"[AUTO-CHAIN] child depth={}, completing",
2593 ctx.delegation_depth);
2594 set_state(ctx, AgentState::COMPLETE);
2599 std::string target = tier_res_.get_tier_param(
2600 ctx.locked_tier,
"auto_chain", tier_res_.user_data);
2602 if (!target.empty()) {
2603 logger->info(
"[AUTO-CHAIN] root, tier change to '{}'", target);
2604 TierChangeDirective tc(target,
"auto_chain");
2606 dir_tier_change(ctx, tc, r);
2608 return !target.empty();
2633std::filesystem::path AgentEngine::get_repo_dir() {
2634 if (repo_dir_checked_) {
2635 return cached_repo_dir_.value_or(std::filesystem::path{});
2637 repo_dir_checked_ =
true;
2638 std::filesystem::path resolved;
2639 if (!project_dir_override_.empty()) {
2640 resolved = project_dir_override_;
2641 logger->info(
"Project dir for sandbox snapshots: {} (configured)",
2644 resolved = std::filesystem::current_path();
2645 logger->info(
"Project dir for sandbox snapshots: {} (cwd fallback)",
2648 cached_repo_dir_ = resolved;
2663void AgentEngine::set_project_dir(
const std::filesystem::path& project_dir) {
2664 project_dir_override_ = project_dir;
2665 cached_repo_dir_.reset();
2666 repo_dir_checked_ =
false;
2686 const std::string& reason,
2687 const std::string& delegation_id) {
2690 m.
content =
"[DELEGATION FAILED: resume_delegation] "
2691 + reason +
" (delegation_id=" + delegation_id +
")";
2692 ctx.
messages.push_back(std::move(m));
2708bool AgentEngine::fetch_resume_payload(
2710 const std::string&
id,
2711 nlohmann::json& parsed) {
2712 std::optional<std::string> error;
2714 if (storage_.load_delegation_with_messages ==
nullptr) {
2715 error =
"storage unavailable";
2716 }
else if (!storage_.load_delegation_with_messages(
2717 id.c_str(), raw, storage_.user_data)) {
2718 error =
"unknown delegation_id";
2720 parsed = nlohmann::json::parse(raw,
nullptr,
false);
2721 if (parsed.is_discarded() || !parsed.is_object()) {
2722 error =
"malformed storage payload";
2725 if (
error.has_value()) {
2726 logger->error(
"resume_delegation '{}': {}",
id, *error);
2741bool AgentEngine::resolve_resume_delegation(
2743 PendingDelegation& pending,
2744 std::vector<Message>& out_history) {
2745 const auto&
id = pending.resume_from_delegation_id;
2747 if (!fetch_resume_payload(ctx,
id, j)) {
2750 auto target = j.value(
"target_tier", std::string{});
2751 if (target.empty()) {
2755 pending.target = target;
2756 if (j.contains(
"messages") && j[
"messages"].is_array()) {
2757 for (
const auto& mj : j[
"messages"]) {
2759 m.role = mj.value(
"role",
"");
2760 m.content = mj.value(
"content",
"");
2761 if (!m.role.empty()) {
2762 out_history.push_back(std::move(m));
2766 logger->info(
"resume_delegation '{}': loaded {} messages, target='{}'",
2767 id, out_history.size(), target);
2785DelegationResult AgentEngine::run_pending_delegation(
2787 const PendingDelegation& pending,
2788 std::vector<Message> resume_history) {
2789 std::optional<int> max_turns;
2790 if (pending.max_turns > 0) { max_turns = pending.max_turns; }
2791 DelegationManager mgr(run_child_loop_trampoline,
this,
2792 tier_res_, get_repo_dir(),
2793 ensure_sandbox_manager());
2794 if (storage_.create_delegation !=
nullptr) {
2795 mgr.set_storage(&storage_);
2797 auto cb_snap = delegation_callbacks_snapshot();
2798 mgr.set_delegation_callbacks(
2799 cb_snap.start, cb_snap.complete, cb_snap.user_data);
2800 if (resume_history.empty()) {
2801 return mgr.execute_delegation(
2802 ctx, pending.target, pending.task, max_turns);
2804 return mgr.execute_resume_delegation(
2805 ctx, pending.target, pending.task,
2806 std::move(resume_history), max_turns);
2816SandboxManager* AgentEngine::ensure_sandbox_manager() {
2818 return &*sandbox_mgr_;
2820 auto repo_dir = get_repo_dir();
2821 if (repo_dir.empty()) {
2824 sandbox_mgr_.emplace(repo_dir);
2825 return &*sandbox_mgr_;
2836void AgentEngine::set_system_prompt(
const std::string& prompt) {
2837 system_prompt_ = prompt;
2847 session_logger_ = log;
2863std::vector<Message> AgentEngine::run_turn(
const std::string& input) {
2868 running_flag_.store(
true);
2869 if (conversation_.empty() && !system_prompt_.empty()) {
2871 sys.
role =
"system";
2873 conversation_.push_back(std::move(sys));
2875 auto result = run_drain_loop(input,
"");
2876 running_flag_.store(
false);
2897std::vector<Message> AgentEngine::run_turn_as(
const std::string& tier,
2898 const std::string& input) {
2899 running_flag_.store(
true);
2900 seed_system_prompt_for_tier(tier);
2901 auto result = run_drain_loop(input, tier);
2902 running_flag_.store(
false);
2918void AgentEngine::seed_system_prompt_for_tier(
const std::string& tier) {
2919 if (!conversation_.empty()) {
return; }
2920 auto it = tier_info_.find(tier);
2921 const std::string& sp =
2922 (it != tier_info_.end() && !it->second.system_prompt.empty())
2923 ? it->second.system_prompt
2925 if (sp.empty()) {
return; }
2927 sys.
role =
"system";
2929 conversation_.push_back(std::move(sys));
2940std::vector<Message> AgentEngine::run_drain_loop(
2941 std::string pending,
const std::string& tier_override) {
2942 std::vector<Message> result;
2947 conversation_.push_back(std::move(usr));
2949 size_t sent_len = conversation_.size();
2950 result = run(conversation_, tier_override);
2951 for (
size_t i = sent_len; i < result.size(); ++i) {
2952 conversation_.push_back(result[i]);
2954 auto next = pop_queued_user_message();
2955 if (!next.has_value()) {
break; }
2956 fire_queue_consumed(*next, user_message_queue_depth());
2957 pending = std::move(*next);
2984void AgentEngine::seed_system_prompt(
2985 const std::vector<Message>& new_messages) {
2986 bool caller_has_system =
false;
2987 for (
const auto& m : new_messages) {
2988 if (m.role ==
"system") { caller_has_system =
true;
break; }
2990 if (conversation_.empty() && !system_prompt_.empty()
2991 && !caller_has_system) {
2993 sys.
role =
"system";
2995 conversation_.push_back(std::move(sys));
3006bool AgentEngine::prepare_next_turn(std::vector<Message>& pending) {
3007 auto next = pop_queued_user_message();
3008 if (!next.has_value()) {
return false; }
3009 fire_queue_consumed(*next, user_message_queue_depth());
3012 usr.
content = std::move(*next);
3014 pending.push_back(std::move(usr));
3025std::vector<Message> AgentEngine::run_turn(std::vector<Message> new_messages) {
3030 running_flag_.store(
true);
3031 seed_system_prompt(new_messages);
3032 std::vector<Message> pending = std::move(new_messages);
3033 std::vector<Message> result;
3035 for (
auto& m : pending) {
3036 conversation_.push_back(std::move(m));
3038 size_t sent_len = conversation_.size();
3039 result = run(conversation_);
3040 for (
size_t i = sent_len; i < result.size(); ++i) {
3041 conversation_.push_back(result[i]);
3043 if (!prepare_next_turn(pending)) {
break; }
3045 running_flag_.store(
false);
3060int AgentEngine::run_streaming(
3061 const std::string& input,
3066 if (session_logger_) {
3067 session_logger_->log_user_input(input);
3071 if (session_logger_ && session_logger_->is_open()) {
3073 SessionLogger::raw_token_callback,
3082 Ctx sctx{&filter, cancel_flag,
this};
3086 auto* c =
static_cast<Ctx*
>(ud);
3087 if (c->cancel && *c->cancel) {
3088 c->engine->interrupt();
3091 c->filter->on_token(t, l);
3093 cbs.user_data = &sctx;
3096 auto result = run_turn(input);
3099 if (session_logger_) {
3100 session_logger_->end_turn();
3102 if (cancel_flag && *cancel_flag) {
return 1; }
3114 const std::vector<Message>& messages) {
3116 for (
const auto& m : messages) {
3117 if (m.role !=
"user" || m.content.empty()) {
continue; }
3118 if (!echo.empty()) { echo +=
'\n'; }
3140int AgentEngine::run_streaming(
3141 std::vector<Message> new_messages,
3146 if (session_logger_) {
3151 if (session_logger_ && session_logger_->is_open()) {
3153 SessionLogger::raw_token_callback,
3162 Ctx sctx{&filter, cancel_flag,
this};
3166 auto* c =
static_cast<Ctx*
>(ud);
3167 if (c->cancel && *c->cancel) {
3168 c->engine->interrupt();
3171 c->filter->on_token(t, l);
3173 cbs.user_data = &sctx;
3176 auto result = run_turn(std::move(new_messages));
3179 if (session_logger_) {
3180 session_logger_->end_turn();
3182 if (cancel_flag && *cancel_flag) {
return 1; }
3191void AgentEngine::clear_conversation() {
3192 conversation_.clear();
3193 logger->info(
"conversation cleared");
3202size_t AgentEngine::message_count()
const {
3203 return conversation_.size();
3212const std::vector<Message>& AgentEngine::get_messages()
const {
3213 return conversation_;
3223bool AgentEngine::queue_user_message(
const std::string& message) {
3224 std::lock_guard lock(queue_mutex_);
3225 int cap = loop_config_.message_queue_capacity;
3226 if (cap < 0) { cap = 0; }
3227 if (user_message_queue_.size()
3228 >=
static_cast<size_t>(cap)) {
3231 user_message_queue_.push_back(message);
3232 logger->info(
"queued mid-gen user message: depth={}",
3233 user_message_queue_.size());
3242size_t AgentEngine::user_message_queue_depth()
const {
3243 std::lock_guard lock(queue_mutex_);
3244 return user_message_queue_.size();
3252void AgentEngine::clear_user_message_queue() {
3253 std::lock_guard lock(queue_mutex_);
3254 size_t dropped = user_message_queue_.size();
3255 user_message_queue_.clear();
3257 logger->info(
"cleared mid-gen queue: dropped={}", dropped);
3266void AgentEngine::set_message_queue_capacity(
int cap) {
3267 std::lock_guard lock(queue_mutex_);
3268 loop_config_.message_queue_capacity = cap < 0 ? 0 : cap;
3276std::optional<std::string>
3277AgentEngine::pop_queued_user_message() {
3278 std::lock_guard lock(queue_mutex_);
3279 if (user_message_queue_.empty()) {
3280 return std::nullopt;
3282 std::string front = std::move(user_message_queue_.front());
3283 user_message_queue_.pop_front();
3292void AgentEngine::fire_queue_consumed(
const std::string& consumed,
3294 if (queue_observer_ !=
nullptr) {
3296 consumed.c_str(), remaining, queue_observer_data_);
3305void AgentEngine::set_queue_observer(
3306 void (*observer)(
const char*,
size_t,
void*),
3308 queue_observer_ = observer;
3309 queue_observer_data_ = user_data;
3323void AgentEngine::set_state_observer(
3324 void (*observer)(
int,
void*),
3326 state_observer_ = observer;
3327 state_observer_data_ = user_data;
3328 response_generator_.set_state_observer(observer, user_data);
3343 const std::vector<const Directive*>& dirs,
3346 ->directive_processor().process(ctx, dirs);
3361void AgentEngine::set_tier_info(
3362 const std::string& name,
3365 tier_info_[name] = info;
3375bool AgentEngine::has_tier(
const std::string& name)
const {
3376 return tier_info_.count(name) > 0;
3387std::vector<std::string> AgentEngine::get_tier_allowed_tools(
3388 const std::string& name)
const {
3389 auto it = tier_info_.find(name);
3390 if (it == tier_info_.end()) {
return {}; }
3391 return it->second.allowed_tools;
3401const std::string& AgentEngine::tier_system_prompt(
3402 const std::string& name)
const {
3403 static const std::string kEmpty;
3404 auto it = tier_info_.find(name);
3405 return it == tier_info_.end() ? kEmpty : it->second.system_prompt;
3414void AgentEngine::set_relay_single_delegate(
const std::string& name) {
3415 relay_single_delegate_tiers_.insert(name);
3425void AgentEngine::set_handoff_rules(
3426 const std::unordered_map<std::string,
3427 std::vector<std::string>>& rules)
3429 handoff_rules_ = rules;
3430 wire_internal_tier_resolution();
3441 const std::string& name,
void* ud) {
3443 auto it = self->tier_info_.find(name);
3445 if (it == self->tier_info_.end()) {
3461bool AgentEngine::tri_tier_exists(
const std::string& name,
void* ud) {
3462 auto* self =
static_cast<AgentEngine*
>(ud);
3463 return self->tier_info_.count(name) > 0;
3475std::vector<std::string> AgentEngine::tri_get_handoff_targets(
3476 const std::string& name,
void* ud) {
3477 auto* self =
static_cast<AgentEngine*
>(ud);
3478 auto it = self->handoff_rules_.find(name);
3479 std::vector<std::string> result;
3480 if (it != self->handoff_rules_.end()) { result = it->second; }
3491static std::string
join_csv(
const std::vector<std::string>& v) {
3493 for (
const auto& s : v) {
3494 if (!out.empty()) { out +=
','; }
3524std::string AgentEngine::tri_get_tier_param(
const std::string& name,
3525 const std::string& param,
void* ud) {
3526 auto* self =
static_cast<AgentEngine*
>(ud);
3527 auto it = self->tier_info_.find(name);
3528 if (it == self->tier_info_.end()) {
return ""; }
3529 const auto& info = it->second;
3531 if (param ==
"explicit_completion") {
3532 result = info.explicit_completion ?
"true" :
"false";
3534 if (param ==
"max_iterations" && info.max_iterations_override >= 0) {
3535 result = std::to_string(info.max_iterations_override);
3537 if (param ==
"max_tool_calls_per_turn"
3538 && info.max_tool_calls_per_turn_override >= 0) {
3539 result = std::to_string(info.max_tool_calls_per_turn_override);
3541 if (param ==
"max_consecutive_empty_turns"
3542 && info.max_consecutive_empty_turns_override >= 0) {
3543 result = std::to_string(info.max_consecutive_empty_turns_override);
3545 if (param ==
"allowed_tools" && !info.allowed_tools.empty()) {
3546 result =
join_csv(info.allowed_tools);
3556void AgentEngine::wire_internal_tier_resolution() {
3557 TierResolutionInterface tri;
3558 tri.resolve_tier = &AgentEngine::tri_resolve_tier;
3559 tri.tier_exists = &AgentEngine::tri_tier_exists;
3560 tri.get_handoff_targets = &AgentEngine::tri_get_handoff_targets;
3561 tri.get_tier_param = &AgentEngine::tri_get_tier_param;
3562 tri.user_data =
this;
3563 set_tier_resolution(tri);
Core agent execution engine.
void run_loop(LoopContext &ctx, bool inherit_interrupt=false)
Run the engine loop on a pre-built context.
AgentEngine(const InferenceInterface &inference, const LoopConfig &loop_config, const CompactionConfig &compaction_config)
Construct an agent engine.
Manages session_model.log for raw streaming content.
Streaming filter that removes <think> blocks from output.
void set_raw_callback(TokenCallback cb, void *ud)
Set optional raw callback (receives ALL tokens unfiltered).
void flush()
Flush any buffered partial tag content.
Testable free function extracted from fire_delegate_complete_hook.
std::string build_delegate_complete_json(const std::string &target, bool success, const std::string &summary)
Build and dump the ON_DELEGATE_COMPLETE hook JSON payload.
DelegationManager — child loop creation and execution.
Core agent execution engine.
ent_decision_t
Consumer decision returned from delegation callbacks.
entropic_directive_type_t
Directive types emitted by MCP tool results.
@ ENTROPIC_DIRECTIVE_TIER_CHANGE
Switch active tier.
@ ENTROPIC_DIRECTIVE_STOP_PROCESSING
Halt directive processing.
@ ENTROPIC_DIRECTIVE_PRUNE_MESSAGES
Prune old tool results.
@ ENTROPIC_DIRECTIVE_INJECT_CONTEXT
Inject message into context.
@ ENTROPIC_DIRECTIVE_COMPLETE
Mark task complete.
@ ENTROPIC_DIRECTIVE_NOTIFY_PRESENTER
Generic UI notification passthrough.
@ ENTROPIC_DIRECTIVE_CLEAR_SELF_TODOS
Clear self-directed todos (engine no-op)
@ ENTROPIC_DIRECTIVE_PHASE_CHANGE
Switch active inference phase.
@ ENTROPIC_DIRECTIVE_PIPELINE
Multi-stage sequential execution.
@ ENTROPIC_DIRECTIVE_CONTEXT_ANCHOR
Replace context anchor.
@ ENTROPIC_DIRECTIVE_DELEGATE
Route to another identity.
entropic_hook_point_t
Hook points in the engine lifecycle.
@ ENTROPIC_HOOK_ON_LOOP_START
20: Agentic loop entry
@ ENTROPIC_HOOK_ON_DELEGATE
8: Delegation to child tier started
@ ENTROPIC_HOOK_ON_LOOP_END
21: Agentic loop exit
@ ENTROPIC_HOOK_ON_STATE_CHANGE
6: Engine state machine transition
@ ENTROPIC_HOOK_PRE_GENERATE
0: Before inference generate call
@ ENTROPIC_HOOK_ON_CONTEXT_ASSEMBLE
10: Context window assembled
@ ENTROPIC_HOOK_ON_LOOP_ITERATION
5: Each agentic loop iteration
@ ENTROPIC_HOOK_ON_ERROR
7: Async error occurred
@ ENTROPIC_HOOK_POST_GENERATE
1: After inference generate returns
@ ENTROPIC_HOOK_ON_DELEGATE_COMPLETE
9: Child delegation completed
@ ENTROPIC_HOOK_ON_COMPLETE
entropic.complete MCP tool called — pre-hook, can cancel.
spdlog initialization and logger access.
ENTROPIC_EXPORT std::shared_ptr< spdlog::logger > get(const std::string &name)
Get or create a named logger.
Activate model on GPU (WARM → ACTIVE).
static std::string build_coverage_gap_message(const std::string &tier, const DelegationResult &result)
Build the [COVERAGE GAP] message body that goes back to lead when a relay-tier child returns coverage...
static void fire_context_assemble_hook(const HookInterface &hooks, const LoopContext &ctx)
Build + fire the ON_CONTEXT_ASSEMBLE info hook.
static std::string extract_system_prompt(const std::vector< Message > &messages)
Extract the first system message content from messages.
static double now_seconds()
Get current time as seconds since epoch.
static ToolCall build_tool_call_from_json(const nlohmann::json &obj, const std::string &tc_str)
Parse tool calls from raw model output.
ToolResultKind
Categorical outcome of a single tool invocation.
@ error
Tool server returned an error payload.
@ delegation_failed
entropic.delegate child failed (terminal_reason or budget). (#7, v2.1.4)
static std::string build_tool_results_json(const std::vector< Message > &messages)
Build a JSON array of tool results from context messages.
static std::vector< ToolCall > decode_tool_calls_json(const std::string &tc_str)
Decode a JSON tool-calls array string into ToolCall vector.
static std::string json_escape_engine(const std::string &s)
JSON-escape a string (no surrounding quotes).
static std::string build_tool_evidence(const std::vector< Message > &messages, size_t max_results, size_t chars_per_result)
Build un-pruned tool-result evidence for the validator.
@ count
Sentinel — MUST remain last.
static void remove_anchor_messages(LoopContext &ctx, const std::string &key)
Remove messages with a specific anchor key.
const char * agent_state_name(AgentState state)
Get the string name for an AgentState value.
static void fold_delegation_summary(LoopContext &ctx, const std::string &summary)
gh#119 (v2.9.17): fold child summary into the lead's empty delegate turn.
static void push_delegation_repeat_blocked(LoopContext &ctx, const std::string &target, int n)
Append a "stop retrying same target" reject message (gh#64).
static void fire_loop_start_hook(const HookInterface &hooks, const LoopContext &ctx)
Build + fire the ON_LOOP_START info hook.
static std::string build_tool_manifest(const std::vector< Message > &messages)
Build tool call manifest from conversation messages.
static std::string format_tool_evidence_entry(const Message &msg, size_t chars_per_result)
Format one tool-result message for the evidence block.
static const std::array< const char *, 256 > & json_escape_table()
256-entry table mapping bytes to JSON escape sequences.
static std::string enrich_manifest_with_history(const std::string &base, const ToolExecutionInterface &tool_exec)
Append executor history (if any) to the tool manifest.
static void push_delegation_result(LoopContext &ctx, const std::string &target, const DelegationResult &result)
Append a delegation result message to the loop context.
static std::string join_csv(const std::vector< std::string > &v)
Join a vector of strings with comma separators.
static void run_child_loop_trampoline(LoopContext &ctx, void *user_data)
Trampoline for DelegationManager to call engine loop.
const char * result_kind_to_string(ToolResultKind kind)
Serialize a ToolResultKind to its wire-stable string form.
static void fire_hook_info(const HookInterface &hooks, entropic_hook_point_t point, const char *json)
Fire an informational hook if the interface is wired.
void(*)(const char *, size_t, void *) TokenCallback
Token callback type matching the C API signature.
static void push_delegation_rejected(LoopContext &ctx)
Append a delegation rejection message to the loop context.
static std::string concat_user_echo(const std::vector< Message > &messages)
Concatenate user-role message text for session-log echo.
AgentState
C++ enum class for agent execution states.
static void fire_loop_iteration_hook(const HookInterface &hooks, const LoopContext &ctx)
Build + fire the ON_LOOP_ITERATION info hook.
static void push_resume_failure(LoopContext &ctx, const std::string &reason, const std::string &delegation_id)
Lazy session-scoped SandboxManager accessor (gh#33, v2.1.6).
static void push_delegation_cycle_rejected(LoopContext &ctx, const std::string &target)
Append a structured cycle-rejection message to the loop context so the model can recover.
static void fire_loop_end_hook(const HookInterface &hooks, const LoopContext &ctx)
Build + fire the ON_LOOP_END info hook.
Filesystem-based sandbox isolation for delegations.
Request describing a delegation that is about to run.
Result of a finalized delegation, delivered to the consumer.
Grouped consumer-registered delegation callbacks (gh#29).
Resolved tier information for building child delegation contexts.
Auto-compaction configuration.
Update a keyed persistent context anchor.
Engine-level hooks called during context management.
Result returned from a child delegation loop.
bool success
Whether child reached COMPLETE via real entropic.complete.
std::string summary
Final summary from child.
std::vector< std::string > suggested_files
Issue #10 (v2.1.4): file paths the lead should inspect to fill the coverage gap.
std::string gap_description
Issue #10 (v2.1.4): concrete description of what the child's answer DOES NOT cover.
Aggregate result of processing a batch of directives.
Callback function pointer types for engine events.
void(* on_stream_chunk)(const char *chunk, size_t len, void *ud)
Per-token streaming.
Configuration for the agentic loop.
Mutable state carried through the agentic loop.
std::vector< std::string > delegation_ancestor_tiers
Tier stack from root to this loop (P1-9, 2.0.6-rc16)
std::string last_failed_delegation_target
gh#64: target tier of the most recent FAILED delegation (DelegationResult.success == false).
LoopMetrics metrics
Timing and counts.
int effective_max_iterations
Per-identity override (-1 = LoopConfig, P3-18)
int consecutive_duplicate_attempts
Stuck-model detector.
int consecutive_failed_delegations
gh#64: count of consecutive failed delegations against last_failed_delegation_target.
std::string conversation_id
Conversation ID for storage (v1.8.8)
int consecutive_errors
Error streak counter.
bool has_pending_tool_results
Tool results awaiting presentation.
std::unordered_map< std::string, std::string > metadata
Runtime metadata.
int effective_max_tool_calls_per_turn
Per-identity override (-1 = LoopConfig, P3-18)
int delegation_depth
0 = root, 1+ = child
AgentState state
Current state.
std::vector< Message > messages
Conversation history.
std::string locked_tier
Tier locked for this loop ("" = none)
double start_time
Loop start (seconds since epoch)
int tokens_used
Total tokens consumed.
int errors
Total errors encountered.
int duration_ms() const
Get loop duration in milliseconds.
int tool_calls
Total tool calls executed.
double end_time
Loop end (seconds since epoch)
int iterations
Total iterations completed.
A message in a conversation.
std::unordered_map< std::string, std::string > metadata
Arbitrary metadata.
std::string content
Message text content (always populated)
std::string role
Message role.
Storage interface for conversation persistence.
Tier resolution callbacks for delegation and auto-chain.
std::string(* get_tier_param)(const std::string &tier_name, const std::string ¶m_name, void *user_data)
Get a string parameter from tier identity frontmatter.
UTF-8 validation + replacement at every system boundary where bytes change ownership.