// Process-local pub/sub for captured actions. Consumers (the SSE feed,
// future WebSocket handlers, in-process risk taggers) call subscribe() and
// receive every successful packet append.

import { logger } from '../core/logger.ts'
import type { AgentPacket } from '../core/packet.ts'

type Listener = (packet: AgentPacket) => void

const listeners = new Set<Listener>()

export function subscribeAction(fn: Listener): () => void {
  listeners.add(fn)
  return () => {
    listeners.delete(fn)
  }
}

export function emitAction(packet: AgentPacket): void {
  for (const fn of listeners) {
    try {
      fn(packet)
    } catch (err) {
      logger.warn(`feed listener threw: ${(err as Error).message}`)
    }
  }
}

export function listenerCount(): number {
  return listeners.size
}
