///|
/// Agent-owned wakeup capability, ledger, and pump (plan step 3).
///
/// The ledger is synchronous and self-contained (safe to touch from task
/// bodies, observers, and hooks): per-(session, tag) pending tickets with
/// latest-requirement generations, a FIFO execution order, and a cond-var
/// notification primitive — new requests, operation settled, and shutdown
/// all broadcast, so the pump never polls. Execution reuses
/// `run_turn_with_mode`: a wakeup is an ordinary turn through the
/// single-operation lifecycle, nothing else. The pump itself is not a public
/// Background Task and produces no outcome.

///|
/// Envelope key cap: tags beyond this length are rejected at admission.
let wakeup_tag_max_length : Int = 80

///|
priv struct WakeupPending {
  ticket : @port.WakeupTicket
  extension_id : String
  mut note : String
  /// Latest requirement generation seen for this key (port §4): refreshed
  /// by every request, so a snapshot taken before the latest request can
  /// never satisfy the entry.
  mut latest_requirement : Int
  /// Activation generation this ticket was admitted under: the binding
  /// generation of the window that owned it. A later binding change makes
  /// the ticket stale forever, even if the same session is bound again.
  activation : Int
}

///|
priv struct AgentWakeupRuntime {
  /// Request admission is open only inside a live `run_scoped` and before
  /// shutdown begins; once closed it never reopens.
  mut scope_active : Bool
  mut closed : Bool
  /// The host-bound active window session. `None` until the host binds one
  /// and cleared again at shutdown; requests may only target it.
  mut active_session : String?
  /// Bumped on every binding change. Pending tickets carry the generation
  /// they were admitted under; the pump rejects mismatched ones, so a
  /// dropped window's tickets can never execute after a switch back.
  mut active_generation : Int
  mut next_ticket_seq : Int
  mut next_operation_seq : Int
  /// Pump notification primitive. Signalled on: new request, operation
  /// settled, shutdown.
  wake : @async.CondVar
  /// Monotonic count of those state changes. Waiters compare epochs instead
  /// of trusting broadcast timing, so a notification landing between a
  /// condition check and `wait` can never be lost.
  mut epoch : Int
  /// Composed observers for the wakeup lifecycle projection. Available at
  /// composition, before `on_compose` delivers the capability.
  observers : Array[&@port.Observer]
  /// Pending tickets in FIFO execution order.
  order : Array[WakeupPending]
  /// session → tag → pending entry (coalescing lookup).
  by_key : Map[String, Map[String, WakeupPending]]
  /// Per-session latest requirement generation, bumped on every request.
  generations : Map[String, Int]
  /// Armed input cutoff of the one ordinary turn currently between its
  /// input snapshot and its committed run start. Keyed by that run's id;
  /// commit and unwind both match the key, so cleanup can never remove
  /// another turn's arm.
  mut cutoff_arm : WakeupCutoffArm?
}

///|
fn AgentWakeupRuntime::AgentWakeupRuntime(
  observers~ : Array[&@port.Observer],
) -> AgentWakeupRuntime {
  {
    scope_active: false,
    closed: false,
    active_session: None,
    active_generation: 0,
    next_ticket_seq: 0,
    next_operation_seq: 0,
    epoch: 0,
    wake: @async.CondVar::Cond(),
    observers,
    order: [],
    by_key: Map::from_array([]),
    generations: Map::from_array([]),
    cutoff_arm: None,
  }
}

///|
/// Bind one extension's capability view: requests carry the extension id for
/// event attribution. The runtime is fully built before `on_compose`, so the
/// capability never exposes an uninitialized object.
fn AgentWakeupRuntime::capability(
  self : AgentWakeupRuntime,
  extension_id : String,
) -> @port.WakeupPort {
  let runtime = self
  @port.WakeupPort::from_request(request=fn(session, tag, note) {
    runtime.request(extension_id, session, tag, note)
  })
}

///|
fn AgentWakeupRuntime::activate(
  self : AgentWakeupRuntime,
) -> Unit raise @error.AgentError {
  if self.closed {
    raise @error.AgentError::Runtime(
      @error.RuntimeError::InvocationFailed("wakeup runtime is closed"),
    )
  }
  if self.scope_active {
    raise @error.AgentError::Runtime(
      @error.RuntimeError::InvocationFailed("wakeup scope is already active"),
    )
  }
  self.scope_active = true
}

///|
/// Host-driven active-window binding (synchronous, usable while idle).
/// `Some(id)` binds, `None` clears. Reaffirming the bound id is a no-op:
/// valid tickets keep their place. A real change retires the previous
/// window — its pending tickets get exactly one `dropped` terminal and are
/// never revived by binding the same session again. An already admitted
/// operation is untouched: it settles through its own lifecycle.
fn AgentWakeupRuntime::set_active_session(
  self : AgentWakeupRuntime,
  session : String?,
) -> Unit {
  if self.closed {
    return
  }
  match (self.active_session, session) {
    (Some(current), Some(next)) if current == next => return
    _ => ()
  }
  self.active_session = session
  self.active_generation = self.active_generation + 1
  let stale : Array[WakeupPending] = []
  let mut index = 0
  while index < self.order.length() {
    let entry = self.order[index]
    match session {
      Some(active) if entry.ticket.session == active => index = index + 1
      _ => {
        stale.push(entry)
        ignore(self.order.remove(index))
      }
    }
  }
  for entry in stale {
    self.forget(entry)
  }
  for entry in stale {
    self.emit_wakeup_event("wakeup_dropped", entry)
  }
}

///|
/// A taken-but-not-yet-admitted ticket is stale once its window is gone:
/// its admission generation no longer matches the binding, or its session
/// is no longer the bound one.
fn AgentWakeupRuntime::pending_stale(
  self : AgentWakeupRuntime,
  entry : WakeupPending,
) -> Bool {
  match self.active_session {
    Some(active) => {
      let stale_generation = entry.activation != self.active_generation
      stale_generation || entry.ticket.session != active
    }
    None => true
  }
}

///|
fn AgentWakeupRuntime::emit_wakeup_event(
  self : AgentWakeupRuntime,
  label : String,
  entry : WakeupPending,
  operation_id? : String? = None,
  envelope? : String? = None,
) -> Unit {
  let ticket = entry.ticket
  let fields : Array[(String, Json)] = [
    ("ticket_id", Json::string(ticket.id)),
    ("session", Json::string(ticket.session)),
    ("tag", Json::string(ticket.tag)),
    ("extension_id", Json::string(entry.extension_id)),
    ("ticket_enqueued_at", Json::number(ticket.enqueued_at.to_double())),
  ]
  match operation_id {
    Some(id) => fields.push(("operation_id", Json::string(id)))
    None => ()
  }
  match envelope {
    Some(text) => fields.push(("envelope", Json::string(text)))
    None => ()
  }
  emit_custom_event(self.observers, None, "posoco.wakeup", label, fields)
}

///|
/// Synchronous admission. Validates scope, session, and tag; coalesces on
/// (session, tag); every request records the latest requirement generation.
/// A pending entry returns the same ticket (the `wakeup_requested` event
/// fired once, at creation); during execution the key is free again, so a
/// same-key request forms the next pending ticket at a new generation.
fn AgentWakeupRuntime::request(
  self : AgentWakeupRuntime,
  extension_id : String,
  session : String,
  tag : String,
  note : String,
) -> Result[@port.WakeupTicket, @port.WakeupError] {
  if self.closed || !self.scope_active {
    return Err(ScopeInactive)
  }
  if session == "" {
    return Err(InvalidSession(reason="session must not be empty"))
  }
  if tag == "" {
    return Err(TagInvalid(reason="tag must not be empty"))
  }
  if tag.length() > wakeup_tag_max_length {
    return Err(
      TagInvalid(
        reason="tag must not exceed \{wakeup_tag_max_length} characters",
      ),
    )
  }
  // A request never binds the active session itself; only the currently
  // bound window session may be targeted, and only while one is bound.
  match self.active_session {
    Some(active) if active == session => ()
    Some(_) =>
      return Err(
        InvalidSession(reason="session does not match the active session"),
      )
    None => return Err(InvalidSession(reason="no active session is bound"))
  }
  let generation = match self.generations.get(session) {
    Some(value) => value + 1
    None => 1
  }
  self.generations[session] = generation
  let session_map = match self.by_key.get(session) {
    Some(existing) => existing
    None => {
      let created : Map[String, WakeupPending] = Map::from_array([])
      self.by_key[session] = created
      created
    }
  }
  match session_map.get(tag) {
    Some(entry) => {
      // Late requirement on the same key: keep the ticket (requested fired
      // once), refresh its note, and move the entry to the latest
      // generation, so a snapshot taken before this request can never
      // satisfy it.
      entry.note = note
      entry.latest_requirement = generation
      Ok(entry.ticket)
    }
    None => {
      self.next_ticket_seq = self.next_ticket_seq + 1
      let ticket : @port.WakeupTicket = {
        id: "wakeup_ticket_\{self.next_ticket_seq}",
        session,
        tag,
        enqueued_at: @async.now(),
      }
      let entry : WakeupPending = {
        ticket,
        extension_id,
        note,
        latest_requirement: generation,
        activation: self.active_generation,
      }
      session_map[tag] = entry
      self.order.push(entry)
      self.emit_wakeup_event("wakeup_requested", entry)
      self.notify_state_change()
      Ok(ticket)
    }
  }
}

///|
/// Commit the armed input cutoff for `session_id`: every still-pending
/// entry whose latest requirement generation is at or before `cutoff` is
/// covered by the qualifying ordinary input that took the snapshot (the
/// satisfaction means the scheduling demand was met by that input
/// opportunity — nothing about outcomes). Entries created or refreshed
/// after the snapshot survive to a later ordinary turn.
fn AgentWakeupRuntime::satisfy_covered(
  self : AgentWakeupRuntime,
  session_id : String,
  cutoff : Int,
  operation_id : String,
) -> Unit {
  let covered : Array[WakeupPending] = []
  let mut index = 0
  while index < self.order.length() {
    let entry = self.order[index]
    if entry.ticket.session == session_id && entry.latest_requirement <= cutoff {
      covered.push(entry)
      ignore(self.order.remove(index))
    } else {
      index = index + 1
    }
  }
  for entry in covered {
    self.forget(entry)
  }
  for entry in covered {
    self.emit_wakeup_event(
      "wakeup_satisfied",
      entry,
      operation_id=Some(operation_id),
    )
  }
}

///|
fn AgentWakeupRuntime::forget(
  self : AgentWakeupRuntime,
  entry : WakeupPending,
) -> Unit {
  match self.by_key.get(entry.ticket.session) {
    Some(session_map) => {
      session_map.remove(entry.ticket.tag)
      if session_map.is_empty() {
        self.by_key.remove(entry.ticket.session)
      }
    }
    None => ()
  }
}

///|
/// Take the oldest pending ticket when the Agent is idle. The check and the
/// subsequent admission happen with no await in between, so a wakeup turn
/// never steals the guard from an operation that already ran its admission.
fn AgentWakeupRuntime::take_next(
  self : AgentWakeupRuntime,
  agent : AgentRuntime,
) -> WakeupPending? {
  for ;; {
    if self.closed || agent.turn_active {
      return None
    }
    match self.order.get(0) {
      Some(entry) => {
        ignore(self.order.remove(0))
        self.forget(entry)
        if self.pending_stale(entry) {
          // The window that owned this ticket is gone: its terminal fired
          // at the switch (this is the defensive second layer). Skip it and
          // keep scanning so a fresher ticket behind it stays reachable.
          continue
        }
        return Some(entry)
      }
      None => return None
    }
  }
}

///|
/// Start the pump after `run_scoped`'s startup lifecycle completed, so the
/// first lifecycle `on_start` never races a pump turn.
fn[G] AgentWakeupRuntime::start_pump(
  self : AgentWakeupRuntime,
  group : @async.TaskGroup[G],
  agent : AgentRuntime,
) -> Unit {
  group.spawn_bg(async fn() { self.pump_loop(agent) })
}

///|
/// Record one pump-relevant state change: advance the epoch, then broadcast.
/// The epoch makes the notification durable for a waiter whose condition
/// check already ran — it re-checks the epoch, not the broadcast timing.
fn AgentWakeupRuntime::notify_state_change(self : AgentWakeupRuntime) -> Unit {
  self.epoch = self.epoch + 1
  self.wake.broadcast()
}

///|
/// Wait for a pump-relevant state change at or after `seen`. Returns
/// immediately when the epoch already moved past `seen`: the notification
/// raced the caller's condition check (e.g. the incumbent operation settled
/// between the Busy rejection and this wait) and must not be slept through.
async fn AgentWakeupRuntime::wait_for_change(
  self : AgentWakeupRuntime,
  seen : Int,
) -> Unit {
  if self.epoch == seen {
    self.wake.wait()
  }
}

///|
/// Wake only on notifications: new request, operation settled, shutdown.
/// `seen` is captured before the condition is read, so every state change
/// that could invalidate the decision to sleep is observable as an epoch
/// advance — the waits never spin and never miss a notification.
async fn AgentWakeupRuntime::pump_loop(
  self : AgentWakeupRuntime,
  agent : AgentRuntime,
) -> Unit {
  for ;; {
    let seen = self.epoch
    match self.take_next(agent) {
      Some(entry) => {
        if self.execute_ticket(agent, entry) {
          // Re-queued behind a live incumbent: wait on the state-change
          // notification (its settled broadcast) — never a no-await
          // hotloop. A settle that already raced the Busy return advanced
          // the epoch past `seen`, so this returns instead of sleeping.
          self.wait_for_change(seen)
        }
        continue
      }
      None => ()
    }
    if self.closed {
      return
    }
    self.wait_for_change(seen)
  }
}

///|
/// Bounded, injection-safe envelope header field: the tag is untrusted
/// extension input, so it is length-capped and newline-flattened, then
/// JSON string-escaped via `Json::stringify` — the envelope's `"` / `[` /
/// `]` delimiters can no longer appear unescaped inside the value, so the
/// header closes exactly once no matter what the tag contains.
fn wakeup_envelope_field(value : String, cap : Int) -> String {
  Json::string(bounded_error_label(value, cap))
  .stringify()
  .replace(old="[", new="\\u005b")
  .replace(old="]", new="\\u005d")
}

///|
/// The envelope head all wakeup shapes share (transcript input and the
/// `wakeup_executing` display envelope).
fn wakeup_envelope_head(tag : String) -> String {
  "[posoco-wakeup tag=\{wakeup_envelope_field(tag, wakeup_tag_max_length)}]"
}

///|
/// The wakeup turn's input is a user message of the task-outcome envelope
/// family: `[posoco-wakeup tag=] `, with the trailing note
/// omitted when empty. The envelope belongs to the session record; the
/// wakeup turn is audited exactly like a host turn.
fn wakeup_envelope_message(tag : String, note : String) -> @kernel.Message {
  let head = wakeup_envelope_head(tag)
  if note == "" {
    @kernel.Message::UserMessage(content=[@kernel.Content::Text(head)])
  } else {
    @kernel.Message::UserMessage(content=[
      @kernel.Content::Text("\{head} \{note}"),
    ])
  }
}

///|
async fn AgentWakeupRuntime::execute_ticket(
  self : AgentWakeupRuntime,
  agent : AgentRuntime,
  entry : WakeupPending,
) -> Bool {
  if self.pending_stale(entry) {
    // Taken before a binding change, never admitted: no operation exists,
    // so rejecting here cannot disturb any in-flight user turn or guard.
    self.emit_wakeup_event("wakeup_dropped", entry)
    return false
  }
  let ticket = entry.ticket
  self.next_operation_seq = self.next_operation_seq + 1
  let operation_id = "agent_wakeup_\{self.next_operation_seq}"
  let head = wakeup_envelope_head(ticket.tag)
  let envelope_text = if entry.note == "" {
    head
  } else {
    "\{head} \{entry.note}"
  }
  // Exactly one terminal per ticket. A raw pump cancellation bypasses the
  // catch below, so this unwind emitter still fires once — after the
  // operation settled cancelled on its own unwind — and the cancellation
  // keeps propagating raw, never translated into an ordinary fallback.
  let terminal : Ref[String?] = Ref(None)
  errdefer (match terminal.val {
    Some(_) => ()
    None =>
      self.emit_wakeup_event(
        "wakeup_cancelled",
        entry,
        operation_id=Some(operation_id),
      )
  })
  let result : Result[@types.TurnResult, @error.AgentError] = Ok(
    agent.run_turn_with_mode(
      wakeup_envelope_message(ticket.tag, entry.note),
      ticket.session,
      false,
      origin=Wakeup,
      operation_label=Some(operation_id),
      on_admitted=Some(fn(admission) {
        ignore(admission)
        self.emit_wakeup_event(
          "wakeup_executing",
          entry,
          operation_id=Some(operation_id),
          envelope=Some(envelope_text),
        )
      }),
    ),
  ) catch {
    error => Err(error)
  }
  let label = match result {
    Ok(_) => "wakeup_completed"
    // Defensive: admission raced a host operation. Re-queue at the head and
    // wait for its settled notification instead of failing the ticket —
    // busy is a queueing signal, not a terminal state; neither a terminal
    // nor an executing event fired.
    Err(@error.AgentError::Busy(..)) => {
      self.requeue_front(entry)
      return true
    }
    Err(@error.AgentError::Cancelled(_)) => "wakeup_cancelled"
    Err(_) => "wakeup_failed"
  }
  terminal.val = Some(label)
  self.emit_wakeup_event(label, entry, operation_id=Some(operation_id))
  false
}

///|
fn AgentWakeupRuntime::requeue_front(
  self : AgentWakeupRuntime,
  entry : WakeupPending,
) -> Unit {
  self.order.insert(0, entry)
  let session_map = match self.by_key.get(entry.ticket.session) {
    Some(existing) => existing
    None => {
      let created : Map[String, WakeupPending] = Map::from_array([])
      self.by_key[entry.ticket.session] = created
      created
    }
  }
  session_map[entry.ticket.tag] = entry
}

///|
fn AgentWakeupRuntime::notify_operation_settled(
  self : AgentWakeupRuntime,
) -> Unit {
  self.notify_state_change()
}

///|
/// Shutdown step 1: close request admission (never reopens) and wake the
/// pump so it stops taking tickets. A ticket already executing finishes its
/// ordinary turn lifecycle on its own.
fn AgentWakeupRuntime::begin_shutdown(self : AgentWakeupRuntime) -> Unit {
  if !self.closed {
    self.closed = true
    self.scope_active = false
    // The window is gone: clear the binding. Pending tickets still receive
    // their `dropped` terminal in the later shutdown step, so the clear here
    // must not emit them a second time.
    self.active_session = None
    self.notify_state_change()
  }
}

///|
/// Shutdown step 2: pending scheduling requests die with the scope. Each
/// pending ticket gets exactly one `dropped` terminal; a ticket already
/// executing is not touched here — its turn settles it.
fn AgentWakeupRuntime::drop_pending(self : AgentWakeupRuntime) -> Unit {
  let pending = self.order.copy()
  self.order.clear()
  for entry in pending {
    self.forget(entry)
    self.emit_wakeup_event("wakeup_dropped", entry)
  }
}