← Files Meetings (Beta)ARCHIVED FILE

scripts/meetings_api_client_transport.py

48.6 KB · Oct 8, 2026 · 12:02 UTC

↓ Download file

#!/usr/bin/env python3
"""Bounded, cookie-free HTTPS transport for direct Meetings API clients."""

from __future__ import annotations

import hmac
import http.client
import json
import math
import random
import re
import socket
import threading
import time
import urllib.error
import urllib.parse
import urllib.request
from collections import deque
from collections.abc import Callable
from concurrent.futures import Future
from concurrent.futures import TimeoutError as FutureTimeoutError
from contextlib import nullcontext
from dataclasses import dataclass
from datetime import datetime, timezone
from email.utils import parsedate_to_datetime
from typing import Protocol, TypeGuard, cast

from codex_auth_client import (
    AuthMaterial,
    ChatGPTAuthProvider,
    CodexAuthAccountUnverified,
    CodexAuthCancelled,
)
from meetings_analytics import parse_analytics_event
from meetings_api_client_common import record_account_fingerprint, record_request_headers
from meetings_connection_migration import MEETINGS_CONNECTION_URL, connection_request_body
from meetings_metrics import measure_notes_lookup_stage

from helpers import create_https_context, is_json, unique_json_object

RECORD_MEETINGS_URL = "https://chatgpt.com/backend-api/meetings/meetings"
RECORD_NOTES_URL = "https://chatgpt.com/backend-api/meetings/notes"
CALENDAR_CONNECTIONS_URL = "https://chatgpt.com/backend-api/aip/connectors/links/list_accessible"
STATSIG_BOOTSTRAP_URL = "https://chatgpt.com/backend-api/wham/statsig/bootstrap"
STATSIG_BOOTSTRAP_BODY = b'{"brand_name":"chatgpt-meetings","window_type":"mcp"}'
MINIMUM_INITIAL_PAGE_LIMIT = 1
MAXIMUM_INITIAL_PAGE_LIMIT = 40
CONTINUATION_PAGE_LIMIT = 15
MAXIMUM_PAGE_RESPONSE_BYTES = 16 * 1024 * 1024
MAXIMUM_CURSOR_BYTES = 8 * 1024
_MAXIMUM_SOCKET_TIMEOUT_SECONDS = 75.0
_HTTP_CANCELLATION_CHECK_SECONDS = 0.025

_CALENDAR_QUERY_TIMESTAMP_PATTERN = re.compile(
    r"^[0-9]{4}-[0-9]{2}-[0-9]{2}T[0-9]{2}:[0-9]{2}:[0-9]{2}\.[0-9]{3}Z$"
)
_MAXIMUM_TIMESTAMP_BYTES = 128
_REQUEST_TIMEOUT_MESSAGE = "Meetings request timed out."


class RecordTransportError(RuntimeError):
    """Base class for bounded Record transport failures."""


class RecordTransportBackendError(RecordTransportError):
    """The Record transport failed or observed an invalid bounded response."""


class RecordTransportTimeout(RecordTransportBackendError):
    """The Record transport exceeded its caller-supplied deadline."""


class RecordTransportCancelled(RecordTransportError):
    """The caller cancelled an in-flight Record transport request."""


class RecordTransportRateLimited(RecordTransportBackendError):
    """The request cannot retry within its remaining attempt or time budget."""

    def __init__(self, retry_after_seconds: float) -> None:
        super().__init__("Meetings is temporarily rate limited.")
        self.retry_after_seconds = retry_after_seconds


@dataclass(frozen=True, slots=True)
class RecordHTTPResult:
    """Bounded HTTP status and response body returned by Record transports."""

    status: int
    body: bytes
    retry_after_seconds: float | None = None


@dataclass(slots=True)
class RecordRateLimitBudget:
    """Quota retries left in one request, including its single auth refresh."""

    remaining_retries: int = 2


@dataclass(slots=True)
class _RateLimitCooldown:
    ready_at: float
    failures: int
    probing: bool = False


class RecordRateLimitRetry:
    """Share transient owner cooldowns without serializing healthy requests."""

    def __init__(
        self,
        *,
        clock: Callable[[], float] = time.monotonic,
        random_fraction: Callable[[], float] = random.random,
        wait: Callable[[float, threading.Event | None], None] | None = None,
    ) -> None:
        self._clock = clock
        self._random_fraction = random_fraction
        self._wait = wait or self._cancelable_wait
        self._lock = threading.Lock()
        self._cooldowns: dict[bytes, _RateLimitCooldown] = {}

    @staticmethod
    def _cancelable_wait(seconds: float, cancellation_event: threading.Event | None) -> None:
        if cancellation_event is None:
            time.sleep(seconds)
        elif cancellation_event.wait(seconds):
            raise RecordTransportCancelled("Meetings request was cancelled.")

    def request(
        self,
        operation: Callable[[float], RecordHTTPResult],
        *,
        owner: bytes,
        method: str,
        timeout_seconds: float,
        cancellation_event: threading.Event | None,
        verify_owner: Callable[[float], None],
        budget: RecordRateLimitBudget | None = None,
    ) -> RecordHTTPResult:
        """Retry only reads or an explicit pre-handler request-quota rejection.

        The existing caller deadline covers the entire retry sequence. One probe
        per throttled owner is released at a time; healthy requests stay concurrent.
        No response bodies, credentials, or resource identifiers enter the cooldown.
        """

        deadline = self._clock() + timeout_seconds
        budget = budget or RecordRateLimitBudget()
        waited = False
        attempt = 0
        while attempt < 3:
            probe: _RateLimitCooldown | None = None
            while True:
                _raise_if_cancelled(cancellation_event)
                now = self._clock()
                remaining = deadline - now
                if remaining <= 0:
                    raise RecordTransportTimeout(_REQUEST_TIMEOUT_MESSAGE)
                with self._lock:
                    cooldown = self._cooldowns.get(owner)
                    delay = max(0.0, cooldown.ready_at - now) if cooldown else 0.0
                    if cooldown is not None and cooldown.probing:
                        delay = max(delay, _HTTP_CANCELLATION_CHECK_SECONDS)
                    elif delay <= 0:
                        if cooldown is not None:
                            cooldown.probing = True
                            probe = cooldown
                        break
                if delay >= remaining:
                    raise RecordTransportRateLimited(delay)
                self._wait(min(delay, _HTTP_CANCELLATION_CHECK_SECONDS), cancellation_event)
                waited = True
            try:
                if waited or attempt:
                    verify_owner(max(0.0, deadline - self._clock()))
                _raise_if_cancelled(cancellation_event)
                with self._lock:
                    # Cached-auth recovery may wait while an in-flight response
                    # extends the shared floor. Reserve that newer cooldown too.
                    if self._cooldowns.get(owner) is not probe:
                        continue
                remaining = deadline - self._clock()
                if remaining <= 0:
                    raise RecordTransportTimeout(_REQUEST_TIMEOUT_MESSAGE)
                attempt += 1
                result = operation(remaining)
                _raise_if_cancelled(cancellation_event)
                if result.status != 429 or (
                    method != "GET" and not _is_request_quota_rejection(result.body)
                ):
                    with self._lock:
                        if self._cooldowns.get(owner) is probe:
                            self._cooldowns.pop(owner, None)
                    return result
                with self._lock:
                    previous = self._cooldowns.get(owner)
                    failures = min(7, previous.failures + 1) if previous else 1
                    fallback = min(60.0, 2.0 ** (failures - 1))
                    delay = max(
                        result.retry_after_seconds or 0.0,
                        fallback * (0.5 + 0.5 * self._random_fraction()),
                    )
                    ready_at = self._clock() + delay
                    if previous is not None:
                        ready_at = max(ready_at, previous.ready_at)
                    # Bounded process-local state; never evict a live owner's floor.
                    if owner not in self._cooldowns and len(self._cooldowns) >= 256:
                        expired = [
                            key
                            for key, item in self._cooldowns.items()
                            if item.ready_at <= self._clock() and not item.probing
                        ]
                        for key in expired:
                            self._cooldowns.pop(key, None)
                        if len(self._cooldowns) >= 256:
                            raise RecordTransportRateLimited(delay)
                    self._cooldowns[owner] = _RateLimitCooldown(ready_at, failures)
                if budget.remaining_retries == 0:
                    raise RecordTransportRateLimited(max(0.0, ready_at - self._clock()))
                budget.remaining_retries -= 1
            finally:
                if probe is not None:
                    with self._lock:
                        probe.probing = False
        raise AssertionError("Rate-limit attempt budget was not enforced")


def _is_request_quota_rejection(body: bytes) -> bool:
    if len(body) > 8192:
        return False
    try:
        payload: object = json.loads(body, object_pairs_hook=unique_json_object)
    except (ValueError, UnicodeError, RecursionError):
        return False
    return is_json(payload) and payload.get("detail") == "record_api_request_quota_exceeded"


PROCESS_RECORD_RATE_LIMIT_RETRY = RecordRateLimitRetry()


class _RequestDeadlineCancellation(threading.Event):
    """Bound this auth waiter without cancelling the shared auth refresh flight."""

    def __init__(self, caller: threading.Event | None, timeout_seconds: float) -> None:
        super().__init__()
        self._caller = caller
        self._deadline = time.monotonic() + timeout_seconds

    def is_set(self) -> bool:
        return (
            super().is_set()
            or (self._caller is not None and self._caller.is_set())
            or time.monotonic() >= self._deadline
        )

    def wait(self, timeout: float | None = None) -> bool:
        deadline = self._deadline
        if timeout is not None:
            deadline = min(deadline, time.monotonic() + timeout)
        while not self.is_set():
            remaining = deadline - time.monotonic()
            if remaining <= 0:
                break
            super().wait(min(remaining, 0.025))
        return self.is_set()


def authenticated_record_request(
    transport: RecordHTTPTransport,
    auth: AuthMaterial,
    auth_client: ChatGPTAuthProvider,
    *,
    url: str,
    timeout_seconds: float,
    cancellation_event: threading.Event | None,
    method: str = "GET",
    body: bytes | None = None,
    statsig_bootstrap: bool = False,
    cleanup: bool = False,
    quota_budget: RecordRateLimitBudget | None = None,
) -> tuple[RecordHTTPResult, AuthMaterial]:
    path = urllib.parse.urlsplit(url).path
    is_notes_page = method == "GET" and path in {
        "/backend-api/meetings/notes",
        "/backend-api/meetings/meetings",
    }
    expected_owner = record_account_fingerprint(auth)
    request_auth = auth

    def verify_owner(remaining: float) -> None:
        nonlocal request_auth
        auth_cancellation = _RequestDeadlineCancellation(cancellation_event, remaining)
        try:
            with measure_notes_lookup_stage("backend_auth") if is_notes_page else nullcontext():
                current = auth_client.get_chatgpt_auth(
                    refresh_token=False, cancellation_event=auth_cancellation
                )
        except CodexAuthCancelled:
            if cancellation_event is not None and cancellation_event.is_set():
                raise
            raise RecordTransportTimeout("Meetings request timed out.") from None
        if not hmac.compare_digest(record_account_fingerprint(current), expected_owner):
            raise CodexAuthAccountUnverified("ChatGPT account changed during the request.")
        request_auth = current

    def operation(remaining: float) -> RecordHTTPResult:
        headers = record_request_headers(request_auth, has_json_body=body is not None)
        if statsig_bootstrap:
            headers["originator"] = "codex_desktop"
        if url in {CALENDAR_CONNECTIONS_URL, MEETINGS_CONNECTION_URL}:
            headers["OAI-Product-Sku"] = "CODEX"
        with measure_notes_lookup_stage("backend_http") if is_notes_page else nullcontext():
            return transport.request(
                url=url,
                headers=headers,
                timeout_seconds=remaining,
                cancellation_event=cancellation_event,
                method=method,
                body=body,
            )

    # Cleanup and preference opt-out stay usable during a read/request cooldown.
    # Native Start/Stop/end/upload recovery has its own recording retry owner.
    if (
        not url.startswith("https://chatgpt.com/backend-api/meetings/")
        or cleanup
        or method == "DELETE"
        or (
            method == "PATCH"
            and path
            in {"/backend-api/meetings/settings", "/backend-api/meetings/settings/preferences"}
        )
    ):
        return operation(timeout_seconds), request_auth
    result = PROCESS_RECORD_RATE_LIMIT_RETRY.request(
        operation,
        owner=expected_owner,
        method=method,
        timeout_seconds=timeout_seconds,
        cancellation_event=cancellation_event,
        verify_owner=verify_owner,
        budget=quota_budget,
    )
    return result, request_auth


class RecordHTTPTransport(Protocol):
    """Structural transport injected into authenticated Record API clients."""

    def request(
        self,
        *,
        url: str,
        headers: dict[str, str],
        timeout_seconds: float,
        cancellation_event: threading.Event | None,
        method: str = "GET",
        body: bytes | None = None,
    ) -> RecordHTTPResult:
        """Perform one bounded Record API request.

        Args:
            url: Absolute allowlisted Record API URL.
            headers: Request headers including current authenticated material.
            timeout_seconds: Remaining request deadline in seconds.
            cancellation_event: Optional caller cancellation signal.
            method: HTTP method for this request.
            body: Optional encoded JSON request body.

        Returns:
            The bounded status and body observed by the transport.
        """

        ...


class _ResponseHeaders(Protocol):
    def get(self, key: str, default: object = None, /) -> object:
        """Return one response-header value."""

        ...


class _HTTPResponse(Protocol):
    @property
    def headers(self) -> _ResponseHeaders:
        """Return the response headers."""

        ...

    def getcode(self) -> int:
        """Return the numeric HTTP status."""

        ...

    def read(self, size: int) -> bytes:
        """Read up to ``size`` response bytes."""

        ...

    def close(self) -> None:
        """Close the response stream."""

        ...


class _HTTPOpener(Protocol):
    def open(
        self,
        request: urllib.request.Request,
        /,
        *,
        timeout: float,
    ) -> _HTTPResponse:
        """Open one request with a bounded socket timeout."""

        ...


class _HTTPAbortTarget(Protocol):
    def close(self) -> None:
        """Interrupt one active connection or response."""

        ...


class _HTTPRequestExecution:
    """Bind one bounded HTTP operation to its current abortable transport."""

    def __init__(self, operation: Callable[[], RecordHTTPResult]) -> None:
        self.operation = operation
        self.future: Future[RecordHTTPResult] = Future()
        self._lock = threading.Lock()
        self._target: _HTTPAbortTarget | None = None
        self._abandoned = False

    def set_abort_target(self, target: _HTTPAbortTarget) -> None:
        with self._lock:
            abandoned = self._abandoned
            if not abandoned:
                self._target = target
        if abandoned:
            try:
                target.close()
            except (OSError, ValueError):
                pass
            raise RecordTransportCancelled("Meetings request was cancelled.")

    def require_active(self) -> None:
        with self._lock:
            abandoned = self._abandoned
        if abandoned:
            raise RecordTransportCancelled("Meetings request was cancelled.")

    def clear_abort_target(self, target: _HTTPAbortTarget) -> None:
        with self._lock:
            if self._target is target:
                self._target = None

    def abandon(self) -> None:
        with self._lock:
            self._abandoned = True
            target = self._target
            self._target = None
        self.future.cancel()
        if target is not None:
            try:
                target.close()
            except (OSError, ValueError):
                pass


class _HTTPExecutionContext(threading.local):
    current: _HTTPRequestExecution | None

    def __init__(self) -> None:
        self.current = None


_HTTP_EXECUTION_CONTEXT = _HTTPExecutionContext()


class _SharedHTTPExecutionPool:
    """Keep blocking DNS, TLS, headers, and body reads off the requesting thread."""

    _MAXIMUM_WORKERS = 4
    _MAXIMUM_PENDING_REQUESTS = 16

    def __init__(self) -> None:
        self._condition = threading.Condition()
        self._pending: deque[_HTTPRequestExecution] = deque()
        self._workers: list[threading.Thread] = []
        self._idle_workers = 0

    def submit(self, operation: Callable[[], RecordHTTPResult]) -> _HTTPRequestExecution:
        execution = _HTTPRequestExecution(operation)
        with self._condition:
            if len(self._pending) >= self._MAXIMUM_PENDING_REQUESTS:
                raise RecordTransportBackendError("Meetings request capacity is unavailable.")
            self._pending.append(execution)
            if (
                self._idle_workers < len(self._pending)
                and len(self._workers) < self._MAXIMUM_WORKERS
            ):
                worker = threading.Thread(
                    target=self._run,
                    name=f"record-meetings-http-{len(self._workers) + 1}",
                    daemon=True,
                )
                self._workers.append(worker)
                worker.start()
            self._condition.notify()
        return execution

    def _run(self) -> None:
        while True:
            with self._condition:
                self._idle_workers += 1
                try:
                    while not self._pending:
                        self._condition.wait()
                    execution = self._pending.popleft()
                finally:
                    self._idle_workers -= 1
            if not execution.future.set_running_or_notify_cancel():
                continue
            _HTTP_EXECUTION_CONTEXT.current = execution
            try:
                result = execution.operation()
            except BaseException as error:
                execution.future.set_exception(error)
            else:
                execution.future.set_result(result)
            finally:
                _HTTP_EXECUTION_CONTEXT.current = None


_PROCESS_HTTP_EXECUTION_POOL = _SharedHTTPExecutionPool()


class _PooledHTTPResponse:
    """Return an exhausted HTTPS connection to the process-local pool."""

    def __init__(
        self,
        response: http.client.HTTPResponse,
        connection: http.client.HTTPSConnection,
        pool: _PooledHTTPSOpener,
    ) -> None:
        self._response = response
        self._connection: http.client.HTTPSConnection | None = connection
        self._pool = pool
        self._exhausted = False
        self._lock = threading.Lock()

    @property
    def headers(self) -> _ResponseHeaders:
        return cast(_ResponseHeaders, self._response.headers)

    def getcode(self) -> int:
        return self._response.getcode()

    def read(self, size: int) -> bytes:
        chunk = self._response.read(size)
        if not chunk:
            self._exhausted = True
        return chunk

    def close(self) -> None:
        with self._lock:
            connection = self._connection
            self._connection = None
        if connection is None:
            return
        reusable = self._exhausted and not self._response.will_close
        self._response.close()
        if reusable:
            self._pool.release(connection)
        else:
            connection.close()


class _PooledHTTPSOpener:
    """Reuse bounded, cookie-free TLS connections for the fixed Record origin."""

    _MAXIMUM_IDLE_CONNECTIONS = 8
    _MAXIMUM_IDLE_SECONDS = 30.0

    def __init__(self) -> None:
        self._idle: list[tuple[float, http.client.HTTPSConnection]] = []
        self._lock = threading.Lock()

    def open(
        self,
        request: urllib.request.Request,
        /,
        *,
        timeout: float,
    ) -> _HTTPResponse:
        parsed = urllib.parse.urlsplit(request.full_url)
        if parsed.hostname is None:
            raise RecordTransportBackendError("Meetings request was invalid.")
        connection: http.client.HTTPSConnection | None = None
        stale_connections: list[http.client.HTTPSConnection] = []
        now = time.monotonic()
        with self._lock:
            while self._idle:
                released_at, candidate = self._idle.pop()
                if now - released_at <= self._MAXIMUM_IDLE_SECONDS:
                    connection = candidate
                    break
                stale_connections.append(candidate)
        for stale_connection in stale_connections:
            stale_connection.close()
        reused_connection = connection is not None
        if connection is None:
            connection = self._new_connection(parsed, timeout=timeout)
        else:
            connection.timeout = timeout
            if connection.sock is not None:
                connection.sock.settimeout(timeout)
        path = parsed.path or "/"
        if parsed.query:
            path += f"?{parsed.query}"
        execution = _HTTP_EXECUTION_CONTEXT.current
        while True:
            if execution is not None:
                execution.set_abort_target(connection)
            try:
                if connection.sock is None:
                    connection.connect()
                    if execution is not None:
                        execution.require_active()
                connection.request(
                    request.get_method(),
                    path,
                    body=request.data,
                    headers=dict(request.header_items()),
                )
                response = connection.getresponse()
            except (OSError, http.client.HTTPException):
                connection.close()
                if execution is not None:
                    execution.clear_abort_target(connection)
                if not reused_connection or request.get_method() != "GET":
                    raise
                # A stale keep-alive can fail after its original server has
                # closed it. Replay only an idempotent read, exactly once.
                reused_connection = False
                connection = self._new_connection(parsed, timeout=timeout)
                continue
            except BaseException:
                connection.close()
                if execution is not None:
                    execution.clear_abort_target(connection)
                raise
            pooled_response = _PooledHTTPResponse(response, connection, self)
            if execution is not None:
                execution.set_abort_target(pooled_response)
            return pooled_response

    @staticmethod
    def _new_connection(
        parsed: urllib.parse.SplitResult, *, timeout: float
    ) -> http.client.HTTPSConnection:
        hostname = parsed.hostname
        if hostname is None:
            raise RecordTransportBackendError("Meetings request was invalid.")
        try:
            context = create_https_context()
        except (OSError, ValueError):
            raise RecordTransportBackendError(
                "Meetings TLS trust could not be configured. Check CODEX_CA_CERTIFICATE or SSL_CERT_FILE."
            ) from None
        proxy = urllib.request.getproxies().get("https")
        if proxy is not None and not urllib.request.proxy_bypass(hostname):
            proxy_address = urllib.parse.urlsplit(proxy)
            if proxy_address.scheme != "http" or proxy_address.hostname is None:
                raise RecordTransportBackendError("Meetings HTTPS proxy is unavailable.")
            connection = http.client.HTTPSConnection(
                proxy_address.hostname,
                port=proxy_address.port or 80,
                timeout=timeout,
                context=context,
            )
            connection.set_tunnel(hostname, port=parsed.port or 443)
            return connection
        return http.client.HTTPSConnection(hostname, timeout=timeout, context=context)

    def release(self, connection: http.client.HTTPSConnection) -> None:
        with self._lock:
            if len(self._idle) < self._MAXIMUM_IDLE_CONNECTIONS:
                self._idle.append((time.monotonic(), connection))
                return
        connection.close()


_PROCESS_HTTPS_OPENER = _PooledHTTPSOpener()


class HTTPSRecordTransport:
    """Process-cached, cookie-free HTTPS transport with a bounded response body."""

    def __init__(self, opener: _HTTPOpener | None = None) -> None:
        # Every authenticated backend client shares one bounded process-local
        # connection pool. Injected openers remain a deterministic test seam.
        self._opener = opener or _PROCESS_HTTPS_OPENER

    def request(
        self,
        *,
        url: str,
        headers: dict[str, str],
        timeout_seconds: float,
        cancellation_event: threading.Event | None,
        method: str = "GET",
        body: bytes | None = None,
    ) -> RecordHTTPResult:
        """Perform one cookie-free HTTPS request.

        Args:
            url: Absolute allowlisted Record API URL.
            headers: Request headers including current authenticated material.
            timeout_seconds: Remaining request deadline in seconds.
            cancellation_event: Optional caller cancellation signal.
            method: HTTP method for this request.
            body: Optional encoded JSON request body.

        Returns:
            The bounded HTTP status and response body.
        """

        if (
            isinstance(timeout_seconds, bool)
            or not math.isfinite(timeout_seconds)
            or timeout_seconds <= 0
        ):
            raise RecordTransportTimeout(_REQUEST_TIMEOUT_MESSAGE)
        if self._opener is _PROCESS_HTTPS_OPENER and (
            not _is_allowed_record_request(method, url, body)
            or not _are_allowed_record_headers(
                headers,
                has_json_body=body is not None,
                allow_statsig_originator=url == STATSIG_BOOTSTRAP_URL,
                allow_calendar_product_sku=url
                in {CALENDAR_CONNECTIONS_URL, MEETINGS_CONNECTION_URL},
            )
        ):
            raise RecordTransportBackendError("Meetings request was invalid.")
        _raise_if_cancelled(cancellation_event)
        deadline = time.monotonic() + timeout_seconds
        execution = _PROCESS_HTTP_EXECUTION_POOL.submit(
            lambda: _inline_https_request(
                self._opener,
                url=url,
                headers=headers,
                timeout_seconds=timeout_seconds,
                cancellation_event=cancellation_event,
                method=method,
                body=body,
                deadline=deadline,
            )
        )
        while True:
            if cancellation_event is not None and cancellation_event.is_set():
                execution.abandon()
                raise RecordTransportCancelled("Meetings request was cancelled.")
            remaining = deadline - time.monotonic()
            if remaining <= 0:
                execution.abandon()
                raise RecordTransportTimeout(_REQUEST_TIMEOUT_MESSAGE)
            try:
                return execution.future.result(
                    timeout=(
                        min(remaining, _HTTP_CANCELLATION_CHECK_SECONDS)
                        if cancellation_event is not None
                        else remaining
                    )
                )
            except FutureTimeoutError:
                continue
            except (OSError, RecordTransportBackendError) as error:
                if cancellation_event is not None and cancellation_event.is_set():
                    raise RecordTransportCancelled("Meetings request was cancelled.") from None
                if time.monotonic() >= deadline:
                    raise RecordTransportTimeout(_REQUEST_TIMEOUT_MESSAGE) from None
                raise error


def _inline_https_request(
    opener: _HTTPOpener,
    *,
    url: str,
    headers: dict[str, str],
    timeout_seconds: float,
    cancellation_event: threading.Event | None,
    method: str = "GET",
    body: bytes | None = None,
    deadline: float | None = None,
) -> RecordHTTPResult:
    """Perform one bounded request using a pooled or explicitly injected opener."""

    _raise_if_cancelled(cancellation_event)
    if deadline is None:
        deadline = time.monotonic() + timeout_seconds
    request = urllib.request.Request(
        url,
        headers=headers,
        data=body,
        method=method,
    )
    response: _HTTPResponse | None = None
    try:
        try:
            response = opener.open(
                request,
                timeout=min(timeout_seconds, _MAXIMUM_SOCKET_TIMEOUT_SECONDS),
            )
        except urllib.error.HTTPError as exc:
            # HTTPError is also the bounded response stream for non-2xx status
            # codes, despite the narrower typeshed protocol for its headers.
            response = cast(_HTTPResponse, exc)
        execution = _HTTP_EXECUTION_CONTEXT.current
        if execution is not None:
            execution.set_abort_target(response)
        _raise_if_cancelled(cancellation_event)
        _raise_if_deadline_expired(deadline)
        status = int(response.getcode())
        declared_length = _content_length(response.headers)
        if declared_length is not None and declared_length > MAXIMUM_PAGE_RESPONSE_BYTES:
            raise RecordTransportBackendError("Meetings returned an oversized response.")
        response_body = _read_bounded_response(
            response,
            cancellation_event,
            deadline=deadline,
        )
        return RecordHTTPResult(
            status=status,
            body=response_body,
            retry_after_seconds=(
                _retry_after_seconds(response.headers) if status in {429, 503} else None
            ),
        )
    except RecordTransportCancelled:
        raise
    except (
        OSError,
        TimeoutError,
        urllib.error.URLError,
        socket.timeout,
        http.client.HTTPException,
    ) as exc:
        _raise_if_cancelled(cancellation_event)
        _raise_if_deadline_expired(deadline)
        raise RecordTransportBackendError("Meetings could not be reached.") from exc
    finally:
        if response is not None:
            try:
                response.close()
            except OSError:
                pass


def _retry_after_seconds(headers: _ResponseHeaders) -> float | None:
    raw = headers.get("Retry-After")
    if not isinstance(raw, str) or len(raw) > 256:
        return None
    value = raw.strip()
    if re.fullmatch(r"[0-9]+(?:\.[0-9]+)?", value):
        seconds = float(value)
        return seconds if math.isfinite(seconds) else None
    try:
        instant = parsedate_to_datetime(value)
        if instant.tzinfo is None:
            return None
        now = time.time()
        response_date = headers.get("Date")
        if isinstance(response_date, str) and len(response_date) <= 256:
            try:
                server_now = parsedate_to_datetime(response_date)
                if server_now.tzinfo is not None:
                    now = server_now.timestamp()
            except (ValueError, TypeError, OverflowError):
                pass
        return max(0.0, instant.timestamp() - now)
    except (ValueError, TypeError, OverflowError):
        return None


def _content_length(headers: _ResponseHeaders | None) -> int | None:
    raw = headers.get("Content-Length") if headers is not None else None
    if isinstance(raw, bool) or not isinstance(raw, (int, str)):
        return None
    try:
        value = int(raw)
    except (TypeError, ValueError):
        return None
    return value if value >= 0 else None


def _read_bounded_response(
    response: _HTTPResponse,
    cancellation_event: threading.Event | None,
    *,
    deadline: float,
) -> bytes:
    body = bytearray()
    while True:
        _raise_if_cancelled(cancellation_event)
        _raise_if_deadline_expired(deadline)
        try:
            chunk = response.read(
                min(
                    64 * 1024,
                    MAXIMUM_PAGE_RESPONSE_BYTES + 1 - len(body),
                )
            )
        except (OSError, ValueError, http.client.HTTPException) as exc:
            _raise_if_cancelled(cancellation_event)
            _raise_if_deadline_expired(deadline)
            raise RecordTransportBackendError("Meetings returned an incomplete response.") from exc
        if not chunk:
            _raise_if_cancelled(cancellation_event)
            _raise_if_deadline_expired(deadline)
            return bytes(body)
        _raise_if_deadline_expired(deadline)
        body.extend(chunk)
        if len(body) > MAXIMUM_PAGE_RESPONSE_BYTES:
            raise RecordTransportBackendError("Meetings returned an oversized response.")


def _raise_if_deadline_expired(deadline: float) -> None:
    if time.monotonic() >= deadline:
        raise RecordTransportTimeout(_REQUEST_TIMEOUT_MESSAGE)


def _raise_if_cancelled(event: threading.Event | None) -> None:
    if event is not None and event.is_set():
        raise RecordTransportCancelled("Meetings request was cancelled.")


_RECORD_INTERACTION_PATH_PATTERN = re.compile(
    r"^/backend-api/meetings/meetings/([A-Za-z0-9_-]{1,512})"
    r"(?:/(summary|transcripts|share/eligibility|share|feedback))?$"
)
_RECORD_NOTE_INTERACTION_PATH_PATTERN = re.compile(
    r"^/backend-api/meetings/notes/([A-Za-z0-9_-]{1,512})"
    r"(?:/(transcripts|share/eligibility|share|feedback|reprocess))?$"
)


def _is_allowed_versioned_statsig_body(body: bytes | None) -> bool:
    if body is None or len(body) > 256:
        return False
    try:
        value = json.loads(body, object_pairs_hook=unique_json_object)
    except (ValueError, UnicodeError):
        return False
    if not is_json(value):
        return False
    if set(value) == {"brand_name", "window_type", "app_version"}:
        allowed_context = value["window_type"] == "mcp"
    elif set(value) == {"brand_name", "window_type", "app_version", "system_name"}:
        allowed_context = (
            isinstance(value["window_type"], str)
            and value["window_type"] in {"web", "desktop", "mobile", "unknown"}
            and isinstance(value["system_name"], str)
            and value["system_name"] in {"Windows", "Darwin", "Linux", "iOS", "Android", "unknown"}
        )
    else:
        return False
    return (
        allowed_context
        and value["brand_name"] == "chatgpt-meetings"
        and isinstance(value["app_version"], str)
        and re.fullmatch(
            r"(?:0|[1-9][0-9]{0,8})\.(?:0|[1-9][0-9]{0,8})\.(?:0|[1-9][0-9]{0,8})"
            r"(?:-(?!0[0-9]+(?:[.+]|$))[0-9A-Za-z-]+"
            r"(?:\.(?!0[0-9]+(?:[.+]|$))[0-9A-Za-z-]+)*)?"
            r"(?:\+[0-9A-Za-z-]+(?:\.[0-9A-Za-z-]+)*)?",
            value["app_version"],
        )
        is not None
    )


def _is_allowed_record_request(
    method: str,
    url: str,
    body: bytes | None,
) -> bool:
    try:
        parsed = urllib.parse.urlsplit(url)
        # Python 3.10 rejects an empty string with strict parsing; no query is
        # valid for the fixed authenticated endpoints below on every runtime.
        query = (
            urllib.parse.parse_qsl(
                parsed.query,
                keep_blank_values=True,
                strict_parsing=True,
                max_num_fields=4,
            )
            if parsed.query
            else []
        )
    except ValueError:
        return False
    if (
        method not in {"GET", "POST", "PATCH", "DELETE"}
        or parsed.scheme != "https"
        or parsed.netloc != "chatgpt.com"
        or parsed.fragment
    ):
        return False
    if parsed.path == "/backend-api/me":
        return url == "https://chatgpt.com/backend-api/me" and method == "GET" and body is None
    if parsed.path == "/backend-api/wham/statsig/bootstrap":
        return (
            url == STATSIG_BOOTSTRAP_URL
            and method == "POST"
            and (body == STATSIG_BOOTSTRAP_BODY or _is_allowed_versioned_statsig_body(body))
        )
    if parsed.path == "/backend-api/meetings/meetings":
        return method == "GET" and body is None and _is_allowed_record_notes_query(query)
    if parsed.path == "/backend-api/meetings/notes":
        return method == "GET" and body is None and _is_allowed_record_notes_query(query)
    if url == CALENDAR_CONNECTIONS_URL:
        return method == "POST" and _decoded_json_object(body) == {
            "principals": [],
            "link_refresh_strategy": "BLOCKING",
        }
    if url == MEETINGS_CONNECTION_URL:
        payload = _decoded_json_object(body)
        return (
            method == "POST"
            and payload is not None
            and payload == connection_request_body()
            and payload.get("ensure_exists") is True
        )
    if parsed.path == "/backend-api/accounts":
        return method == "GET" and body is None and not query
    if parsed.path == "/backend-api/meetings/calendar/events":
        return method == "GET" and body is None and _is_allowed_record_calendar_query(query)
    if parsed.path == "/backend-api/meetings/track":
        payload = _decoded_json_object(body)
        if (
            method != "POST"
            or query
            or payload is None
            or set(payload)
            != {
                "clientEventId",
                "eventName",
                "metadata",
                "pluginVersion",
                "codexAppVersion",
                "operatingSystem",
            }
        ):
            return False
        version = payload["pluginVersion"]
        if (
            not isinstance(version, str)
            or len(version) > 80
            or re.fullmatch(r"[A-Za-z0-9][A-Za-z0-9._+-]*", version) is None
        ):
            return False
        codex_version = payload["codexAppVersion"]
        if (
            not isinstance(codex_version, str)
            or len(codex_version) > 80
            or re.fullmatch(
                r"(?:unknown|[0-9]{1,9}\.[0-9]{1,9}\.[0-9]{1,9}(?:[-+][A-Za-z0-9][A-Za-z0-9.+-]{0,39})?)",
                codex_version,
            )
            is None
            or payload["operatingSystem"] not in ("macos", "windows", "linux", "unknown")
        ):
            return False
        try:
            parse_analytics_event(
                {
                    key: value
                    for key, value in payload.items()
                    if key not in {"pluginVersion", "codexAppVersion", "operatingSystem"}
                }
            )
        except ValueError:
            return False
        return True
    if parsed.path == "/backend-api/meetings/activity":
        return method == "GET" and body is None and not query
    if parsed.path == "/backend-api/meetings/activity/viewed":
        return method == "POST" and body is None and not query
    if parsed.path == "/backend-api/meetings/settings/preferences":
        if query:
            return False
        if method == "GET":
            return body is None
        return method == "PATCH" and _is_allowed_settings_body(body)
    note_path_match = _RECORD_NOTE_INTERACTION_PATH_PATTERN.fullmatch(parsed.path)
    if note_path_match is not None:
        suffix = note_path_match.group(2)
        if query:
            return False
        if suffix == "share":
            return method == "POST" and _is_allowed_share_body(body)
        if suffix == "share/eligibility":
            return method == "POST" and _is_allowed_share_eligibility_body(body)
        if suffix == "feedback":
            return method == "POST" and _is_allowed_feedback_body(body)
        if suffix == "reprocess":
            return method == "POST" and body is None
        if method == "DELETE":
            return suffix is None and body is None
        return method == "GET" and body is None and suffix in {None, "transcripts"}
    path_match = _RECORD_INTERACTION_PATH_PATTERN.fullmatch(parsed.path)
    if path_match is None:
        return False
    suffix = path_match.group(2)
    if query:
        return False
    if suffix in {"share", "share/eligibility"}:
        return method == "POST" and _is_allowed_share_body(body)
    if suffix == "feedback":
        return method == "POST" and _is_allowed_feedback_body(body)
    if method == "DELETE":
        return suffix is None and body is None
    return method == "GET" and body is None and suffix in {None, "summary", "transcripts"}


def _decoded_json_object(body: bytes | None) -> dict[str, object] | None:
    if body is None or not body or len(body) > 16 * 1024:
        return None
    try:
        value: object = json.loads(body.decode("utf-8"))
    except (UnicodeDecodeError, json.JSONDecodeError, RecursionError):
        return None
    return value if is_json(value) else None


def _is_allowed_settings_body(body: bytes | None) -> bool:
    value = _decoded_json_object(body)
    return bool(
        value
        and set(value).issubset(
            {
                "featureEnabled",
                "autoRecordEnabled",
                "slackNotificationsEnabled",
            }
        )
        and all(isinstance(item, bool) for item in value.values())
    )


def _is_allowed_share_eligibility_body(body: bytes | None) -> bool:
    value = _decoded_json_object(body)
    return value == {} or _is_allowed_share_body(body)


def _is_allowed_share_body(body: bytes | None) -> bool:
    value = _decoded_json_object(body)
    if value is None or set(value) != {"email"}:
        return False
    email = value.get("email")
    return (
        isinstance(email, str)
        and 3 <= len(email) <= 320
        and email.count("@") == 1
        and not email.startswith("@")
        and not email.endswith("@")
        and not any(character.isspace() for character in email)
        and not _contains_header_control_character(email)
    )


_FEEDBACK_CATEGORIES = frozenset(
    {
        "good_bot",
        "bad_bot",
        "other",
        "bot_did_not_join",
        "bot_disconnected",
        "recording_missing",
        "transcript_issue",
        "summary_issue",
        "calendar_schedule_issue",
    }
)


def _is_allowed_feedback_body(body: bytes | None) -> bool:
    value = _decoded_json_object(body)
    if value is None or not {"rating", "includeDebugArtifacts"}.issubset(value):
        return False
    if set(value) - {
        "rating",
        "feedbackText",
        "includeDebugArtifacts",
        "feedbackContinuationToken",
    }:
        return False
    text = value.get("feedbackText")
    continuation_token = value.get("feedbackContinuationToken")
    return (
        value.get("rating") in _FEEDBACK_CATEGORIES
        and isinstance(value.get("includeDebugArtifacts"), bool)
        and (text is None or isinstance(text, str) and len(text) <= 2_000)
        and (
            continuation_token is None
            or isinstance(continuation_token, str)
            and 1 <= len(continuation_token) <= 2_048
            and not _contains_header_control_character(continuation_token)
            and not any(character.isspace() for character in continuation_token)
        )
    )


def _is_allowed_record_notes_query(query: list[tuple[str, str]]) -> bool:
    if len(query) not in {1, 2, 3} or query[0][0] != "limit":
        return False
    try:
        limit = int(query[0][1])
    except ValueError:
        return False
    if str(limit) != query[0][1]:
        return False
    remaining = query[1:]
    if remaining[:1] == [("include_unsuccessful", "true")]:
        remaining = remaining[1:]
    if remaining and remaining[0][0] == "recording_started_at":
        if len(remaining) != 1 or _canonical_calendar_query_datetime(remaining[0][1]) is None:
            return False
        remaining = remaining[1:]
    if not remaining:
        return MINIMUM_INITIAL_PAGE_LIMIT <= limit <= MAXIMUM_INITIAL_PAGE_LIMIT
    if (
        len(remaining) != 1
        or limit != CONTINUATION_PAGE_LIMIT
        or remaining[0][0] != "cursor"
        or not remaining[0][1]
        or _contains_header_control_character(remaining[0][1])
    ):
        return False
    try:
        return len(remaining[0][1].encode("utf-8")) <= MAXIMUM_CURSOR_BYTES
    except UnicodeEncodeError:
        return False


def _is_allowed_record_calendar_query(query: list[tuple[str, str]]) -> bool:
    if len(query) not in {3, 4}:
        return False
    if [key for key, _value in query[:3]] != ["time_min", "time_max", "limit"]:
        return False
    if query[2][1] != "200":
        return False
    time_min = _canonical_calendar_query_datetime(query[0][1])
    time_max = _canonical_calendar_query_datetime(query[1][1])
    if time_min is None or time_max is None:
        return False
    local_start = time_min.astimezone()
    local_end = time_max.astimezone()
    if (
        local_start.timetz().replace(tzinfo=None) != datetime.min.time()
        or local_end.timetz().replace(tzinfo=None) != datetime.min.time()
        or (local_end.date() - local_start.date()).days != 1
    ):
        return False
    if len(query) == 3:
        return True
    option_key, option_value = query[3]
    if option_key == "refresh":
        return option_value == "true"
    if (
        option_key != "cursor"
        or not option_value
        or _contains_header_control_character(option_value)
    ):
        return False
    try:
        return len(option_value.encode("utf-8")) <= MAXIMUM_CURSOR_BYTES
    except UnicodeEncodeError:
        return False


def _canonical_calendar_query_datetime(value: str) -> datetime | None:
    try:
        encoded = value.encode("utf-8")
    except UnicodeEncodeError:
        return None
    if (
        len(encoded) > _MAXIMUM_TIMESTAMP_BYTES
        or _CALENDAR_QUERY_TIMESTAMP_PATTERN.fullmatch(value) is None
    ):
        return None
    try:
        parsed = datetime.strptime(value, "%Y-%m-%dT%H:%M:%S.%fZ")
    except ValueError:
        return None
    return parsed.replace(tzinfo=timezone.utc)


def _are_allowed_record_headers(
    value: object,
    *,
    has_json_body: bool = False,
    allow_statsig_originator: bool = False,
    allow_calendar_product_sku: bool = False,
) -> TypeGuard[dict[str, str]]:
    expected_headers = {
        "Accept",
        "Authorization",
        "Cache-Control",
        "ChatGPT-Account-Id",
        "User-Agent",
    }
    if has_json_body:
        expected_headers.add("Content-Type")
    if allow_statsig_originator:
        expected_headers.add("originator")
    if allow_calendar_product_sku:
        expected_headers.add("OAI-Product-Sku")
    if not is_json(value) or set(value) != expected_headers:
        return False
    if not all(isinstance(item, str) for item in value.values()):
        return False
    authorization = value["Authorization"]
    account_id = value["ChatGPT-Account-Id"]
    if not isinstance(authorization, str) or not isinstance(account_id, str):
        return False
    try:
        authorization_bytes = authorization.encode("ascii")
        account_id_bytes = account_id.encode("utf-8")
    except UnicodeEncodeError:
        return False
    return (
        value["Accept"] == "application/json"
        and value["Cache-Control"] == "no-store"
        and value["User-Agent"] == "ChatGPT Meetings/1.0"
        and (not has_json_body or value.get("Content-Type") == "application/json")
        and (not allow_statsig_originator or value.get("originator") == "codex_desktop")
        and (not allow_calendar_product_sku or value.get("OAI-Product-Sku") == "CODEX")
        and authorization.startswith("Bearer ")
        and len(authorization_bytes) > len("Bearer ")
        and len(authorization_bytes) <= 64 * 1024 + len("Bearer ")
        and 0 < len(account_id_bytes) <= 256
        and not _contains_header_control_character(authorization)
        and not _contains_header_control_character(account_id)
    )


def _contains_header_control_character(value: str) -> bool:
    return any(ord(character) < 0x20 or ord(character) == 0x7F for character in value)


__all__ = [
    "CONTINUATION_PAGE_LIMIT",
    "HTTPSRecordTransport",
    "MAXIMUM_CURSOR_BYTES",
    "MAXIMUM_INITIAL_PAGE_LIMIT",
    "MAXIMUM_PAGE_RESPONSE_BYTES",
    "MINIMUM_INITIAL_PAGE_LIMIT",
    "RECORD_MEETINGS_URL",
    "RECORD_NOTES_URL",
    "RecordHTTPResult",
    "RecordHTTPTransport",
    "RecordTransportBackendError",
    "RecordTransportCancelled",
    "RecordTransportError",
    "RecordTransportTimeout",
]

SHA-256: 1a062fe0fe97252430e6770c18e9cd29465f9b1ae146b1d35d8db4f98e2a5388