Packages

A Fair Source multi-agent runtime with deterministic agent scoring and replayable run history.

Current section

Files

Jump to
syntropy lib syntropy event_recorder.ex
Raw

lib/syntropy/event_recorder.ex

defmodule Syntropy.EventRecorder do
@moduledoc """
Retains lattice runtime events and broadcasts them on Phoenix PubSub.
Recording is asynchronous: `record/2` builds the event in the caller,
casts it to the recorder, and returns immediately. The recorder
broadcasts the event on `"lattice:events"` as soon as it is received
and persists events in batches (`Repo.insert_all`) on a short flush
interval, so producers such as the task scheduler never block on
Postgres writes.
The recorder also subscribes to `"lattice:events"` itself: events that
originate on other cluster nodes (PG PubSub fans them out cluster-wide)
are merged into the in-memory history with uid-based deduplication, so
`recent/1` — and therefore `GET /api/events` and the cockpit event
panel — reflect cluster-wide history. Only locally-recorded events are
persisted; each event reaches the shared database exactly once, from
its origin node.
"""
use GenServer
alias Syntropy.{ClusterInfo, Persistence}
@history_limit 100
@flush_interval_ms 100
@flush_batch_size 50
@type event :: %{
id: String.t(),
uid: String.t(),
node_id: String.t(),
node_name: String.t(),
payload: map(),
timestamp: DateTime.t()
}
@spec start_link(keyword()) :: GenServer.on_start()
def start_link(opts \\ []) do
GenServer.start_link(__MODULE__, opts, name: __MODULE__)
end
@spec record(String.t(), map()) :: event()
def record(event_id, payload) do
event = %{
id: event_id,
uid: Ecto.UUID.generate(),
node_id: ClusterInfo.node_id(),
node_name: ClusterInfo.node_name(),
payload: payload,
timestamp: DateTime.utc_now()
}
GenServer.cast(__MODULE__, {:record, event})
event
end
@spec recent(non_neg_integer()) :: [event()]
def recent(limit \\ @history_limit) do
GenServer.call(__MODULE__, {:recent, limit})
end
@doc """
Synchronously persists any batched events that have not been flushed yet.
Intended for tests and shutdown paths that need write durability before
inspecting the database.
"""
@spec flush() :: :ok
def flush do
GenServer.call(__MODULE__, :flush)
end
@spec reset() :: :ok
def reset do
GenServer.call(__MODULE__, :reset)
end
@impl true
def init(_opts) do
Process.flag(:trap_exit, true)
Phoenix.PubSub.subscribe(Syntropy.PubSub, "lattice:events")
{:ok,
%{
events: Persistence.load_recent_events(Persistence.hydrate_limit()),
pending: [],
flush_timer: nil
}}
end
@impl true
def handle_cast({:record, event}, state) do
Phoenix.PubSub.broadcast(Syntropy.PubSub, "lattice:events", {:lattice_event, event})
events = Enum.take([event | state.events], @history_limit)
pending = [event | state.pending]
state = %{state | events: events, pending: pending}
if length(pending) >= @flush_batch_size do
{:noreply, flush_pending(state)}
else
{:noreply, ensure_flush_timer(state)}
end
end
@impl true
def handle_call({:recent, limit}, _from, state) do
{:reply, Enum.take(state.events, max(limit, 0)), state}
end
@impl true
def handle_call(:flush, _from, state) do
{:reply, :ok, flush_pending(state)}
end
@impl true
def handle_call(:reset, _from, state) do
cancel_flush_timer(state)
{:reply, :ok, %{events: [], pending: [], flush_timer: nil}}
end
@impl true
def handle_info(:flush_pending, state) do
{:noreply, flush_pending(%{state | flush_timer: nil})}
end
@impl true
def handle_info({:lattice_event, event}, state) do
if remote_event?(event) do
{:noreply, %{state | events: merge_remote_event(state.events, event)}}
else
# Local events (including our own broadcast echoing back) are already
# in the buffer from handle_cast.
{:noreply, state}
end
end
@impl true
def handle_info(_message, state), do: {:noreply, state}
@impl true
def terminate(_reason, state) do
flush_pending(state)
:ok
end
defp ensure_flush_timer(%{flush_timer: nil} = state) do
%{state | flush_timer: Process.send_after(self(), :flush_pending, @flush_interval_ms)}
end
defp ensure_flush_timer(state), do: state
defp flush_pending(%{pending: []} = state), do: cancel_flush_timer(state)
defp flush_pending(state) do
batch = Enum.reverse(state.pending)
started_at = System.monotonic_time(:millisecond)
batch
|> Persistence.persist_events()
|> Persistence.tolerate_write("runtime_event_batch_insert")
:ok =
Syntropy.Telemetry.event_recorder_flush(
length(batch),
System.monotonic_time(:millisecond) - started_at
)
cancel_flush_timer(%{state | pending: []})
end
defp cancel_flush_timer(%{flush_timer: nil} = state), do: state
defp cancel_flush_timer(%{flush_timer: timer} = state) do
Process.cancel_timer(timer)
%{state | flush_timer: nil}
end
defp remote_event?(event) do
Map.get(event, :node_name, ClusterInfo.node_name()) != ClusterInfo.node_name()
end
defp merge_remote_event(events, event) do
uid = Map.get(event, :uid)
if uid && Enum.any?(events, &(Map.get(&1, :uid) == uid)) do
events
else
[event | events]
|> Enum.sort_by(& &1.timestamp, {:desc, DateTime})
|> Enum.take(@history_limit)
end
end
end