// Upstream-request construction and execution for routed model swaps. The
// generation URL, the model-field rewrite, and provider auth back the Stage 3
// same-protocol tee path; `execAttempt` is the buffered unit the chain
// orchestrator and cross-protocol swaps are built from — one upstream call
// with request translation and response→IR parsing.

import type { AgentId } from '../../core/agent.ts'
import { logger } from '../../core/logger.ts'
import type { ModelApi, ProviderAuth, ProviderEntry } from '../../core/model-registry.ts'
import type { DecoderKind } from '../../core/routes.ts'
import { getSecret } from '../../core/secrets.ts'
import { dynamicHeadersFor } from '../proxy/dynamic-headers.ts'
import { filterRequestHeaders } from '../proxy/forward.ts'
import { getSecrets } from '../routing/load.ts'
import {
  adaptForClaudeCodeOAuth,
  adaptForCodexBackend,
  isCodexChatGptBackend,
} from '../translate/codex-compat.ts'
import { type IRResponse, parseResponseToIR, translateRequest } from '../translate/index.ts'
import { getAgentToken } from './oauth/index.ts'
import type { ResolvedMember } from './resolve.ts'

/** Does this path target the API's model-generation endpoint? Guards the
 *  orchestrator from rerouting non-generation calls (e.g. GET /v1/models). */
export function isGenerationPath(api: ModelApi, path: string): boolean {
  if (api === 'anthropic-messages') return path.endsWith('/messages')
  if (api === 'openai-chat') return path.endsWith('/chat/completions')
  return path.endsWith('/responses')
}

/** The canonical generation endpoint for a wire API, appended to a provider
 *  base URL. Provider base URLs follow the OpenAI/Anthropic SDK convention of
 *  ending at `/v1`; the version segment is only added when it is absent.
 *
 *  Exception: codex's ChatGPT session backend (`chatgpt.com/backend-api/codex`)
 *  serves `/responses` directly — no `/v1` segment. Wrapping that base with
 *  `/v1/responses` lands on a 403 HTML login page. Codex's own native flow
 *  hits the right path because it constructs it itself; when an off-agent
 *  route (openclaw, hermes, …) sends through this provider via Thomas, we
 *  rebuild the URL here and must match codex's convention. */
export function generationUrl(baseUrl: string, api: ModelApi): string {
  const base = baseUrl.replace(/\/+$/, '')
  const rest =
    api === 'anthropic-messages'
      ? 'messages'
      : api === 'openai-chat'
        ? 'chat/completions'
        : 'responses'
  if (isCodexChatGptBackend(base)) return `${base}/${rest}`
  return base.endsWith('/v1') ? `${base}/${rest}` : `${base}/v1/${rest}`
}

/** Rewrite the `model` field of a request body. All three wire formats carry
 *  the target model in a top-level `model` string, so one rewrite covers all.
 *  An unparseable body is forwarded unchanged. */
export function rewriteModel(body: ArrayBuffer, modelId: string): ArrayBuffer {
  try {
    const json: unknown = JSON.parse(new TextDecoder().decode(body))
    if (json && typeof json === 'object' && !Array.isArray(json)) {
      ;(json as Record<string, unknown>).model = modelId
      return encodeJson(json)
    }
  } catch {
    // Unparseable body — forward as-is; the model swap simply doesn't apply.
  }
  return body
}

/** True if this provider's base URL points at api.anthropic.com — the
 *  only place where Anthropic's server tools (web_search_20250305,
 *  code_execution_*, computer_*, etc.) actually execute. Used to gate
 *  server-tool stripping: if we're forwarding elsewhere (DeepSeek's
 *  /anthropic endpoint, an OpenAI-compatible provider, …), those tools
 *  have no executor on the other side and the model just emits broken
 *  tool_use calls. */
export function isAnthropicNative(provider: ProviderEntry): boolean {
  try {
    const host = new URL(provider.baseUrl).hostname
    return host === 'api.anthropic.com' || host.endsWith('.anthropic.com')
  } catch {
    return false
  }
}

/** Return a copy of the request body with Anthropic server tools removed
 *  from the top-level `tools` array. Server tools are identified by a
 *  `type` field that isn't the literal `'custom'` (web_search_20250305,
 *  code_execution_20250825, computer_*, bash_*, text_editor_*, …). A
 *  regular custom tool has only `name` + `input_schema` (no `type`, or
 *  `type:'custom'`) and is left untouched. MCP-namespaced tool names
 *  (`mcp__server__tool`) are custom tools too — they pass through.
 *
 *  Returns the same buffer when nothing was stripped (cheap fast path
 *  for the common case where the client sent no server tools). */
export function stripAnthropicServerToolsFromBody(body: ArrayBuffer): ArrayBuffer {
  let parsed: unknown
  try {
    parsed = JSON.parse(new TextDecoder().decode(body))
  } catch {
    return body
  }
  if (!parsed || typeof parsed !== 'object' || Array.isArray(parsed)) return body
  const obj = parsed as Record<string, unknown>
  const filtered = withoutServerTools(obj)
  if (filtered === obj) return body
  return encodeJson(filtered)
}

/** In-memory variant — same filter, applied to a parsed request body.
 *  Returns the same reference when nothing was stripped; otherwise a
 *  shallow copy with `tools` replaced (or removed if empty). */
export function withoutServerTools(body: Record<string, unknown>): Record<string, unknown> {
  const tools = body.tools
  if (!Array.isArray(tools) || tools.length === 0) return body
  const kept = tools.filter((t) => {
    if (typeof t !== 'object' || t === null) return true
    const type = (t as { type?: unknown }).type
    if (typeof type !== 'string') return true
    return type === 'custom'
  })
  if (kept.length === tools.length) return body
  const next: Record<string, unknown> = { ...body }
  if (kept.length === 0) delete next.tools
  else next.tools = kept
  return next
}

/** Apply an auth spec to forwarded request headers. Passthrough auth
 *  leaves the client's own headers intact; otherwise every auth header the
 *  client may have sent is dropped before the configured one is injected.
 *  `agent-oauth` providers inject a subscription token Thomas reads — and
 *  co-refreshes — from the owning agent's own credential store. `id` is
 *  used purely for log labelling (provider id or route key). */
export async function applyAuth(
  headers: Headers,
  auth: ProviderAuth,
  id: string,
): Promise<void> {
  if (auth.kind === 'passthrough') return

  if (auth.kind === 'agent-oauth') {
    try {
      const tok = await getAgentToken(auth.agent)
      stripClientAuth(headers)
      headers.set('authorization', `Bearer ${tok.token}`)
      if (auth.agent === 'codex') {
        if (tok.accountId) headers.set('chatgpt-account-id', tok.accountId)
        // chatgpt.com/backend-api/codex is a Codex-CLI session endpoint —
        // not a generic OpenAI API. It checks the `originator` header for
        // a first-party Codex identity and rejects everything else with
        // a 403 HTML login page. The User-Agent must also claim Codex CLI
        // shape (the `codex_cli_rs/<ver>` prefix is the part that
        // matters; the rest is hand-waved). Setting both lets a non-codex
        // agent (hermes / openclaw) route through the same subscription.
        headers.set('originator', 'codex_cli_rs')
        headers.set('user-agent', CODEX_CLIENT_UA)
      }
      if (auth.agent === 'claude-code') {
        // Anthropic only accepts a Claude.ai subscription (OAuth) token on
        // the Messages API when the request claims Claude Code shape:
        //   • claude-code-20250219 — the Claude-Code-flavored beta gate
        //   • oauth-2025-04-20    — the OAuth-token beta gate
        // Both are required; without claude-code-20250219 the upstream
        // throttles hard (529 / aggressive rate-limit) even when auth itself
        // succeeds. Pattern matches openclaw's anthropic-transport-stream.
        const CC_OAUTH_BETAS = 'claude-code-20250219,oauth-2025-04-20'
        const beta = headers.get('anthropic-beta')
        headers.set(
          'anthropic-beta',
          beta ? `${beta},${CC_OAUTH_BETAS}` : CC_OAUTH_BETAS,
        )
        // Claude Code CLI's User-Agent is checked by the throttle layer.
        // Stamp it consistently so non-claude-code agents routing through
        // a Claude.ai subscription don't hit "unknown client" rate caps.
        headers.set('user-agent', CLAUDE_CODE_CLIENT_UA)
      }
    } catch (err) {
      // Token unavailable (locked Keychain, not logged in) — leave the
      // client's own auth in place and let it through as passthrough
      // rather than sending an unauthenticated request.
      logger.warn(
        `routing: ${auth.agent} OAuth unavailable for ${id}; ` +
          `passing the client's own auth through — ${(err as Error).message}`,
      )
    }
    return
  }

  stripClientAuth(headers)
  const value = getSecret(getSecrets(), auth.valueRef)
  if (!value) {
    logger.warn(`routing: secret "${auth.valueRef}" not set for ${id}`)
    return
  }
  if (auth.kind === 'bearer') {
    headers.set('authorization', `Bearer ${value}`)
  } else {
    headers.set(auth.header, value)
  }
}

/** Drop every auth header the client may have sent — used before injecting
 *  Thomas-managed credentials so the wrong agent's key can't leak upstream. */
function stripClientAuth(headers: Headers): void {
  headers.delete('authorization')
  headers.delete('x-api-key')
  headers.delete('api-key')
}

// User-Agent strings the upstream throttle / gate checks for first-party
// CLI shape. The `<originator>/<version>` prefix is what matters; the
// suffix mirrors codex/claude-code's own format so the request looks
// reasonable in their server logs. Bumping versions is harmless — the
// upstream only checks the prefix and (loosely) the shape.
const CODEX_CLIENT_UA = 'codex_cli_rs/0.81.0 (Thomas; openthomas.com)'
const CLAUDE_CODE_CLIENT_UA = 'claude-cli/2.0.0 (Thomas; openthomas.com)'

function encodeJson(value: unknown): ArrayBuffer {
  const bytes = new TextEncoder().encode(JSON.stringify(value))
  const out = new ArrayBuffer(bytes.byteLength)
  new Uint8Array(out).set(bytes)
  return out
}

/** A wire API → the decoder kind, used for pricing lookup and the captured
 *  packet's `protocol` field. */
export function apiToDecoder(api: ModelApi): DecoderKind {
  if (api === 'anthropic-messages') return 'anthropic'
  if (api === 'openai-chat') return 'openai-chat'
  return 'openai-responses'
}

/** The outcome of one buffered upstream call. */
export type AttemptResult =
  | {
      ok: true
      status: number
      json: unknown
      ir: IRResponse
      durMs: number
      startedAtWall: number
      upstreamUrl: string
    }
  | {
      ok: false
      status: number
      errorText: string
      /** Raw upstream response headers captured on the error path. Useful for
       *  diagnosing 400/401/etc. where the body alone (often empty or
       *  unhelpfully terse) doesn't say what went wrong. */
      errorHeaders?: Record<string, string>
      durMs: number
      startedAtWall: number
      upstreamUrl: string
    }

/** Build the upstream HTTP request for a routed member: translate the client
 *  request into the member's wire format, force the model and the stream flag,
 *  apply provider auth, and target the member API's generation endpoint. */
async function buildUpstreamRequest(
  member: ResolvedMember,
  clientApi: ModelApi,
  clientRequest: Record<string, unknown>,
  ctx: { agent: AgentId; reqHeaders: Headers },
  upstreamUrl: string,
  stream: boolean,
): Promise<Request> {
  // Strip Anthropic server tools before translation when the destination
  // can't execute them. Anthropic's web_search / code_execution / etc.
  // only work when the API call lands on api.anthropic.com; forwarding
  // them to a third-party (even an Anthropic-API-compatible upstream
  // like DeepSeek's /anthropic) lets the model emit tool_use calls that
  // never get results. Filtering at the source keeps the IR clean and
  // prevents the bogus tool from showing up after translation either.
  const sourceRequest = isAnthropicNative(member.provider)
    ? clientRequest
    : withoutServerTools(clientRequest)
  const translated = translateRequest(clientApi, member.api, sourceRequest) as Record<
    string,
    unknown
  >
  const upstreamBody: Record<string, unknown> = {
    ...translated,
    model: member.model.id,
    stream,
  }

  // Claude.ai subscription routes need both the beta flags (set in
  // applyAuth) and an exact-match `system` field — the adapter rewrites
  // the body to match what Claude Code CLI sends.
  if (
    member.api === 'anthropic-messages' &&
    member.provider.auth.kind === 'agent-oauth' &&
    member.provider.auth.agent === 'claude-code'
  ) {
    adaptForClaudeCodeOAuth(upstreamBody, member.provider.reasoningEffort)
  }

  // codex's ChatGPT session endpoint speaks a session-bound subset of the
  // public Responses API; the body adapter mirrors codex-CLI's own request
  // shape (store, reasoning.effort, include, text.verbosity, …).
  if (
    member.api === 'openai-responses' &&
    isCodexChatGptBackend(member.provider.baseUrl)
  ) {
    adaptForCodexBackend(upstreamBody, member.provider.reasoningEffort)
  }

  const headers = filterRequestHeaders(ctx.reqHeaders, new URL(upstreamUrl).host)
  headers.set('content-type', 'application/json')
  headers.set('accept', stream ? 'text/event-stream' : 'application/json')
  if (member.api === 'anthropic-messages' && !headers.has('anthropic-version')) {
    headers.set('anthropic-version', '2023-06-01')
  }
  if (member.provider.auth.kind === 'passthrough') {
    const extra = await dynamicHeadersFor(ctx.agent)
    for (const [k, v] of Object.entries(extra)) headers.set(k, v)
  } else {
    await applyAuth(headers, member.provider.auth, member.provider.id)
  }

  return new Request(upstreamUrl, {
    method: 'POST',
    headers,
    body: JSON.stringify(upstreamBody),
  })
}

/** Run one buffered upstream call against a resolved member: translate the
 *  client request into the member's wire format, force the model and
 *  `stream:false`, apply auth, fetch, and parse the response into IR. Never
 *  throws — a network error or a non-2xx status becomes an `ok:false` result
 *  the orchestrator fails over from. */
export async function execAttempt(
  member: ResolvedMember,
  clientApi: ModelApi,
  clientRequest: Record<string, unknown>,
  ctx: { agent: AgentId; reqHeaders: Headers },
): Promise<AttemptResult> {
  const startedAtWall = Date.now()
  const t0 = performance.now()
  const upstreamUrl = generationUrl(member.provider.baseUrl, member.api)

  try {
    const request = await buildUpstreamRequest(
      member,
      clientApi,
      clientRequest,
      ctx,
      upstreamUrl,
      false,
    )
    const res = await fetch(request)
    const text = await res.text()
    const durMs = performance.now() - t0

    if (res.status >= 400) {
      return {
        ok: false,
        status: res.status,
        errorText: errorExcerpt(text) || `HTTP ${res.status}`,
        errorHeaders: snapshotHeaders(res),
        durMs,
        startedAtWall,
        upstreamUrl,
      }
    }

    let json: unknown
    try {
      json = JSON.parse(text)
    } catch {
      return {
        ok: false,
        status: res.status,
        errorText: `non-JSON response: ${errorExcerpt(text)}`,
        errorHeaders: snapshotHeaders(res),
        durMs,
        startedAtWall,
        upstreamUrl,
      }
    }

    return {
      ok: true,
      status: res.status,
      json,
      ir: parseResponseToIR(member.api, json),
      durMs,
      startedAtWall,
      upstreamUrl,
    }
  } catch (err) {
    return {
      ok: false,
      status: 0,
      errorText: `upstream call failed: ${(err as Error).message}`,
      durMs: performance.now() - t0,
      startedAtWall,
      upstreamUrl,
    }
  }
}

/** The outcome of one streaming upstream call. A success hands back the live
 *  SSE body unread, for the caller to tee; a failure carries the read body. */
export type StreamAttempt =
  | {
      ok: true
      status: number
      body: ReadableStream<Uint8Array>
      durMs: number
      startedAtWall: number
      upstreamUrl: string
    }
  | {
      ok: false
      status: number
      errorText: string
      errorHeaders?: Record<string, string>
      durMs: number
      startedAtWall: number
      upstreamUrl: string
    }

/** Run one streaming upstream call against a resolved member: translate the
 *  client request, force the model and `stream:true`, apply auth, fetch, and
 *  return the live SSE body. Never throws — a network error or a non-2xx
 *  status becomes an `ok:false` result. */
export async function execStream(
  member: ResolvedMember,
  clientApi: ModelApi,
  clientRequest: Record<string, unknown>,
  ctx: { agent: AgentId; reqHeaders: Headers },
): Promise<StreamAttempt> {
  const startedAtWall = Date.now()
  const t0 = performance.now()
  const upstreamUrl = generationUrl(member.provider.baseUrl, member.api)

  try {
    const request = await buildUpstreamRequest(
      member,
      clientApi,
      clientRequest,
      ctx,
      upstreamUrl,
      true,
    )
    const res = await fetch(request)
    const durMs = performance.now() - t0

    if (res.status >= 400 || !res.body) {
      const text = res.body ? await res.text() : ''
      return {
        ok: false,
        status: res.status,
        errorText: errorExcerpt(text) || `HTTP ${res.status}`,
        errorHeaders: snapshotHeaders(res),
        durMs,
        startedAtWall,
        upstreamUrl,
      }
    }

    return { ok: true, status: res.status, body: res.body, durMs, startedAtWall, upstreamUrl }
  } catch (err) {
    return {
      ok: false,
      status: 0,
      errorText: `upstream call failed: ${(err as Error).message}`,
      durMs: performance.now() - t0,
      startedAtWall,
      upstreamUrl,
    }
  }
}

/** Trim an upstream error body to a short single-line excerpt — strips HTML
 *  tags and collapses whitespace, for logs and the captured `error` field. */
function errorExcerpt(text: string): string {
  return text
    .replace(/<[^>]+>/g, ' ')
    .replace(/\s+/g, ' ')
    .trim()
    .slice(0, 200)
}

/** Snapshot an upstream Response's headers into a plain object, for the
 *  diagnostic `errorHeaders` field on a failed attempt. `set-cookie` is
 *  dropped — it would leak the upstream's session state and is never
 *  load-bearing for debugging a 4xx; values longer than 1 KB are truncated
 *  so a misbehaving header can't bloat the captured payload. */
function snapshotHeaders(res: Response): Record<string, string> {
  const out: Record<string, string> = {}
  res.headers.forEach((value, key) => {
    if (key.toLowerCase() === 'set-cookie') return
    out[key] = value.length > 1024 ? `${value.slice(0, 1024)}…` : value
  })
  return out
}
