///|
/// Agent-facing projection of committed Puppet envelopes onto observers,
/// plus the per-turn stream-chunk callback wiring.
///|
/// Correctness-critical persistence projection for committed model/tool facts.
/// Agent checkpoints the admitted input before Puppet starts; this subscriber
/// then awaits SessionStore checkpoints after each committed assistant/tool
/// boundary. A transcript rewrite marks the session dirty and pauses
/// incremental projection until the terminal full save re-establishes a
/// canonical prefix.
priv struct SessionCheckpointSubscriber {
stores : Array[&@port.SessionStore]
cursors : Map[String, Int]
dirty : Map[String, Bool]
}
///|
fn SessionCheckpointSubscriber::SessionCheckpointSubscriber(
stores : Array[&@port.SessionStore],
cursors : Map[String, Int],
dirty : Map[String, Bool],
) -> SessionCheckpointSubscriber {
{ stores, cursors, dirty, }
}
///|
impl @puppetry.EventSubscriber for SessionCheckpointSubscriber with fn subscriber_id(
_self : SessionCheckpointSubscriber,
) -> String {
"agent_session_checkpoint"
}
///|
impl @puppetry.EventSubscriber for SessionCheckpointSubscriber with fn provenance(
_self : SessionCheckpointSubscriber,
) -> String {
"posoco.agent.session"
}
///|
async fn SessionCheckpointSubscriber::append_committed_message(
self : SessionCheckpointSubscriber,
session_id : String,
message : @kernel.Message,
) -> Unit raise @puppetry.SubscriberError {
if self.stores.is_empty() {
return
}
let cursor = match self.cursors.get(session_id) {
Some(value) => value
None => 0
}
let results : Array[Result[Unit, @error.SessionError]] = @async.all(
self.stores.map(fn(store) {
() => {
let checkpoint : @port.SessionCheckpoint = {
from_index: cursor,
messages: [snapshot_message(message)],
metadata: None,
}
Ok(store.checkpoint(session_id, checkpoint)) catch {
error => Err(error)
}
}
}),
) catch {
error =>
raise @puppetry.SubscriberError::Internal(
subscriber_id="agent_session_checkpoint",
detail="checkpoint dispatch failed: " + error.to_string(),
)
}
for result in results {
match result {
Ok(_) => ()
Err(error) =>
raise @puppetry.SubscriberError::Internal(
subscriber_id="agent_session_checkpoint",
detail="session_id=" +
safe_error_label(session_id) +
"; " +
error.to_string(),
)
}
}
self.cursors[session_id] = cursor + 1
}
///|
impl @puppetry.EventSubscriber for SessionCheckpointSubscriber with fn on_event(
self : SessionCheckpointSubscriber,
envelope : @puppetry.EventEnvelope,
) -> Unit raise @puppetry.SubscriberError {
let session_id = envelope.session_id().to_string()
match envelope.event() {
RunStarted(..) => self.dirty.remove(session_id)
TranscriptRewritten(..) => self.dirty[session_id] = true
ModelStepCompleted(completion~, ..) =>
if !self.dirty.contains(session_id) {
let payload = completion.message
self.append_committed_message(
session_id,
AssistantMessage(
content=payload.content,
tool_calls=payload.tool_calls,
reasoning=payload.reasoning,
finish_reason=payload.finish_reason,
),
)
}
ToolCompleted(call~, outcome~) =>
if !self.dirty.contains(session_id) {
self.append_committed_message(
session_id,
ToolMessage(call_id=call.call_id, tool_name=call.name, outcome~),
)
}
_ => ()
}
}
///|
/// Agent-facing projection of committed Puppet envelopes. It never scans a
/// final transcript. The pending call cache is correlation state populated by
/// `ToolBatchStarted` and consumed by completion events. A committed
/// Pending events retain the model's original call; completion events carry
/// the canonical call that reached the host.
priv struct AgentEventSubscriber {
observers : Array[&@port.Observer]
pending_calls : Map[String, @kernel.ToolCall]
chunk_dispatcher : Ref[ChunkDispatcher?]
/// Late-bound live context-state lookup (session id -> puppet state),
/// wired after the runtime exists in `AgentRuntime::compose`. `None` (or
/// a lookup returning `None`) projects nothing.
context_state_of : Ref[((String) -> @types.ContextState?)?]
}
///|
fn AgentEventSubscriber::AgentEventSubscriber(
observers : Array[&@port.Observer],
chunk_dispatcher : Ref[ChunkDispatcher?],
context_state_of? : Ref[((String) -> @types.ContextState?)?] = Ref(None),
) -> AgentEventSubscriber {
{
observers,
pending_calls: Map::from_array([]),
chunk_dispatcher,
context_state_of,
}
}
///|
impl @puppetry.EventSubscriber for AgentEventSubscriber with fn subscriber_id(
_self : AgentEventSubscriber,
) -> String {
"agent_committed_event_projection"
}
///|
impl @puppetry.EventSubscriber for AgentEventSubscriber with fn provenance(
_self : AgentEventSubscriber,
) -> String {
"posoco.agent"
}
///|
fn AgentEventSubscriber::emit(
self : AgentEventSubscriber,
scope : @types.EventScope?,
event : @types.TurnEvent,
) -> Unit {
// Piggyback flush: catch up buffered chunk telemetry before every committed
// event so observers see chunks before ModelResponseReceived, ToolCallPending,
// etc., and so they catch up at boundaries during an unbroken SSE burst.
match self.chunk_dispatcher.val {
Some(dispatcher) => dispatcher.flush()
None => ()
}
for observer in self.observers {
observer.on_event_at(scope, snapshot_turn_event(event))
}
}
///|
fn AgentEventSubscriber::emit_tool_result(
self : AgentEventSubscriber,
scope : @types.EventScope,
call : @kernel.ToolCall,
outcome : @kernel.ToolOutcome,
) -> Unit raise @puppetry.SubscriberError {
let key = call.call_id.to_string()
if !self.pending_calls.contains(key) {
raise ContractViolation(
subscriber_id="agent_committed_event_projection",
detail="ToolCompleted without committed ToolBatchStarted for call " + key,
)
}
self.pending_calls.remove(key)
self.emit(Some(scope), @types.TurnEvent::tool_call_result(call, outcome))
}
///|
impl @puppetry.EventSubscriber for AgentEventSubscriber with fn on_event(
self : AgentEventSubscriber,
envelope : @puppetry.EventEnvelope,
) -> Unit raise @puppetry.SubscriberError {
let event = envelope.event()
// Envelope-projected events always carry the envelope's run identity as
// their attribution scope.
let scope : @types.EventScope = {
session_id: envelope.session_id(),
run_id: envelope.run_id(),
turn_id: envelope.turn_id(),
}
match event {
RunStarted(..) => self.pending_calls.clear()
ModelStepCompleted(completion~, ..) => {
let payload = completion.message
let message : @kernel.Message = AssistantMessage(
content=payload.content,
tool_calls=payload.tool_calls,
reasoning=payload.reasoning,
finish_reason=payload.finish_reason,
)
self.emit(
Some(scope),
ModelResponseReceived(message~, usage=completion.usage),
)
// The puppet recorded this step's usage into the session's context
// state before the boundary committed, so the lookup is fresh.
// Projecting per model step keeps the ctx indicator at the same
// cadence as cache — every round, not only at turn boundaries.
match self.context_state_of.val {
Some(lookup) =>
match lookup(envelope.session_id().to_string()) {
Some(state) => self.emit(Some(scope), ContextStateUpdated(state~))
None => ()
}
None => ()
}
}
ToolBatchStarted(calls~, ..) =>
for call in calls {
let snapshot = snapshot_tool_call(call)
self.pending_calls[snapshot.call_id.to_string()] = snapshot
self.emit(Some(scope), ToolCallPending(snapshot))
}
ToolCompleted(call~, outcome~) =>
self.emit_tool_result(scope, call, outcome)
_ => ()
}
}
///|
/// Wrap the puppet's `HostChunkCallback` so each `StreamChunk` is re-emitted
/// to observers as `StreamChunkReceived`. The chunk is already the canonical
/// typed shape; no JSON decoding is needed here.
fn create_agent_stream_callback(
observers : Array[&@port.Observer],
dispatcher_ref : Ref[ChunkDispatcher?],
) -> @kernel_exec.HostChunkCallback? {
if observers.is_empty() {
return None
}
let cb = fn(chunk : @types.StreamChunk) {
// Route to the active turn's dispatcher. If no turn is active (should not
// happen for chunks produced by the Puppet), drop silently.
match dispatcher_ref.val {
Some(dispatcher) => dispatcher.enqueue(chunk)
None => ()
}
}
Some(cb)
}