Current section
Files
Jump to
Current section
Files
src/pig@agent@runtime.erl
-module(pig@agent@runtime).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/pig/agent/runtime.gleam").
-export([start_with_state/2, start/1, run/3, run_continue/2, stop/1, history/2, try_run/3, try_run_continue/2, supervised/5, supervised_with_session_store/4]).
-export_type([runtime_config/0, runtime_msg/0, post_commit_disposition/0, session_state/0, runtime_state/0, blocked_tool/0, tool_result/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(
" Sans-IO runtime interpreter for the pig agent.\n"
"\n"
" The runtime is an OTP actor that:\n"
" 1. Receives prompts (Run) or control messages (Stop)\n"
" 2. Calls `update.update(state, msg)` — pure state machine\n"
" 3. For each effect, applies hooks then executes\n"
" 4. Produces SessionEvent values and sends to dispatcher\n"
" 5. Feeds effect results back as new AgentMsg values\n"
"\n"
" The core logic (update.gleam) is pure. This module is all IO.\n"
).
-type runtime_config() :: {runtime_config,
fun((list(pig_protocol@message:message()), list(pig_protocol@tool_definition:tool_definition())) -> {ok,
pig_protocol@inference:inference_result()} |
{error, pig_protocol@error:ai_error()}),
pig@tool:tool_registry(),
list(pig@hooks:hooks()),
gleam@erlang@process:subject(pig@obs@dispatcher:dispatcher_message()),
binary(),
integer()}.
-type runtime_msg() :: {run,
binary(),
gleam@erlang@process:subject({ok, pig_protocol@message:message()} |
{error, pig@run_error:run_error()})} |
{continue,
gleam@erlang@process:subject({ok, pig_protocol@message:message()} |
{error, pig@run_error:run_error()})} |
{get_history,
gleam@erlang@process:subject(list(pig_protocol@message:message()))} |
stop.
-type post_commit_disposition() :: resume_from_history |
{return_message, pig_protocol@message:message()} |
{return_ai_error, pig_protocol@error:ai_error()}.
-type session_state() :: session_disabled |
{session_ready,
pig@session_store:session_store(),
gleam@option:option(binary())} |
{session_pending,
pig@session_store:session_store(),
pig@session_store:session_commit(),
pig@agent@state:agent_state(),
post_commit_disposition()}.
-type runtime_state() :: {runtime_state,
pig@agent@state:agent_state(),
runtime_config(),
session_state()}.
-type blocked_tool() :: {blocked_tool,
pig_protocol@message:tool_call(),
binary(),
binary()}.
-type tool_result() :: {tool_result,
binary(),
{ok, gleam@json:json()} | {error, pig@tool:tool_error()},
integer()}.
-file("src/pig/agent/runtime.gleam", 1006).
-spec find_result(list(tool_result()), binary()) -> {ok,
{{ok, gleam@json:json()} | {error, pig@tool:tool_error()}, integer()}} |
{error, nil}.
find_result(Results, Call_id) ->
case gleam@list:find(
Results,
fun(R) -> erlang:element(2, R) =:= Call_id end
) of
{ok, R@1} ->
{ok, {erlang:element(3, R@1), erlang:element(4, R@1)}};
{error, nil} ->
{error, nil}
end.
-file("src/pig/agent/runtime.gleam", 999).
-spec find_blocked(list(blocked_tool()), binary()) -> {ok, blocked_tool()} |
{error, nil}.
find_blocked(Blocked, Call_id) ->
gleam@list:find(
Blocked,
fun(B) -> erlang:element(2, erlang:element(2, B)) =:= Call_id end
).
-file("src/pig/agent/runtime.gleam", 1023).
-spec emit_hook_acted_list(
gleam@erlang@process:subject(pig@obs@dispatcher:dispatcher_message()),
list(binary()),
pig@obs@events:hook_point(),
binary(),
binary()
) -> nil.
emit_hook_acted_list(Disp, Transformer_names, Hook, Action_type, Description) ->
gleam@list:each(
Transformer_names,
fun(Name) ->
pig@obs@emit:to_dispatcher(
Disp,
{hook_acted,
Name,
Hook,
{hook_action_detail, Action_type, Description}}
)
end
).
-file("src/pig/agent/runtime.gleam", 1016).
-spec is_error({ok, any()} | {error, any()}) -> boolean().
is_error(Result) ->
case Result of
{ok, _} ->
false;
{error, _} ->
true
end.
-file("src/pig/agent/runtime.gleam", 960).
-spec spawn_and_collect(
runtime_config(),
list(pig_protocol@message:tool_call())
) -> list(tool_result()).
spawn_and_collect(Config, Calls) ->
Disp = erlang:element(5, Config),
Pairs = gleam@list:map(
Calls,
fun(Call) ->
Reply_subject = gleam@erlang@process:new_subject(),
Pid = proc_lib:spawn_link(
fun() ->
pig@obs@emit:to_dispatcher(Disp, {tool_started, Call}),
Start_time = pig@obs@events:system_time(),
Result = pig@tool@execution:execute_tool(
erlang:element(3, Config),
Call
),
Duration = pig@obs@events:system_time() - Start_time,
gleam@erlang@process:send(
Reply_subject,
{tool_result, erlang:element(2, Call), Result, Duration}
)
end
),
{Pid, Reply_subject, erlang:element(2, Call)}
end
),
Timeout_ms = 5000,
gleam@list:map(
Pairs,
fun(Pair) ->
{Pid@1, Subject, Call_id} = Pair,
case gleam@erlang@process:'receive'(Subject, Timeout_ms) of
{ok, Result@1} ->
Result@1;
{error, nil} ->
gleam@erlang@process:kill(Pid@1),
logging:log(
error,
<<<<"Tool execution timed out after "/utf8,
(erlang:integer_to_binary(Timeout_ms))/binary>>/binary,
"ms"/utf8>>
),
{tool_result,
Call_id,
{error,
{tool_error, <<"Tool execution timed out"/utf8>>}},
Timeout_ms}
end
end
).
-file("src/pig/agent/runtime.gleam", 936).
-spec partition_by_hook_decision(
list(pig@hooks:hooks()),
list(pig_protocol@message:tool_call())
) -> {list(blocked_tool()), list(pig_protocol@message:tool_call())}.
partition_by_hook_decision(Hooks_list, Calls) ->
gleam@list:fold(
Calls,
{[], []},
fun(Acc, Call) ->
{Blocked_acc, Allowed_acc} = Acc,
Hook_event = {tool_call_event,
erlang:element(3, Call),
erlang:element(2, Call),
erlang:element(4, Call)},
case pig@hooks:decide_tool_call(Hooks_list, Hook_event) of
tool_allowed ->
{Blocked_acc, lists:append(Allowed_acc, [Call])};
{tool_blocked, Hook_name, Reason} ->
{lists:append(
Blocked_acc,
[{blocked_tool, Call, Hook_name, Reason}]
),
Allowed_acc}
end
end
).
-file("src/pig/agent/runtime.gleam", 820).
-spec execute_valid_tools_effect(
runtime_config(),
pig@agent@state:agent_state(),
list(pig_protocol@message:tool_call()),
fun((list({pig_protocol@message:tool_call(),
{ok, gleam@json:json()} | {error, pig@tool:tool_error()}})) -> pig@agent@msg:agent_msg())
) -> {pig@agent@state:agent_state(), pig@agent@msg:agent_msg()}.
execute_valid_tools_effect(Config, Agent_st, Calls, On_results) ->
Disp = erlang:element(5, Config),
{Blocked, Allowed} = partition_by_hook_decision(
erlang:element(4, Config),
Calls
),
gleam@list:each(
Blocked,
fun(B) ->
pig@obs@emit:to_dispatcher(
Disp,
{tool_blocked,
erlang:element(2, B),
erlang:element(3, B),
erlang:element(4, B)}
),
pig@obs@emit:to_dispatcher(
Disp,
{hook_acted,
erlang:element(3, B),
before_tool_call,
{hook_action_detail,
<<"block"/utf8>>,
<<"Blocked tool: "/utf8, (erlang:element(4, B))/binary>>}}
)
end
),
Results = spawn_and_collect(Config, Allowed),
_ = gleam@list:map(
Allowed,
fun(Call) ->
{Result@1, Duration@1} = case find_result(
Results,
erlang:element(2, Call)
) of
{ok, {Result, Duration}} -> {Result, Duration};
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"pig/agent/runtime"/utf8>>,
function => <<"execute_valid_tools_effect"/utf8>>,
line => 862,
value => _assert_fail,
start => 27397,
'end' => 27463,
pattern_start => 27408,
pattern_end => 27431})
end,
Raw_content = case Result@1 of
{ok, Json_result} ->
gleam@json:to_string(Json_result);
{error, Tool_err} ->
<<"Tool error: "/utf8,
(pig@tool:error_message(Tool_err))/binary>>
end,
Result_event = {tool_result_event,
erlang:element(3, Call),
erlang:element(2, Call),
Raw_content,
is_error(Result@1),
Duration@1},
Final_content = case pig@hooks:decide_tool_result(
erlang:element(4, Config),
Result_event
) of
{result_unchanged, _} ->
Raw_content;
{result_transformed, Final_event, Transformers} ->
emit_hook_acted_list(
Disp,
Transformers,
after_tool_call,
<<"transform"/utf8>>,
<<"Transformed result"/utf8>>
),
erlang:element(4, Final_event)
end,
pig@obs@emit:to_dispatcher(
Disp,
{tool_executed, Call, Final_content, Duration@1}
)
end
),
All_results = gleam@list:map(
Calls,
fun(Call@1) -> case find_blocked(Blocked, erlang:element(2, Call@1)) of
{ok, B@1} ->
{Call@1,
{error,
{tool_error,
<<<<<<"Tool blocked by '"/utf8,
(erlang:element(3, B@1))/binary>>/binary,
"': "/utf8>>/binary,
(erlang:element(4, B@1))/binary>>}}};
{error, nil} ->
Res@1 = case find_result(Results, erlang:element(2, Call@1)) of
{ok, {Res, _}} -> Res;
_assert_fail@1 ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"pig/agent/runtime"/utf8>>,
function => <<"execute_valid_tools_effect"/utf8>>,
line => 913,
value => _assert_fail@1,
start => 28911,
'end' => 28970,
pattern_start => 28922,
pattern_end => 28938})
end,
{Call@1, Res@1}
end end
),
{Agent_st, On_results(All_results)}.
-file("src/pig/agent/runtime.gleam", 800).
?DOC(
" Execute ExecuteTools: apply tool call hooks, execute allowed tools\n"
" in parallel, apply result hooks, emit events.\n"
).
-spec execute_tools_effect(
runtime_config(),
pig@agent@state:agent_state(),
list(pig_protocol@message:tool_call()),
fun((list({pig_protocol@message:tool_call(),
{ok, gleam@json:json()} | {error, pig@tool:tool_error()}})) -> pig@agent@msg:agent_msg())
) -> {pig@agent@state:agent_state(), pig@agent@msg:agent_msg()}.
execute_tools_effect(Config, Agent_st, Calls, On_results) ->
case pig@tool@execution:validate_tool_calls(Calls) of
{error, Batch_error} ->
Rejected = gleam@list:map(
Calls,
fun(Call) ->
{Call, {error, {invalid_tool_call_batch, Batch_error}}}
end
),
{Agent_st, On_results(Rejected)};
{ok, nil} ->
execute_valid_tools_effect(Config, Agent_st, Calls, On_results)
end.
-file("src/pig/agent/runtime.gleam", 718).
?DOC(
" Execute CallProvider: apply before_inference hooks, call provider,\n"
" emit events, fire notification hooks.\n"
).
-spec execute_call_provider(
runtime_config(),
pig@agent@state:agent_state(),
list(pig_protocol@message:message()),
list(pig_protocol@tool_definition:tool_definition()),
fun(({ok, pig_protocol@inference:inference_result()} |
{error, pig_protocol@error:ai_error()}) -> pig@agent@msg:agent_msg())
) -> {pig@agent@state:agent_state(), pig@agent@msg:agent_msg()}.
execute_call_provider(Config, Agent_st, Messages, Tools, On_response) ->
Disp = erlang:element(5, Config),
Model = erlang:element(6, Config),
Before_event = {before_inference_event, Model, Messages},
Final_msgs = case pig@hooks:decide_messages(
erlang:element(4, Config),
Before_event
) of
{messages_unchanged, _} ->
Messages;
{messages_replaced, Final_messages, Transformers} ->
emit_hook_acted_list(
Disp,
Transformers,
before_inference,
<<"transform"/utf8>>,
<<"Transformed messages before inference"/utf8>>
),
Final_messages
end,
Msg_count = erlang:length(Final_msgs),
pig@obs@emit:to_dispatcher(Disp, {inference_started, Model, Msg_count}),
Start_time = pig@obs@events:system_time(),
Result = case (erlang:element(2, Config))(Final_msgs, Tools) of
{ok, Inference_result} ->
Msg = erlang:element(2, Inference_result),
Meta = erlang:element(3, Inference_result),
Duration = pig@obs@events:system_time() - Start_time,
Response_model = case erlang:element(3, Meta) of
{some, _} = M ->
M;
none ->
{some, Model}
end,
pig@obs@emit:to_dispatcher(
Disp,
{inference_completed,
Msg,
erlang:element(2, Meta),
Response_model,
erlang:element(4, Meta),
erlang:element(5, Meta),
erlang:element(6, Meta),
Duration,
erlang:element(3, Agent_st)}
),
pig@hooks:notify_after_inference(
erlang:element(4, Config),
{after_inference_event, Model, Msg, Duration}
),
{ok, Inference_result};
{error, E} ->
Duration@1 = pig@obs@events:system_time() - Start_time,
pig@obs@emit:to_dispatcher(
Disp,
{inference_failed, E, Duration@1, erlang:element(3, Agent_st)}
),
pig@hooks:notify_error(
erlang:element(4, Config),
{error_event, Model, E}
),
{error, E}
end,
{Agent_st, On_response(Result)}.
-file("src/pig/agent/runtime.gleam", 703).
?DOC(" Execute a single effect: apply hooks, execute, emit events.\n").
-spec execute_effect(
runtime_config(),
pig@agent@state:agent_state(),
pig@agent@effect:effect(pig@agent@msg:agent_msg())
) -> {pig@agent@state:agent_state(), pig@agent@msg:agent_msg()}.
execute_effect(Config, Agent_st, Eff) ->
case Eff of
{call_provider, Messages, Tools, On_response} ->
execute_call_provider(
Config,
Agent_st,
Messages,
Tools,
On_response
);
{execute_tools, Calls, On_results} ->
execute_tools_effect(Config, Agent_st, Calls, On_results)
end.
-file("src/pig/agent/runtime.gleam", 493).
-spec accept_committed_session(
pig@session_store:session_store(),
pig@session_store:session_commit(),
pig@agent@state:agent_state(),
post_commit_disposition(),
pig@session_store:session()
) -> {ok, session_state()} |
{error, {session_state(), pig@session_store:session_error()}}.
accept_committed_session(Store, Commit, Candidate, Disposition, Committed) ->
case erlang:element(2, Committed) =:= {some, erlang:element(2, Commit)} of
true ->
{ok, {session_ready, Store, erlang:element(2, Committed)}};
false ->
{error,
{{session_pending, Store, Commit, Candidate, Disposition},
{parent_conflict,
{some, erlang:element(2, Commit)},
erlang:element(2, Committed)}}}
end.
-file("src/pig/agent/runtime.gleam", 473).
-spec commit_messages(
pig@session_store:session_store(),
gleam@option:option(binary()),
list(pig_protocol@message:message()),
pig@agent@state:agent_state(),
post_commit_disposition()
) -> {ok, session_state()} |
{error, {session_state(), pig@session_store:session_error()}}.
commit_messages(Store, Head, Messages, Candidate, Disposition) ->
{session_store, _, Store_commit} = Store,
Commit = pig@session_store:new_commit(Head, Messages),
case Store_commit(Commit) of
{ok, Committed} ->
accept_committed_session(
Store,
Commit,
Candidate,
Disposition,
Committed
);
{error, Session_error} ->
{error,
{{session_pending, Store, Commit, Candidate, Disposition},
Session_error}}
end.
-file("src/pig/agent/runtime.gleam", 456).
-spec commit_transition(
session_state(),
pig@agent@state:agent_state(),
pig@agent@state:agent_state(),
post_commit_disposition()
) -> {ok, session_state()} |
{error, {session_state(), pig@session_store:session_error()}}.
commit_transition(Session, Previous, Candidate, Disposition) ->
Delta = gleam@list:drop(
erlang:element(3, Candidate),
erlang:length(erlang:element(3, Previous))
),
case {Session, Delta} of
{{session_pending, _, _, _, _}, _} ->
erlang:error(#{gleam_error => panic,
message => <<"cannot commit while a session commit is pending"/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"pig/agent/runtime"/utf8>>,
function => <<"commit_transition"/utf8>>,
line => 465});
{_, []} ->
{ok, Session};
{session_disabled, _} ->
{ok, Session};
{{session_ready, Store, Head}, Messages} ->
commit_messages(Store, Head, Messages, Candidate, Disposition)
end.
-file("src/pig/agent/runtime.gleam", 563).
-spec execute_effects(
runtime_config(),
pig@agent@state:agent_state(),
session_state(),
list(pig@agent@effect:effect(pig@agent@msg:agent_msg()))
) -> {pig@agent@state:agent_state(),
session_state(),
{ok, pig_protocol@message:message()} | {error, pig@run_error:run_error()}}.
execute_effects(Config, Agent_st, Session, Effs) ->
{Updated_st, Response_msgs} = gleam@list:fold(
Effs,
{Agent_st, []},
fun(Acc, Eff) ->
{St_acc, Msgs_acc} = Acc,
{New_st_acc, Response_msg} = execute_effect(Config, St_acc, Eff),
{New_st_acc, lists:append(Msgs_acc, [Response_msg])}
end
),
case Response_msgs of
[First_msg | _] ->
do_loop(Config, Updated_st, Session, First_msg);
[] ->
{Updated_st,
Session,
{error, {runtime, <<"no response from effects"/utf8>>}}}
end.
-file("src/pig/agent/runtime.gleam", 546).
-spec complete_disposition(
runtime_config(),
pig@agent@state:agent_state(),
session_state(),
post_commit_disposition()
) -> {pig@agent@state:agent_state(),
session_state(),
{ok, pig_protocol@message:message()} | {error, pig@run_error:run_error()}}.
complete_disposition(Config, Candidate, Session, Disposition) ->
case Disposition of
{return_message, Message} ->
{Candidate, Session, {ok, Message}};
{return_ai_error, Error} ->
{Candidate, Session, {error, {inference, Error}}};
resume_from_history ->
resume_from_history(Config, Candidate, Session)
end.
-file("src/pig/agent/runtime.gleam", 438).
-spec commit_or_pending(
runtime_config(),
session_state(),
pig@agent@state:agent_state(),
pig@agent@state:agent_state(),
post_commit_disposition()
) -> {pig@agent@state:agent_state(),
session_state(),
{ok, pig_protocol@message:message()} | {error, pig@run_error:run_error()}}.
commit_or_pending(Config, Session, Previous, Candidate, Disposition) ->
case commit_transition(Session, Previous, Candidate, Disposition) of
{ok, Next_session} ->
complete_disposition(Config, Candidate, Next_session, Disposition);
{error, {Pending, Session_error}} ->
{Previous, Pending, {error, {session, Session_error}}}
end.
-file("src/pig/agent/runtime.gleam", 400).
-spec do_loop(
runtime_config(),
pig@agent@state:agent_state(),
session_state(),
pig@agent@msg:agent_msg()
) -> {pig@agent@state:agent_state(),
session_state(),
{ok, pig_protocol@message:message()} | {error, pig@run_error:run_error()}}.
do_loop(Config, Agent_st, Session, M) ->
Transition = pig@agent@update:update(Agent_st, M),
case Transition of
{done, Final_st, Final_message} ->
commit_or_pending(
Config,
Session,
Agent_st,
Final_st,
{return_message, Final_message}
);
{failed, Final_st@1, Inference_error} ->
commit_or_pending(
Config,
Session,
Agent_st,
Final_st@1,
{return_ai_error, Inference_error}
);
{continue, New_st, Effs} ->
case commit_transition(
Session,
Agent_st,
New_st,
resume_from_history
) of
{error, {Pending, Session_error}} ->
{Agent_st, Pending, {error, {session, Session_error}}};
{ok, Next_session} ->
execute_effects(Config, New_st, Next_session, Effs)
end
end.
-file("src/pig/agent/runtime.gleam", 642).
?DOC(
" Decide how to resume based on the last assistant message's stop_reason\n"
" and tool_calls.\n"
).
-spec resume_from_assistant(
runtime_config(),
pig@agent@state:agent_state(),
session_state(),
list(pig_protocol@message:tool_call()),
gleam@option:option(pig_protocol@stop_reason:stop_reason())
) -> {pig@agent@state:agent_state(),
session_state(),
{ok, pig_protocol@message:message()} | {error, pig@run_error:run_error()}}.
resume_from_assistant(Config, St, Session, Tool_calls, Sr) ->
case Sr of
{some, tool_use} ->
{_, Agent_msg} = execute_tools_effect(
Config,
St,
Tool_calls,
fun(Results) -> {tool_results, Results} end
),
do_loop(Config, St, Session, Agent_msg);
{some, stop} ->
Msg@1 = case gleam@list:last(erlang:element(3, St)) of
{ok, Msg} -> Msg;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"pig/agent/runtime"/utf8>>,
function => <<"resume_from_assistant"/utf8>>,
line => 661,
value => _assert_fail,
start => 21027,
'end' => 21069,
pattern_start => 21038,
pattern_end => 21045})
end,
{St, Session, {ok, Msg@1}};
{some, length} ->
{St_after, Provider_msg} = execute_call_provider(
Config,
St,
pig@agent@state:messages_for_provider(St),
pig@agent@state:tool_definitions(St),
fun(R) ->
{provider_responded,
gleam@result:map(
R,
fun(Ir) -> erlang:element(2, Ir) end
)}
end
),
do_loop(Config, St_after, Session, Provider_msg);
{some, error} ->
{St_after, Provider_msg} = execute_call_provider(
Config,
St,
pig@agent@state:messages_for_provider(St),
pig@agent@state:tool_definitions(St),
fun(R) ->
{provider_responded,
gleam@result:map(
R,
fun(Ir) -> erlang:element(2, Ir) end
)}
end
),
do_loop(Config, St_after, Session, Provider_msg);
{some, {unknown, _}} ->
{St_after, Provider_msg} = execute_call_provider(
Config,
St,
pig@agent@state:messages_for_provider(St),
pig@agent@state:tool_definitions(St),
fun(R) ->
{provider_responded,
gleam@result:map(
R,
fun(Ir) -> erlang:element(2, Ir) end
)}
end
),
do_loop(Config, St_after, Session, Provider_msg);
none ->
case Tool_calls of
[] ->
Msg@3 = case gleam@list:last(erlang:element(3, St)) of
{ok, Msg@2} -> Msg@2;
_assert_fail@1 ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"pig/agent/runtime"/utf8>>,
function => <<"resume_from_assistant"/utf8>>,
line => 685,
value => _assert_fail@1,
start => 21844,
'end' => 21886,
pattern_start => 21855,
pattern_end => 21862})
end,
{St, Session, {ok, Msg@3}};
Calls ->
{_, Agent_msg@1} = execute_tools_effect(
Config,
St,
Calls,
fun(Results@1) -> {tool_results, Results@1} end
),
do_loop(Config, St, Session, Agent_msg@1)
end
end.
-file("src/pig/agent/runtime.gleam", 591).
?DOC(
" Resume the agent loop from its current history.\n"
"\n"
" Determines the entry point by inspecting the last message in history.\n"
" This enables the durability pattern: an external system checkpoints\n"
" messages, and on retry, rebuilds history from those checkpoints.\n"
).
-spec resume_from_history(
runtime_config(),
pig@agent@state:agent_state(),
session_state()
) -> {pig@agent@state:agent_state(),
session_state(),
{ok, pig_protocol@message:message()} | {error, pig@run_error:run_error()}}.
resume_from_history(Config, St, Session) ->
case gleam@list:last(erlang:element(3, St)) of
{error, _} ->
{St, Session, {error, {runtime, <<"no history to continue"/utf8>>}}};
{ok, Last_msg} ->
case Last_msg of
{assistant, _, Tool_calls, _, Sr} ->
resume_from_assistant(Config, St, Session, Tool_calls, Sr);
{user, _} ->
{St_after, Provider_msg} = execute_call_provider(
Config,
St,
pig@agent@state:messages_for_provider(St),
pig@agent@state:tool_definitions(St),
fun(R) ->
{provider_responded,
gleam@result:map(
R,
fun(Ir) -> erlang:element(2, Ir) end
)}
end
),
do_loop(Config, St_after, Session, Provider_msg);
{tool, _, _} ->
{St_after, Provider_msg} = execute_call_provider(
Config,
St,
pig@agent@state:messages_for_provider(St),
pig@agent@state:tool_definitions(St),
fun(R) ->
{provider_responded,
gleam@result:map(
R,
fun(Ir) -> erlang:element(2, Ir) end
)}
end
),
do_loop(Config, St_after, Session, Provider_msg);
{system, _} ->
{St,
Session,
{error,
{runtime,
<<"unexpected system message at end of history"/utf8>>}}}
end
end.
-file("src/pig/agent/runtime.gleam", 513).
-spec retry_pending(
runtime_config(),
pig@agent@state:agent_state(),
session_state()
) -> {pig@agent@state:agent_state(),
session_state(),
{ok, pig_protocol@message:message()} | {error, pig@run_error:run_error()}}.
retry_pending(Config, Current, Pending) ->
{Store@1, Commit@1, Candidate@1, Disposition@1} = case Pending of
{session_pending, Store, Commit, Candidate, Disposition} -> {
Store,
Commit,
Candidate,
Disposition};
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"pig/agent/runtime"/utf8>>,
function => <<"retry_pending"/utf8>>,
line => 518,
value => _assert_fail,
start => 16483,
'end' => 16561,
pattern_start => 16494,
pattern_end => 16551})
end,
{session_store, _, Store_commit} = Store@1,
case Store_commit(Commit@1) of
{error, Session_error} ->
{Current,
{session_pending, Store@1, Commit@1, Candidate@1, Disposition@1},
{error, {session, Session_error}}};
{ok, Committed} ->
case accept_committed_session(
Store@1,
Commit@1,
Candidate@1,
Disposition@1,
Committed
) of
{ok, Ready} ->
complete_disposition(
Config,
Candidate@1,
Ready,
Disposition@1
);
{error, {Still_pending, Session_error@1}} ->
{Current,
Still_pending,
{error, {session, Session_error@1}}}
end
end.
-file("src/pig/agent/runtime.gleam", 391).
?DOC(
" Execute the sans-IO loop: call update, interpret effects, feed back.\n"
" Returns the final agent state and the result.\n"
).
-spec execute_loop(
runtime_config(),
pig@agent@state:agent_state(),
session_state(),
pig@agent@msg:agent_msg()
) -> {pig@agent@state:agent_state(),
session_state(),
{ok, pig_protocol@message:message()} | {error, pig@run_error:run_error()}}.
execute_loop(Config, Agent_st, Session, Initial_msg) ->
do_loop(Config, Agent_st, Session, Initial_msg).
-file("src/pig/agent/runtime.gleam", 318).
-spec handle_message(runtime_state(), runtime_msg()) -> gleam@otp@actor:next(runtime_state(), runtime_msg()).
handle_message(St, M) ->
case M of
{run, Prompt, Reply_to} ->
case erlang:element(4, St) of
{session_pending, _, _, _, _} ->
gleam@erlang@process:send(
Reply_to,
{error,
{runtime,
<<"cannot run a new prompt while a session commit is pending"/utf8>>}}
),
gleam@otp@actor:continue(St);
_ ->
Agent_st = {agent_state,
erlang:element(2, erlang:element(2, St)),
erlang:element(3, erlang:element(2, St)),
0},
Result = execute_loop(
erlang:element(3, St),
Agent_st,
erlang:element(4, St),
{user_prompt, Prompt}
),
{Final_state, Final_session, Outcome} = Result,
gleam@erlang@process:send(Reply_to, Outcome),
gleam@otp@actor:continue(
{runtime_state,
Final_state,
erlang:element(3, St),
Final_session}
)
end;
{continue, Reply_to@1} ->
{Final_state@1, Final_session@1, Outcome@1} = case erlang:element(
4,
St
) of
{session_pending, _, _, _, _} = Pending ->
retry_pending(
erlang:element(3, St),
erlang:element(2, St),
Pending
);
_ ->
Agent_st@1 = {agent_state,
erlang:element(2, erlang:element(2, St)),
erlang:element(3, erlang:element(2, St)),
0},
resume_from_history(
erlang:element(3, St),
Agent_st@1,
erlang:element(4, St)
)
end,
gleam@erlang@process:send(Reply_to@1, Outcome@1),
gleam@otp@actor:continue(
{runtime_state,
Final_state@1,
erlang:element(3, St),
Final_session@1}
);
{get_history, Reply_to@2} ->
gleam@erlang@process:send(
Reply_to@2,
erlang:element(3, erlang:element(2, St))
),
gleam@otp@actor:continue(St);
stop ->
gleam@otp@actor:stop()
end.
-file("src/pig/agent/runtime.gleam", 134).
?DOC(
" Start the runtime actor with a pre-built state.\n"
" Used by `pig.gleam` when session replay needs to happen before start.\n"
).
-spec start_with_state(runtime_config(), runtime_state()) -> {ok,
gleam@erlang@process:subject(runtime_msg())} |
{error, gleam@otp@actor:start_error()}.
start_with_state(_, Initial_state) ->
Builder = begin
_pipe = gleam@otp@actor:new(Initial_state),
gleam@otp@actor:on_message(_pipe, fun handle_message/2)
end,
case gleam@otp@actor:start(Builder) of
{ok, Started} ->
{ok, erlang:element(3, Started)};
{error, E} ->
{error, E}
end.
-file("src/pig/agent/runtime.gleam", 299).
-spec runtime_state(
pig@agent@state:agent_config(),
runtime_config(),
list(pig_protocol@message:message()),
session_state()
) -> runtime_state().
runtime_state(Agent_config, Config, Initial_history, Session) ->
{runtime_state,
gleam@list:fold(
Initial_history,
pig@agent@state:new(Agent_config),
fun pig@agent@state:add_message/2
),
Config,
Session}.
-file("src/pig/agent/runtime.gleam", 110).
?DOC(" Start the runtime actor with the given configuration.\n").
-spec start(runtime_config()) -> {ok,
gleam@erlang@process:subject(runtime_msg())} |
{error, gleam@otp@actor:start_error()}.
start(Config) ->
Agent_config = {agent_config,
erlang:element(2, Config),
erlang:element(3, Config),
none,
erlang:element(7, Config),
erlang:element(6, Config),
none,
none,
none,
none,
none,
none},
Initial_state = runtime_state(Agent_config, Config, [], session_disabled),
start_with_state(Config, Initial_state).
-file("src/pig/agent/runtime.gleam", 148).
?DOC(" Send a prompt to the runtime and wait for a response.\n").
-spec run(gleam@erlang@process:subject(runtime_msg()), binary(), integer()) -> {ok,
pig_protocol@message:message()} |
{error, pig@run_error:run_error()}.
run(Subject, Prompt, Timeout) ->
gleam@otp@actor:call(
Subject,
Timeout,
fun(Reply_to) -> {run, Prompt, Reply_to} end
).
-file("src/pig/agent/runtime.gleam", 163).
?DOC(
" Resume the agent loop from its current history.\n"
"\n"
" Looks at the last message in history to determine the entry point:\n"
" - User/Tool message → call the provider\n"
" - Assistant with stop_reason=ToolUse → execute pending tool calls\n"
" - Assistant with stop_reason=Stop → return immediately\n"
" - Assistant with stop_reason=Length/Error → re-call provider\n"
).
-spec run_continue(gleam@erlang@process:subject(runtime_msg()), integer()) -> {ok,
pig_protocol@message:message()} |
{error, pig@run_error:run_error()}.
run_continue(Subject, Timeout) ->
gleam@otp@actor:call(
Subject,
Timeout,
fun(Reply_to) -> {continue, Reply_to} end
).
-file("src/pig/agent/runtime.gleam", 171).
?DOC(" Send a stop message to the runtime actor.\n").
-spec stop(gleam@erlang@process:subject(runtime_msg())) -> nil.
stop(Subject) ->
gleam@otp@actor:send(Subject, stop).
-file("src/pig/agent/runtime.gleam", 176).
?DOC(" Get the agent's current message history.\n").
-spec history(gleam@erlang@process:subject(runtime_msg()), integer()) -> list(pig_protocol@message:message()).
history(Subject, Timeout) ->
gleam@otp@actor:call(
Subject,
Timeout,
fun(Reply_to) -> {get_history, Reply_to} end
).
-file("src/pig/agent/runtime.gleam", 185).
?DOC(
" Send a prompt to the runtime and wait for a response.\n"
" Returns `Error(Nil)` if the call times out or the runtime crashes.\n"
).
-spec try_run(gleam@erlang@process:subject(runtime_msg()), binary(), integer()) -> {ok,
{ok, pig_protocol@message:message()} |
{error, pig@run_error:run_error()}} |
{error, nil}.
try_run(Subject, Prompt, Timeout) ->
pig_agent_try_call_ffi:try_call(
Subject,
Timeout,
fun(Reply_to) -> {run, Prompt, Reply_to} end
).
-file("src/pig/agent/runtime.gleam", 195).
?DOC(
" Resume the agent loop from its current history and wait for a response.\n"
" Returns `Error(Nil)` if the call times out or the runtime crashes.\n"
).
-spec try_run_continue(gleam@erlang@process:subject(runtime_msg()), integer()) -> {ok,
{ok, pig_protocol@message:message()} |
{error, pig@run_error:run_error()}} |
{error, nil}.
try_run_continue(Subject, Timeout) ->
pig_agent_try_call_ffi:try_call(
Subject,
Timeout,
fun(Reply_to) -> {continue, Reply_to} end
).
-file("src/pig/agent/runtime.gleam", 285).
-spec supervised_runtime_config(
pig@agent@state:agent_config(),
gleam@erlang@process:name(pig@obs@dispatcher:dispatcher_message())
) -> runtime_config().
supervised_runtime_config(Agent_config, Dispatcher_name) ->
{runtime_config,
erlang:element(2, Agent_config),
erlang:element(3, Agent_config),
[],
gleam@erlang@process:named_subject(Dispatcher_name),
erlang:element(6, Agent_config),
erlang:element(5, Agent_config)}.
-file("src/pig/agent/runtime.gleam", 265).
-spec start_named_runtime(
pig@agent@state:agent_config(),
gleam@erlang@process:name(pig@obs@dispatcher:dispatcher_message()),
gleam@erlang@process:name(runtime_msg()),
list(pig_protocol@message:message()),
session_state()
) -> {ok, gleam@otp@actor:started(nil)} | {error, gleam@otp@actor:start_error()}.
start_named_runtime(
Agent_config,
Dispatcher_name,
Name,
Initial_history,
Session
) ->
Runtime_config = supervised_runtime_config(Agent_config, Dispatcher_name),
Initial_state = runtime_state(
Agent_config,
Runtime_config,
Initial_history,
Session
),
Builder = begin
_pipe = gleam@otp@actor:new(Initial_state),
_pipe@1 = gleam@otp@actor:on_message(_pipe, fun handle_message/2),
gleam@otp@actor:named(_pipe@1, Name)
end,
case gleam@otp@actor:start(Builder) of
{ok, Started} ->
{ok, {started, erlang:element(2, Started), nil}};
{error, Error} ->
{error, Error}
end.
-file("src/pig/agent/runtime.gleam", 215).
?DOC(
" Create a ChildSpecification for a named runtime actor.\n"
"\n"
" `initial_history` becomes current before the actor starts, and `session`\n"
" carries the durable commit head associated with it. The Subject can be\n"
" recovered after supervisor start with `process.named_subject(name)`.\n"
).
-spec supervised(
pig@agent@state:agent_config(),
gleam@erlang@process:name(pig@obs@dispatcher:dispatcher_message()),
gleam@erlang@process:name(runtime_msg()),
list(pig_protocol@message:message()),
session_state()
) -> gleam@otp@supervision:child_specification(nil).
supervised(Agent_config, Dispatcher_name, Name, Initial_history, Session) ->
gleam@otp@supervision:worker(
fun() ->
start_named_runtime(
Agent_config,
Dispatcher_name,
Name,
Initial_history,
Session
)
end
).
-file("src/pig/agent/runtime.gleam", 237).
?DOC(
" Create a durable ChildSpecification which reloads the session on every start.\n"
"\n"
" A load failure during a later OTP restart is reported as actor initialisation\n"
" failure, allowing the supervisor to apply its usual restart policy.\n"
).
-spec supervised_with_session_store(
pig@agent@state:agent_config(),
gleam@erlang@process:name(pig@obs@dispatcher:dispatcher_message()),
gleam@erlang@process:name(runtime_msg()),
pig@session_store:session_store()
) -> gleam@otp@supervision:child_specification(nil).
supervised_with_session_store(Agent_config, Dispatcher_name, Name, Store) ->
gleam@otp@supervision:worker(
fun() ->
{session_store, Load, _} = Store,
case Load() of
{ok, Loaded} ->
start_named_runtime(
Agent_config,
Dispatcher_name,
Name,
pig@agent@state:strip_system_messages(
erlang:element(3, Loaded)
),
{session_ready, Store, erlang:element(2, Loaded)}
);
{error, Error} ->
logging:log(
error,
<<"Durable session reload failed: "/utf8,
(gleam@string:inspect(Error))/binary>>
),
{error, {init_failed, gleam@string:inspect(Error)}}
end
end
).