///|
/// 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~),
        )
      }
    ToolCallsPreRejected(calls~, outcomes~, ..) =>
      // Pre-pass rejects are reducer-transcript facts (NotExecuted tool
      // messages); mirror them exactly like folded completions.
      if !self.dirty.contains(session_id) {
        for i in 0.. ()
  }
}

///|
/// 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)
  // One outcome shape → one terminal tool event. The variant IS the node:
  // success, failure and rejection never share an event.
  let event : @types.TurnEvent = match outcome {
    Success(content~, structured~) =>
      ToolCallSucceeded(call~, content~, structured~, attachments=[])
    SuccessWithAttachments(content~, structured~, attachments~) =>
      ToolCallSucceeded(call~, content~, structured~, attachments~)
    ToolReportedError(content~, structured~) =>
      ToolCallFailed(call~, failure=ReportedError(content~, structured~))
    RuntimeFailure(error_category~, message~) =>
      ToolCallFailed(call~, failure=Runtime(category=error_category, message~))
    NotExecuted(reason~, ..) => ToolCallRejected(call~, reason~)
  }
  self.emit(Some(scope), event)
}

///|
/// Abandon every still-pending call. Used when a run ends without folding
/// their completions (cancel, terminal hook failure, suspension) and
/// defensively when a new run starts over leftovers. With `scope` the drain
/// attributes to the run that owned the calls; `None` marks out-of-run
/// cleanup.
fn AgentEventSubscriber::drain_pending(
  self : AgentEventSubscriber,
  scope : @types.EventScope?,
  reason : String,
) -> Unit {
  if self.pending_calls.is_empty() {
    return
  }
  for _key, call in self.pending_calls {
    self.emit(scope, ToolCallAbandoned(call~, reason~))
  }
  self.pending_calls.clear()
}

///|
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.drain_pending(None, "run restarted before completion")
    RunSuspended(..) => self.drain_pending(Some(scope), "run suspended")
    RunCompleted(..) | RunFailed(..) | RunCancelled(..) =>
      // Terminal-only drain: cancelled tool effects still receive their
      // closing ToolCompleted events between CancelRequested and the
      // terminal event, so draining any earlier would break pairing.
      self.drain_pending(Some(scope), "run terminal")
    ToolCallsPreRejected(calls~, outcomes~, ..) =>
      // Pre-pass rejects never entered the pipeline: emit the rejection
      // directly, no pending marker, nothing to drain.
      for i in 0..
            self.emit(Some(scope), ToolCallRejected(call=calls[i], reason~))
          _ =>
            raise ContractViolation(
              subscriber_id="agent_committed_event_projection",
              detail="ToolCallsPreRejected carries a non-NotExecuted outcome",
            )
        }
      }
    ToolCallApproved(model_step_id=_, call~, consent_scope~) =>
      self.emit(Some(scope), ToolCallApproved(call~, consent_scope~))
    ToolWaveStarted(calls~, ..) =>
      for call in calls {
        self.emit(Some(scope), ToolCallStarted(call=snapshot_tool_call(call)))
      }
    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)
}