import type {
  DurableObjectState,
  MessageBatch,
  Queue,
  WebSocket,
} from '@cloudflare/workers-types'
import type { RealtimeEvent } from './index.js'

/** Hoisted per R3 — TextEncoder is expensive native, never reconstruct per call.
 *  verifyInternalSecret runs on the internal-broadcast trust path (per fan-out). */
const utf8 = new TextEncoder()

/** Stamp id + timestamp, fire-and-forget to a CF Queue. Never mutates the caller's object.
 *  Call via ctx.waitUntil(...) in a route handler. */
export async function publishEvent<E extends RealtimeEvent>(
  queue: Queue,
  event: Omit<E, 'id' | 'timestamp'>,
): Promise<void> {
  const stamped = {
    ...event,
    id: crypto.randomUUID(),
    timestamp: new Date().toISOString(),
  } as E
  await queue.send(stamped)
}

/** Hibernatable accept. Supports BOTH accept mechanisms — host picks, module forces neither:
 *  - tags       → ctx.acceptWebSocket(ws, tags)
 *  - attachment → ws.serializeAttachment(attachment) after accept */
export function acceptHibernatable(
  ctx: DurableObjectState,
  ws: WebSocket,
  opts?: { tags?: string[]; attachment?: unknown },
): void {
  if (opts?.tags) {
    ctx.acceptWebSocket(ws, opts.tags)
  } else {
    ctx.acceptWebSocket(ws)
  }
  if (opts?.attachment !== undefined) {
    ws.serializeAttachment(opts.attachment)
  }
}

/** getWebSockets → safe send (ignore closed). targetTag scopes via getWebSockets(targetTag) only. */
export function broadcastFrame(
  ctx: DurableObjectState,
  frame: string,
  opts?: { targetTag?: string; exclude?: WebSocket },
): void {
  const sockets = opts?.targetTag ? ctx.getWebSockets(opts.targetTag) : ctx.getWebSockets()
  for (const ws of sockets) {
    if (opts?.exclude && ws === opts.exclude) continue
    try {
      ws.send(frame)
    } catch {
      // closed socket — ignore
    }
  }
}

/** Drain a queue batch → dispatch each event → ack on success, retry on throw. */
export function makeQueueConsumer<E extends RealtimeEvent>(
  dispatch: (event: E, env: unknown) => Promise<void>,
): (batch: MessageBatch<E>, env: unknown) => Promise<void> {
  return async (batch, env) => {
    for (const msg of batch.messages) {
      try {
        await dispatch(msg.body, env)
        msg.ack()
      } catch {
        msg.retry()
      }
    }
  }
}

/** Constant-time internal-secret check. Length-guard before timingSafeEqual. */
export function verifyInternalSecret(provided: string | null, expected: string): boolean {
  if (provided == null || provided.length !== expected.length) return false
  type CfSubtle = SubtleCrypto & {
    timingSafeEqual(a: ArrayBuffer | ArrayBufferView, b: ArrayBuffer | ArrayBufferView): boolean
  }
  return (crypto.subtle as CfSubtle).timingSafeEqual(
    utf8.encode(provided),
    utf8.encode(expected),
  )
}

/** Participants minus those currently connected (mechanism-agnostic). */
export function computeOffline(
  participantIds: readonly string[],
  online: ReadonlySet<string>,
): string[] {
  return participantIds.filter((id) => !online.has(id))
}
