← Files Meetings (Beta)ARCHIVED FILE
scripts/control_transport/live_client.py
66.4 KB · Oct 9, 2026 · 12:23 UTC
"""Portable, owner-bound JSON-RPC transport for native Meetings control.
Canonical: Microwave chatgpt-meetings-core/e2e/canonical_mcp_transport/live_client.py
Vendored: Microwave plugins/chatgpt-meetings-dev/scripts/live_client.py
Vendored: OpenAI chatgpt/oai-maintained-plugins/plugins/chatgpt-meetings/mcp/control_transport/live_client.py
Sync: python3 chatgpt-meetings-core/scripts/generate_recording_control_stream_contract.py --openai-root /path/to/openai
Check: python3 chatgpt-meetings-core/scripts/generate_recording_control_stream_contract.py --openai-root /path/to/openai --check
This canonical module intentionally depends only on the Python standard library.
Business contracts, generated schemas, owner provenance, and status caches are
validated through caller-supplied callbacks and never imported here.
"""
from __future__ import annotations
import copy
import ctypes
import errno
import io
import json
import os
import select
import socket
import stat
import struct
import sys
import threading
import time
from collections import OrderedDict, deque
from collections.abc import Callable
from dataclasses import dataclass, field
from enum import Enum
from pathlib import Path
from typing import TYPE_CHECKING, Generic, Literal, Protocol, TypeVar
if TYPE_CHECKING:
from typing import TypeGuard
_DEFAULT_MAXIMUM_FRAME_BYTES = 1_048_576
_DEFAULT_MAXIMUM_PENDING_REQUESTS = 16
_DEFAULT_MAXIMUM_CACHED_RESPONSES = 128
_CANCELLATION_CHECK_SECONDS = 0.1
_MAXIMUM_READ_CHUNK_BYTES = 64 * 1024
_WINDOWS_NAMED_PIPE_PREFIX = "\\\\.\\pipe\\"
_WINDOWS_PIPE_INPUT_POLL_SECONDS = 0.05
_WINDOWS_THREAD_TERMINATE = 0x0001
_WINDOWS_RETRYABLE_PIPE_OPEN_ERRNOS = frozenset(
{errno.EBUSY, errno.EINVAL, errno.ENOENT, errno.EPIPE}
)
_WINDOWS_RETRYABLE_PIPE_OPEN_ERRORS = frozenset({2, 231, 232, 233})
_LOCAL_SOCKET_PROTOCOL = 0
_LOCAL_PEER_CREDENTIALS = 0x001
_LOCAL_PEER_PROCESS_ID = 0x002
_MAXIMUM_PEER_CREDENTIAL_BYTES = 128
_WindowsOperationResult = TypeVar("_WindowsOperationResult")
class LiveTransportFailureKind(Enum):
"""Recovery meaning supplied by the failing operation, independent of its message."""
UNAVAILABLE = "unavailable"
TIMED_OUT = "timed-out"
CANCELLED = "cancelled"
REJECTED = "rejected"
class LiveTransportFailureReason(Enum):
"""Concrete operation origin; consumers own any recovery decision."""
ENDPOINT_REJECTED = "endpoint-rejected"
ENDPOINT_MISSING = "endpoint-missing"
class LiveTransportError(RuntimeError):
"""One bounded local transport, peer identity, or JSON-RPC failure."""
def __init__(
self,
kind: LiveTransportFailureKind,
message: str,
*,
io_attempted: bool = False,
diagnostic_reason: LiveTransportFailureReason | None = None,
) -> None:
super().__init__(message)
self.kind = kind
self.diagnostic_reason = diagnostic_reason
# Local validation/queue deadlines are not evidence of an unresponsive peer.
# This describes the failed operation, not permission to replay or replace it.
self.io_attempted = io_attempted
def _io_failure_kind(error: OSError | ValueError) -> LiveTransportFailureKind:
winerror = getattr(error, "winerror", None)
if isinstance(error, TimeoutError) or winerror == 121: # ERROR_SEM_TIMEOUT
return LiveTransportFailureKind.TIMED_OUT
# Windows' exact pipe state takes precedence over its coarse CRT errno.
if winerror in _WINDOWS_RETRYABLE_PIPE_OPEN_ERRORS or winerror == 109:
return LiveTransportFailureKind.UNAVAILABLE
if (
isinstance(error, (PermissionError, NotADirectoryError)) or winerror == 5
): # ERROR_ACCESS_DENIED
return LiveTransportFailureKind.REJECTED
return LiveTransportFailureKind.UNAVAILABLE
def _run_callback(callback: Callable[[], None], *, failure_message: str) -> None:
try:
callback()
except LiveTransportError:
raise
except Exception as exc:
raise LiveTransportError(LiveTransportFailureKind.REJECTED, failure_message) from exc
@dataclass(frozen=True, repr=False)
class LiveOwner:
"""The exact independently discovered native process and OS endpoint."""
endpoint: str
pid: int
transport: Literal["unix-domain-socket", "windows-named-pipe"]
class _WindowsCrtApi(Protocol):
def get_osfhandle(self, descriptor: int) -> int: ...
class _WindowsCtypesApi(Protocol):
def WinDLL(self, name: str, *, use_last_error: bool) -> ctypes.CDLL: ...
def get_last_error(self) -> int: ...
if TYPE_CHECKING:
windows_crt: _WindowsCrtApi
_windows_ctypes: _WindowsCtypesApi
elif os.name == "nt":
import msvcrt as windows_crt
_windows_ctypes = ctypes
def windows_kernel32() -> ctypes.CDLL:
"""Return the Windows process and named-pipe identity API."""
return _windows_ctypes.WinDLL("kernel32", use_last_error=True)
def windows_advapi32() -> ctypes.CDLL:
"""Return the Windows process-token and SID comparison API."""
return _windows_ctypes.WinDLL("advapi32", use_last_error=True)
class _WindowsSidAndAttributes(ctypes.Structure):
_fields_ = [("sid", ctypes.c_void_p), ("attributes", ctypes.c_uint32)]
class _WindowsTokenUser(ctypes.Structure):
_fields_ = [("user", _WindowsSidAndAttributes)]
def windows_process_has_same_user(kernel: ctypes.CDLL, owner_pid: int) -> bool:
"""Compare the reviewed native process and this process's TokenUser SIDs."""
advapi = windows_advapi32()
open_process = kernel.OpenProcess
open_process.argtypes = [ctypes.c_uint32, ctypes.c_int, ctypes.c_uint32]
open_process.restype = ctypes.c_void_p
current_process = kernel.GetCurrentProcess
current_process.argtypes = []
current_process.restype = ctypes.c_void_p
close_handle = kernel.CloseHandle
close_handle.argtypes = [ctypes.c_void_p]
close_handle.restype = ctypes.c_int
open_token = advapi.OpenProcessToken
open_token.argtypes = [
ctypes.c_void_p,
ctypes.c_uint32,
ctypes.POINTER(ctypes.c_void_p),
]
open_token.restype = ctypes.c_int
token_information = advapi.GetTokenInformation
token_information.argtypes = [
ctypes.c_void_p,
ctypes.c_int,
ctypes.c_void_p,
ctypes.c_uint32,
ctypes.POINTER(ctypes.c_uint32),
]
token_information.restype = ctypes.c_int
equal_sid = advapi.EqualSid
equal_sid.argtypes = [ctypes.c_void_p, ctypes.c_void_p]
equal_sid.restype = ctypes.c_int
process = open_process(0x1000, 0, owner_pid)
if not process:
return False
owner_token = ctypes.c_void_p()
client_token = ctypes.c_void_p()
try:
if not open_token(process, 0x0008, ctypes.byref(owner_token)) or not open_token(
current_process(),
0x0008,
ctypes.byref(client_token),
):
return False
owner_size = ctypes.c_uint32()
client_size = ctypes.c_uint32()
token_information(owner_token, 1, None, 0, ctypes.byref(owner_size))
token_information(client_token, 1, None, 0, ctypes.byref(client_size))
if not owner_size.value or not client_size.value:
return False
owner_buffer = ctypes.create_string_buffer(owner_size.value)
client_buffer = ctypes.create_string_buffer(client_size.value)
if not token_information(
owner_token,
1,
owner_buffer,
owner_size.value,
ctypes.byref(owner_size),
) or not token_information(
client_token,
1,
client_buffer,
client_size.value,
ctypes.byref(client_size),
):
return False
owner_sid = _WindowsTokenUser.from_buffer(owner_buffer).user.sid
client_sid = _WindowsTokenUser.from_buffer(client_buffer).user.sid
return bool(
owner_sid
and client_sid
and equal_sid(ctypes.c_void_p(owner_sid), ctypes.c_void_p(client_sid))
)
finally:
if owner_token.value is not None:
close_handle(owner_token)
if client_token.value is not None:
close_handle(client_token)
close_handle(process)
@dataclass(repr=False)
class _WindowsOperation(Generic[_WindowsOperationResult]):
started: threading.Event = field(default_factory=threading.Event)
completed: threading.Event = field(default_factory=threading.Event)
abandoned: threading.Event = field(default_factory=threading.Event)
publication_lock: threading.Lock = field(default_factory=threading.Lock)
result: _WindowsOperationResult | None = None
failure: BaseException | None = None
def cancel_windows_thread(kernel: ctypes.CDLL, thread_id: int | None) -> None:
if thread_id is None:
return
try:
open_thread = kernel.OpenThread
open_thread.argtypes = [ctypes.c_uint32, ctypes.c_int, ctypes.c_uint32]
open_thread.restype = ctypes.c_void_p
cancel = kernel.CancelSynchronousIo
cancel.argtypes = [ctypes.c_void_p]
cancel.restype = ctypes.c_int
close_handle = kernel.CloseHandle
close_handle.argtypes = [ctypes.c_void_p]
close_handle.restype = ctypes.c_int
handle = open_thread(_WINDOWS_THREAD_TERMINATE, 0, thread_id)
if handle:
try:
cancel(handle)
finally:
close_handle(handle)
except (AttributeError, OSError, ValueError):
return
def _raise_if_cancelled(cancellation_event: threading.Event | None) -> None:
if cancellation_event is not None and cancellation_event.is_set():
raise LiveTransportError(
LiveTransportFailureKind.CANCELLED,
"native control stream request was cancelled",
)
def _remaining_seconds(deadline: float, *, io_attempted: bool = False) -> float:
remaining = deadline - time.monotonic()
if remaining <= 0:
raise LiveTransportError(
LiveTransportFailureKind.TIMED_OUT,
"native control stream request timed out",
io_attempted=io_attempted,
)
return remaining
def run_bounded_windows_operation(
operation: Callable[[], _WindowsOperationResult],
*,
kernel: ctypes.CDLL,
deadline: float,
cancellation_event: threading.Event | None,
) -> _WindowsOperationResult:
result: _WindowsOperation[_WindowsOperationResult] = _WindowsOperation()
def perform() -> None:
try:
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline)
result.started.set()
value = operation()
with result.publication_lock:
abandoned = result.abandoned.is_set()
if not abandoned:
result.result = value
if abandoned and isinstance(value, io.IOBase):
value.close()
except BaseException as exc:
result.failure = exc
finally:
result.completed.set()
worker = threading.Thread(
target=perform,
daemon=True,
name="meetings-windows-control-io",
)
worker.start()
try:
while True:
_raise_if_cancelled(cancellation_event)
remaining = _remaining_seconds(deadline, io_attempted=result.started.is_set())
wait = (
remaining
if cancellation_event is None
else min(remaining, _CANCELLATION_CHECK_SECONDS)
)
if result.completed.wait(wait):
break
except LiveTransportError:
with result.publication_lock:
result.abandoned.set()
abandoned_value = result.result
result.result = None
if isinstance(abandoned_value, io.IOBase):
abandoned_value.close()
cancel_windows_thread(kernel, worker.native_id)
raise
if result.failure is not None:
if isinstance(result.failure, LiveTransportError):
raise result.failure
kind = (
_io_failure_kind(result.failure)
if isinstance(result.failure, OSError)
else LiveTransportFailureKind.REJECTED
)
raise LiveTransportError(
kind,
"native Windows control stream I/O failed",
io_attempted=isinstance(result.failure, OSError)
and result.failure.errno != errno.EBADF,
) from result.failure
if result.result is None:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native Windows control stream I/O failed",
)
return result.result
def _retryable_windows_pipe_open_error(exc: OSError) -> bool:
"""Recognize the CRT translations for an advisory named-pipe wait race."""
winerror = getattr(exc, "winerror", None)
return (
winerror in _WINDOWS_RETRYABLE_PIPE_OPEN_ERRORS
or winerror is None
and exc.errno in _WINDOWS_RETRYABLE_PIPE_OPEN_ERRNOS
)
def _open_bounded_windows_named_pipe(
endpoint: str,
*,
kernel: ctypes.CDLL,
deadline: float,
cancellation_event: threading.Event | None,
) -> io.FileIO:
"""Open after an advisory wait, retrying only another client's claim race."""
wait_for_pipe = kernel.WaitNamedPipeW
wait_for_pipe.argtypes = [ctypes.c_wchar_p, ctypes.c_uint32]
wait_for_pipe.restype = ctypes.c_int
def wait_until_available(wait_millis: int) -> bool:
if wait_for_pipe(endpoint, wait_millis):
return True
# Last-error is thread-local: inspect it inside the same bounded worker.
error_code = _windows_ctypes.get_last_error()
if error_code == 121: # ERROR_SEM_TIMEOUT
kind = LiveTransportFailureKind.TIMED_OUT
elif error_code in _WINDOWS_RETRYABLE_PIPE_OPEN_ERRORS or error_code == 109:
kind = LiveTransportFailureKind.UNAVAILABLE
else:
kind = LiveTransportFailureKind.REJECTED
raise LiveTransportError(
kind,
"native Windows control stream is unavailable",
io_attempted=True,
diagnostic_reason=(
LiveTransportFailureReason.ENDPOINT_REJECTED if error_code == 5 else None
),
)
while True:
wait_millis = max(1, int(_remaining_seconds(deadline) * 1000))
run_bounded_windows_operation(
lambda wait_millis=wait_millis: wait_until_available(wait_millis),
kernel=kernel,
deadline=deadline,
cancellation_event=cancellation_event,
)
try:
return run_bounded_windows_operation(
lambda: io.FileIO(endpoint, mode="r+"),
kernel=kernel,
deadline=deadline,
cancellation_event=cancellation_event,
)
except LiveTransportError as exc:
failure = exc.__cause__
if exc.kind is LiveTransportFailureKind.REJECTED and isinstance(failure, OSError):
exc.diagnostic_reason = LiveTransportFailureReason.ENDPOINT_REJECTED
if (
exc.kind is not LiveTransportFailureKind.UNAVAILABLE
or not isinstance(failure, OSError)
or not _retryable_windows_pipe_open_error(failure)
):
raise
# WaitNamedPipeW is advisory: another concurrent client may consume
# the available server instance before the CRT open. Windows ARM64
# has surfaced that race as EINVAL; repeat the wait under the same
# caller deadline instead of rejecting a healthy multi-client owner.
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline, io_attempted=True)
class _Connection(Protocol):
def recv(self, size: int, /) -> bytes: ...
def sendall(self, data: bytes, /) -> None: ...
def close(self) -> None: ...
@dataclass(repr=False)
class _WindowsQueuedWrite:
data: bytes
deadline: float | None
cancellation_event: threading.Event | None
before_send: Callable[[], None] | None
completed: threading.Event = field(default_factory=threading.Event)
state: Literal["queued", "writing", "finished", "abandoned"] = "queued"
failure: BaseException | None = None
preflight_failed: bool = False
io_attempted: bool = False
class WindowsNamedPipeConnection:
"""Bound a local Windows pipe to server PID, user SID, and logon session."""
_pipe: io.FileIO
_kernel: ctypes.CDLL
_read_deadline: float | None
_cancellation_event: threading.Event | None
_reader_thread_id: int | None
_post_initialization: bool
_io_condition: threading.Condition
_queued_writes: deque[_WindowsQueuedWrite]
_active_write: _WindowsQueuedWrite | None
_closed: bool
def __init__(
self,
endpoint: str,
*,
owner_pid: int,
deadline: float | None = None,
cancellation_event: threading.Event | None = None,
) -> None:
if os.name != "nt":
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native Windows control stream is unavailable",
)
name = endpoint[len(_WINDOWS_NAMED_PIPE_PREFIX) :]
if (
not endpoint.startswith(_WINDOWS_NAMED_PIPE_PREFIX)
or not 0 < len(name) <= 256
or "\\" in name
or "/" in name
or "\x00" in name
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native Windows control stream endpoint is malformed",
)
self._kernel = windows_kernel32()
self._read_deadline = deadline
self._cancellation_event = cancellation_event
self._reader_thread_id = None
self._post_initialization = False
self._io_condition = threading.Condition()
self._queued_writes = deque()
self._active_write = None
self._closed = False
try:
if deadline is None:
self._pipe = io.FileIO(endpoint, mode="r+")
else:
self._pipe = _open_bounded_windows_named_pipe(
endpoint,
kernel=self._kernel,
deadline=deadline,
cancellation_event=cancellation_event,
)
except LiveTransportError:
raise
except (AttributeError, OSError, ValueError) as exc:
kind = (
_io_failure_kind(exc)
if isinstance(exc, OSError)
else LiveTransportFailureKind.REJECTED
)
raise LiveTransportError(
kind,
"native Windows control stream is unavailable",
io_attempted=isinstance(exc, OSError) and exc.errno != errno.EBADF,
diagnostic_reason=(
LiveTransportFailureReason.ENDPOINT_REJECTED
if kind is LiveTransportFailureKind.REJECTED and isinstance(exc, OSError)
else None
),
) from exc
try:
self.require_owner(owner_pid)
except Exception:
self._pipe.close()
raise
def require_owner(self, owner_pid: int) -> None:
try:
self._require_owner(owner_pid)
except (AttributeError, OSError, ValueError) as exc:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native Windows control stream peer owner does not match",
) from exc
def _require_owner(self, owner_pid: int) -> None:
kernel = self._kernel
query_peer = kernel.GetNamedPipeServerProcessId
query_peer.argtypes = [ctypes.c_void_p, ctypes.POINTER(ctypes.c_uint32)]
query_peer.restype = ctypes.c_int
process_session = kernel.ProcessIdToSessionId
process_session.argtypes = [ctypes.c_uint32, ctypes.POINTER(ctypes.c_uint32)]
process_session.restype = ctypes.c_int
peer_pid = ctypes.c_uint32()
peer_session = ctypes.c_uint32()
client_session = ctypes.c_uint32()
handle = ctypes.c_void_p(windows_crt.get_osfhandle(self._pipe.fileno()))
if (
not query_peer(handle, ctypes.byref(peer_pid))
or peer_pid.value != owner_pid
or not process_session(peer_pid.value, ctypes.byref(peer_session))
or not process_session(os.getpid(), ctypes.byref(client_session))
or peer_session.value != client_session.value
or not windows_process_has_same_user(kernel, owner_pid)
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native Windows control stream peer owner does not match",
)
def clear_initialization_deadline(self) -> None:
with self._io_condition:
self._read_deadline = None
self._cancellation_event = None
# One synchronous Windows pipe handle cannot make progress when a
# blocking reader and writer use it from separate threads. After
# the sequential initialization exchange, recv's thread owns all
# kernel I/O and request threads submit bounded writes to it.
self._post_initialization = True
self._io_condition.notify_all()
def recv(self, size: int) -> bytes:
deadline = self._read_deadline
if self._post_initialization:
return self._recv_post_initialization(size)
if deadline is None:
self._reader_thread_id = threading.get_native_id()
received = self._pipe.read(size)
else:
received = run_bounded_windows_operation(
lambda: self._pipe.read(size),
kernel=self._kernel,
deadline=deadline,
cancellation_event=self._cancellation_event,
)
return b"" if received is None else received
def sendall(self, data: bytes) -> None:
if self._post_initialization:
self._enqueue_write(
data,
deadline=None,
cancellation_event=None,
before_send=None,
)
return
self._write_direct(data)
def _write_direct(self, data: bytes) -> None:
view = memoryview(data)
while view:
written = self._pipe.write(view)
if written is None or written <= 0:
raise OSError("native Windows control stream write failed")
view = view[written:]
def sendall_before(
self,
data: bytes,
*,
deadline: float,
cancellation_event: threading.Event | None,
before_send: Callable[[], None] | None = None,
) -> None:
if self._post_initialization:
self._enqueue_write(
data,
deadline=deadline,
cancellation_event=cancellation_event,
before_send=before_send,
)
return
if before_send is not None:
_run_callback(before_send, failure_message="native control stream preflight failed")
view = memoryview(data)
while view:
written = run_bounded_windows_operation(
lambda remaining=view: self._pipe.write(remaining),
kernel=self._kernel,
deadline=deadline,
cancellation_event=cancellation_event,
)
if written <= 0:
raise LiveTransportError(
LiveTransportFailureKind.UNAVAILABLE,
"native Windows control stream write failed",
io_attempted=True,
)
view = view[written:]
def _enqueue_write(
self,
data: bytes,
*,
deadline: float | None,
cancellation_event: threading.Event | None,
before_send: Callable[[], None] | None,
) -> None:
request = _WindowsQueuedWrite(
bytes(data),
deadline,
cancellation_event,
before_send,
)
with self._io_condition:
if self._closed:
raise LiveTransportError(
LiveTransportFailureKind.UNAVAILABLE,
"native Windows control stream is disconnected",
)
# LiveJsonRpcClient._write_lock admits one frame at a time. Keep
# that bound explicit if this internal connection is called
# directly so a stalled peer cannot accumulate an unbounded queue.
if self._queued_writes or self._active_write is not None:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native Windows control stream write is unavailable",
)
self._queued_writes.append(request)
self._io_condition.notify_all()
try:
while not request.completed.is_set():
_raise_if_cancelled(cancellation_event)
if deadline is None:
wait = None if cancellation_event is None else _CANCELLATION_CHECK_SECONDS
else:
remaining = _remaining_seconds(deadline, io_attempted=request.io_attempted)
wait = (
remaining
if cancellation_event is None
else min(remaining, _CANCELLATION_CHECK_SECONDS)
)
request.completed.wait(wait)
except LiveTransportError as exc:
with self._io_condition:
if request.state == "finished" and request.failure is None:
return
started = request.state == "writing"
if request.state in {"queued", "writing"}:
request.state = "abandoned"
request.failure = exc
request.completed.set()
self._io_condition.notify_all()
if started:
cancel_windows_thread(self._kernel, self._reader_thread_id)
raise
with self._io_condition:
failure = request.failure
state = request.state
if failure is not None:
if request.preflight_failed or isinstance(failure, LiveTransportError):
raise failure
kind = (
_io_failure_kind(failure)
if isinstance(failure, OSError)
else LiveTransportFailureKind.REJECTED
)
raise LiveTransportError(
kind,
"native Windows control stream write failed",
io_attempted=request.io_attempted,
) from failure
if state != "finished":
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native Windows control stream write failed",
)
def _take_queued_write(self) -> _WindowsQueuedWrite | None:
with self._io_condition:
if self._closed:
raise LiveTransportError(
LiveTransportFailureKind.UNAVAILABLE,
"native Windows control stream is closed",
)
while self._queued_writes:
request = self._queued_writes.popleft()
if request.state != "queued":
request.completed.set()
continue
try:
_raise_if_cancelled(request.cancellation_event)
if request.deadline is not None:
_remaining_seconds(request.deadline)
except LiveTransportError as exc:
request.state = "abandoned"
request.failure = exc
request.completed.set()
continue
request.state = "writing"
self._active_write = request
return request
return None
def _perform_queued_write(self, request: _WindowsQueuedWrite) -> None:
failure: BaseException | None = None
preflight_failed = False
try:
view = memoryview(request.data)
with self._io_condition:
if request.state != "writing":
return
_raise_if_cancelled(request.cancellation_event)
if request.deadline is not None:
_remaining_seconds(request.deadline)
if request.before_send is not None:
try:
_run_callback(
request.before_send,
failure_message="native control stream preflight failed",
)
except BaseException:
preflight_failed = True
raise
with self._io_condition:
if request.state != "writing":
return
_raise_if_cancelled(request.cancellation_event)
if request.deadline is not None:
_remaining_seconds(request.deadline)
while view:
with self._io_condition:
if request.state != "writing":
return
_raise_if_cancelled(request.cancellation_event)
if request.deadline is not None:
_remaining_seconds(request.deadline, io_attempted=request.io_attempted)
request.io_attempted = True
written = self._pipe.write(view)
if written is None or written <= 0:
raise OSError("native Windows control stream write failed")
view = view[written:]
except BaseException as exc:
failure = exc
finally:
with self._io_condition:
if request.state == "writing":
request.state = "finished"
request.failure = failure
request.preflight_failed = preflight_failed
if self._active_write is request:
self._active_write = None
request.completed.set()
self._io_condition.notify_all()
def _peek_available_bytes(self) -> int:
peek = self._kernel.PeekNamedPipe
peek.argtypes = [
ctypes.c_void_p,
ctypes.c_void_p,
ctypes.c_uint32,
ctypes.POINTER(ctypes.c_uint32),
ctypes.POINTER(ctypes.c_uint32),
ctypes.POINTER(ctypes.c_uint32),
]
peek.restype = ctypes.c_int
available = ctypes.c_uint32()
handle = ctypes.c_void_p(windows_crt.get_osfhandle(self._pipe.fileno()))
if not peek(handle, None, 0, None, ctypes.byref(available), None):
raise OSError(
_windows_ctypes.get_last_error(),
"native Windows control stream peek failed",
)
return available.value
def _recv_post_initialization(self, size: int) -> bytes:
self._reader_thread_id = threading.get_native_id()
while True:
request = self._take_queued_write()
if request is not None:
self._perform_queued_write(request)
continue
# PeekNamedPipe is nonblocking. Reading only bytes the kernel has
# already reported prevents an idle read from owning the
# synchronous file object while a request is waiting to write.
available = self._peek_available_bytes()
if available > 0:
received = self._pipe.read(min(size, available))
return b"" if received is None else received
with self._io_condition:
if self._closed:
raise LiveTransportError(
LiveTransportFailureKind.UNAVAILABLE,
"native Windows control stream is closed",
)
if not self._queued_writes:
self._io_condition.wait(_WINDOWS_PIPE_INPUT_POLL_SECONDS)
def close(self) -> None:
failure = LiveTransportError(
LiveTransportFailureKind.UNAVAILABLE,
"native Windows control stream disconnected",
)
with self._io_condition:
if self._closed:
return
self._closed = True
pending = tuple(self._queued_writes)
self._queued_writes.clear()
active = self._active_write
for request in (*pending, active):
if request is not None and request.state in {"queued", "writing"}:
request.state = "abandoned"
request.failure = failure
request.completed.set()
self._io_condition.notify_all()
cancel_windows_thread(self._kernel, self._reader_thread_id)
try:
cancel = self._kernel.CancelIoEx
cancel.argtypes = [ctypes.c_void_p, ctypes.c_void_p]
cancel.restype = ctypes.c_int
handle = ctypes.c_void_p(windows_crt.get_osfhandle(self._pipe.fileno()))
cancel(handle, None)
except (AttributeError, OSError, ValueError):
pass
self._pipe.close()
def _unix_endpoint_failure_reason(error: OSError) -> LiveTransportFailureReason | None:
if isinstance(error, (PermissionError, NotADirectoryError)):
return LiveTransportFailureReason.ENDPOINT_REJECTED
if isinstance(error, FileNotFoundError):
return LiveTransportFailureReason.ENDPOINT_MISSING
return None
def _socket_metadata(path: Path) -> os.stat_result:
try:
return path.lstat()
except OSError as exc:
raise LiveTransportError(
_io_failure_kind(exc),
"native control stream endpoint is unavailable",
diagnostic_reason=_unix_endpoint_failure_reason(exc),
) from exc
def _require_private_socket(endpoint: Path) -> os.stat_result:
parent_metadata = _socket_metadata(endpoint.parent)
if (
not stat.S_ISDIR(parent_metadata.st_mode)
or stat.S_IMODE(parent_metadata.st_mode) != 0o700
or parent_metadata.st_uid != os.geteuid()
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream endpoint is unsafe",
diagnostic_reason=LiveTransportFailureReason.ENDPOINT_REJECTED,
)
metadata = _socket_metadata(endpoint)
if (
not stat.S_ISSOCK(metadata.st_mode)
or stat.S_IMODE(metadata.st_mode) != 0o600
or metadata.st_uid != os.geteuid()
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream endpoint is unsafe",
diagnostic_reason=LiveTransportFailureReason.ENDPOINT_REJECTED,
)
try:
if endpoint.resolve(strict=True) != endpoint.absolute():
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream endpoint is unsafe",
diagnostic_reason=LiveTransportFailureReason.ENDPOINT_REJECTED,
)
except OSError as exc:
raise LiveTransportError(
_io_failure_kind(exc),
"native control stream endpoint is unavailable",
diagnostic_reason=_unix_endpoint_failure_reason(exc),
) from exc
return metadata
def require_unix_peer(connection: socket.socket, *, owner_pid: int) -> None:
if sys.platform != "darwin":
return
try:
peer_pid = struct.unpack(
"=i",
connection.getsockopt(_LOCAL_SOCKET_PROTOCOL, _LOCAL_PEER_PROCESS_ID, 4),
)[0]
credentials = connection.getsockopt(
_LOCAL_SOCKET_PROTOCOL,
_LOCAL_PEER_CREDENTIALS,
_MAXIMUM_PEER_CREDENTIAL_BYTES,
)
if len(credentials) < 8:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream peer credentials are unavailable",
)
peer_uid = struct.unpack_from("=I", credentials, 4)[0]
except (OSError, struct.error) as exc:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream peer credentials are unavailable",
) from exc
if peer_pid != owner_pid or peer_uid != os.geteuid():
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream peer owner does not match",
)
def _connect_unix_socket(endpoint: str, *, owner_pid: int, deadline: float) -> socket.socket:
path = Path(endpoint)
if not path.is_absolute():
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream endpoint is malformed",
)
before = _require_private_socket(path)
connection = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
io_attempted = False
try:
connection.settimeout(_remaining_seconds(deadline))
io_attempted = True
connection.connect(str(path))
io_attempted = False
require_unix_peer(connection, owner_pid=owner_pid)
after = _require_private_socket(path)
if (before.st_dev, before.st_ino) != (after.st_dev, after.st_ino):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream endpoint changed",
)
except (LiveTransportError, OSError) as exc:
connection.close()
if isinstance(exc, LiveTransportError):
raise
raise LiveTransportError(
_io_failure_kind(exc),
"native control stream is unavailable",
io_attempted=io_attempted,
diagnostic_reason=_unix_endpoint_failure_reason(exc),
) from exc
return connection
def _require_connection_owner(connection: _Connection, *, owner_pid: int) -> None:
if isinstance(connection, socket.socket):
require_unix_peer(connection, owner_pid=owner_pid)
elif isinstance(connection, WindowsNamedPipeConnection):
connection.require_owner(owner_pid)
def _close_connection(connection: _Connection) -> None:
if isinstance(connection, socket.socket):
try:
connection.shutdown(socket.SHUT_RDWR)
except OSError:
pass
try:
connection.close()
except (OSError, ValueError):
pass
def _unique_json_object(pairs: list[tuple[str, object]]) -> dict[str, object]:
result: dict[str, object] = {}
for key, value in pairs:
if key in result:
raise ValueError("duplicate JSON object key")
result[key] = value
return result
def _is_json_object(value: object) -> TypeGuard[dict[str, object]]:
"""Narrow objects produced by the typed, duplicate-free JSON decoder."""
return isinstance(value, dict)
def _recv_unix_socket_before(
connection: socket.socket,
count: int,
*,
deadline: float,
cancellation_event: threading.Event | None,
) -> bytes:
io_attempted = False
while True:
_raise_if_cancelled(cancellation_event)
remaining = _remaining_seconds(deadline, io_attempted=io_attempted)
wait = (
remaining if cancellation_event is None else min(remaining, _CANCELLATION_CHECK_SECONDS)
)
try:
io_attempted = True
readable, _, _ = select.select([connection], [], [], wait)
except InterruptedError:
continue
if not readable:
continue
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline, io_attempted=True)
try:
return connection.recv(count, socket.MSG_DONTWAIT)
except (BlockingIOError, InterruptedError):
continue
def _read_exact(
connection: _Connection,
count: int,
*,
deadline: float | None = None,
cancellation_event: threading.Event | None = None,
) -> bytes:
result = bytearray()
while len(result) < count:
_raise_if_cancelled(cancellation_event)
if deadline is not None:
_remaining_seconds(deadline)
try:
chunk_size = min(count - len(result), _MAXIMUM_READ_CHUNK_BYTES)
if isinstance(connection, socket.socket) and deadline is not None:
chunk = _recv_unix_socket_before(
connection,
chunk_size,
deadline=deadline,
cancellation_event=cancellation_event,
)
else:
chunk = connection.recv(chunk_size)
except (OSError, ValueError) as exc:
raise LiveTransportError(
_io_failure_kind(exc),
"native control stream disconnected",
io_attempted=isinstance(exc, OSError) and exc.errno != errno.EBADF,
) from exc
if not chunk:
raise LiveTransportError(
LiveTransportFailureKind.UNAVAILABLE,
"native control stream disconnected",
io_attempted=True,
)
result.extend(chunk)
return bytes(result)
def read_message(
connection: _Connection,
*,
maximum_frame_bytes: int,
deadline: float | None = None,
cancellation_event: threading.Event | None = None,
) -> dict[str, object]:
frame_length = int.from_bytes(
_read_exact(connection, 4, deadline=deadline, cancellation_event=cancellation_event),
byteorder="little",
)
if not 0 < frame_length <= maximum_frame_bytes:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream frame exceeds its limit",
)
try:
value: object = json.loads(
_read_exact(
connection,
frame_length,
deadline=deadline,
cancellation_event=cancellation_event,
).decode("utf-8"),
object_pairs_hook=_unique_json_object,
parse_constant=int,
)
except (UnicodeError, ValueError, json.JSONDecodeError) as exc:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream message is malformed",
) from exc
if not _is_json_object(value):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream message is malformed",
)
return dict(value)
def canonical_json(value: object) -> str:
try:
return json.dumps(
value,
ensure_ascii=False,
separators=(",", ":"),
sort_keys=True,
allow_nan=False,
)
except (OverflowError, RecursionError, TypeError, UnicodeError, ValueError) as exc:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream message is malformed",
) from exc
def _send_unix_socket_before(
connection: socket.socket,
data: bytes,
*,
deadline: float,
cancellation_event: threading.Event | None,
before_send: Callable[[], None] | None,
) -> None:
view = memoryview(data)
wrote_any = False
io_attempted = False
while view:
_raise_if_cancelled(cancellation_event)
remaining = _remaining_seconds(deadline, io_attempted=io_attempted)
wait = (
remaining if cancellation_event is None else min(remaining, _CANCELLATION_CHECK_SECONDS)
)
try:
io_attempted = True
_, writable, _ = select.select([], [connection], [], wait)
except InterruptedError:
continue
if not writable:
continue
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline, io_attempted=True)
if not wrote_any and before_send is not None:
_run_callback(before_send, failure_message="native control stream preflight failed")
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline)
try:
written = connection.send(view, socket.MSG_DONTWAIT)
except (BlockingIOError, InterruptedError):
continue
if written <= 0:
raise OSError("native Unix control stream write failed")
wrote_any = True
view = view[written:]
def write_message(
connection: _Connection,
message: dict[str, object],
*,
maximum_frame_bytes: int,
deadline: float | None = None,
cancellation_event: threading.Event | None = None,
before_send: Callable[[], None] | None = None,
) -> None:
encoded = canonical_json(message).encode("utf-8")
if not 0 < len(encoded) <= maximum_frame_bytes:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream frame exceeds its limit",
)
try:
payload = len(encoded).to_bytes(4, byteorder="little") + encoded
if isinstance(connection, socket.socket) and deadline is not None:
_send_unix_socket_before(
connection,
payload,
deadline=deadline,
cancellation_event=cancellation_event,
before_send=before_send,
)
elif isinstance(connection, WindowsNamedPipeConnection) and deadline is not None:
connection.sendall_before(
payload,
deadline=deadline,
cancellation_event=cancellation_event,
before_send=before_send,
)
else:
if before_send is not None:
_run_callback(
before_send,
failure_message="native control stream preflight failed",
)
connection.sendall(payload)
except (OSError, ValueError) as exc:
raise LiveTransportError(
_io_failure_kind(exc),
"native control stream disconnected",
io_attempted=isinstance(exc, OSError) and exc.errno != errno.EBADF,
) from exc
@dataclass(repr=False)
class _PendingResponse:
request_id: str
fingerprint: str
ready: threading.Event = field(default_factory=threading.Event)
response: dict[str, object] | None = None
failure: RuntimeError | None = None
@dataclass(frozen=True, repr=False)
class _CompletedResponse:
fingerprint: str
response: dict[str, object]
class LiveJsonRpcClient:
"""A persistent local peer, initialization exchange, and JSON-RPC router."""
def __init__(
self,
owner: LiveOwner,
*,
maximum_frame_bytes: int = _DEFAULT_MAXIMUM_FRAME_BYTES,
verify_owner: Callable[[], None] | None = None,
on_notification: Callable[[dict[str, object]], None] | None = None,
validate_response: Callable[[dict[str, object]], None] | None = None,
on_disconnect: Callable[[], None] | None = None,
) -> None:
if (
type(owner.pid) is not int
or not 0 < owner.pid <= 2_147_483_647
or type(maximum_frame_bytes) is not int
or not 0 < maximum_frame_bytes <= _DEFAULT_MAXIMUM_FRAME_BYTES
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream owner is malformed",
)
self._owner = owner
self._maximum_frame_bytes = maximum_frame_bytes
self._verify_owner = verify_owner
self._on_notification = on_notification
self._validate_response = validate_response
self._on_disconnect = on_disconnect
self._connect_lock = threading.RLock()
self._lock = threading.RLock()
self._write_lock = threading.Lock()
self._connection: _Connection | None = None
self._connection_on_disconnect: Callable[[], None] | None = None
self._reader: threading.Thread | None = None
self._reader_processing: tuple[_Connection, bool] | None = None
self._pending: dict[str, _PendingResponse] = {}
self._completed: OrderedDict[str, _CompletedResponse] = OrderedDict()
self._ambiguous: OrderedDict[str, None] = OrderedDict()
self._abandoned: OrderedDict[str, None] = OrderedDict()
@property
def connected(self) -> bool:
return self._connection is not None
def connect(
self,
initialization: dict[str, object],
*,
validate_result: Callable[[dict[str, object]], None] | None = None,
deadline: float,
cancellation_event: threading.Event | None = None,
on_initialized: Callable[[dict[str, object]], None] | None = None,
on_notification: Callable[[dict[str, object]], None] | None = None,
on_disconnect: Callable[[], None] | None = None,
) -> dict[str, object]:
"""Initialize one connection with optional connection-bound callbacks.
A callback may finish after close or a later initialization. Consumers
that cache state should capture their connection generation in these
callbacks and reject publication from a retired generation.
"""
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline)
with self._connect_lock:
if self._connection is not None:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream is already initialized",
)
owner = self._owner
if owner.transport == "unix-domain-socket":
connection: _Connection = _connect_unix_socket(
owner.endpoint,
owner_pid=owner.pid,
deadline=deadline,
)
elif owner.transport == "windows-named-pipe":
connection = WindowsNamedPipeConnection(
owner.endpoint,
owner_pid=owner.pid,
deadline=deadline,
cancellation_event=cancellation_event,
)
else:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream transport is unsupported",
)
try:
self._check_owner(connection)
request_id = initialization.get("id")
if (
initialization.get("jsonrpc") != "2.0"
or not isinstance(request_id, str)
or not isinstance(initialization.get("method"), str)
or not isinstance(initialization.get("params"), dict)
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream initialization is malformed",
)
write_message(
connection,
initialization,
maximum_frame_bytes=self._maximum_frame_bytes,
deadline=deadline,
cancellation_event=cancellation_event,
)
response = read_message(
connection,
maximum_frame_bytes=self._maximum_frame_bytes,
deadline=deadline,
cancellation_event=cancellation_event,
)
result_value = response.get("result")
if (
response.get("jsonrpc") != "2.0"
or response.get("id") != request_id
or "error" in response
or not _is_json_object(result_value)
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream initialization failed",
)
result = dict(result_value)
if validate_result is not None:
validate_result(result)
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline)
self._check_owner(connection)
if on_initialized is not None:
on_initialized(result)
if isinstance(connection, socket.socket):
connection.settimeout(None)
else:
connection.clear_initialization_deadline()
reader = threading.Thread(
target=self._receive_loop,
args=(connection, on_notification or self._on_notification),
daemon=True,
name="meetings-native-control-stream",
)
with self._lock:
self._connection = connection
self._connection_on_disconnect = on_disconnect or self._on_disconnect
self._reader = reader
reader.start()
return result
except Exception as exc:
_close_connection(connection)
if isinstance(exc, LiveTransportError):
raise
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream initialization failed",
) from exc
def request(
self,
message: dict[str, object],
*,
deadline: float,
cancellation_event: threading.Event | None = None,
before_send: Callable[[], None] | None = None,
on_sent: Callable[[], None] | None = None,
) -> dict[str, object]:
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline)
request_id = message.get("id")
if (
message.get("jsonrpc") != "2.0"
or not isinstance(request_id, str)
or not isinstance(message.get("method"), str)
or not isinstance(message.get("params"), dict)
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream request is malformed",
)
fingerprint = canonical_json(message)
pending = _PendingResponse(request_id, fingerprint)
with self._lock:
connection = self._connection
if connection is None:
raise LiveTransportError(
LiveTransportFailureKind.UNAVAILABLE,
"native control stream is disconnected",
)
completed = self._completed.get(request_id)
if completed is not None:
if completed.fingerprint != fingerprint:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream request identity was reused",
)
self._check_owner(connection)
if self._connection is not connection:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream owner changed",
)
return copy.deepcopy(completed.response)
if request_id in self._ambiguous:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream request outcome is ambiguous",
)
if (
request_id in self._pending
or len(self._pending) >= _DEFAULT_MAXIMUM_PENDING_REQUESTS
):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream request identity is unavailable",
)
self._pending[request_id] = pending
sent = False
writing = False
try:
with self._write_lock:
_raise_if_cancelled(cancellation_event)
_remaining_seconds(deadline)
with self._lock:
if self._connection is not connection:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream owner changed",
)
self._check_owner(connection)
writing = True
write_message(
connection,
message,
maximum_frame_bytes=self._maximum_frame_bytes,
deadline=deadline,
cancellation_event=cancellation_event,
before_send=before_send,
)
sent = True
if on_sent is not None:
_run_callback(on_sent, failure_message="native control stream request failed")
return self._await_response(
pending,
connection=connection,
deadline=deadline,
cancellation_event=cancellation_event,
)
except Exception as exc:
if (
isinstance(exc, LiveTransportError)
and exc.kind
in {
LiveTransportFailureKind.CANCELLED,
LiveTransportFailureKind.TIMED_OUT,
}
and (sent or not writing)
):
if sent:
with self._lock:
self._abandoned[request_id] = None
while len(self._abandoned) > _DEFAULT_MAXIMUM_CACHED_RESPONSES:
self._abandoned.popitem(last=False)
raise
self._remember_ambiguous(request_id)
self._disconnect(connection, exc)
raise
finally:
with self._lock:
self._pending.pop(request_id, None)
def close(self) -> None:
with self._connect_lock:
connection = self._connection
if connection is not None:
self._disconnect(connection)
def _check_owner(self, connection: _Connection) -> None:
try:
_require_connection_owner(connection, owner_pid=self._owner.pid)
if self._verify_owner is not None:
self._verify_owner()
except LiveTransportError:
raise
except Exception as exc:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED, "native control stream owner changed"
) from exc
def _await_response(
self,
pending: _PendingResponse,
*,
connection: _Connection,
deadline: float,
cancellation_event: threading.Event | None,
) -> dict[str, object]:
io_attempted = False
reader_processing = self._reader_processing
while True:
_raise_if_cancelled(cancellation_event)
# A local callback can prevent the reader from consuming a reply.
# Keep the original marker so even a completed callback invalidates
# peer-timeout evidence for this wait. Reads must not wait on a lock
# held by that local processing beyond this request's deadline.
current_processing = self._reader_processing
local_processing = current_processing is not reader_processing or (
current_processing is not None
and current_processing[0] is connection
and current_processing[1]
)
remaining = _remaining_seconds(
deadline,
io_attempted=io_attempted and not local_processing and self._connection is connection,
)
wait = (
remaining
if cancellation_event is None
else min(remaining, _CANCELLATION_CHECK_SECONDS)
)
io_attempted = True
if not pending.ready.wait(wait):
continue
if pending.failure is not None:
raise pending.failure
if pending.response is None:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream response is unavailable",
)
return pending.response
def _receive_loop(
self,
connection: _Connection,
on_notification: Callable[[dict[str, object]], None] | None,
) -> None:
try:
while True:
with self._lock:
if self._connection is not connection:
return
message = read_message(
connection,
maximum_frame_bytes=self._maximum_frame_bytes,
)
# The read can finish after another thread closes this stream.
# In-flight callbacks still retain this connection's consumer.
with self._lock:
if self._connection is not connection:
return
processing = (connection, True)
self._reader_processing = processing
try:
if message.get("jsonrpc") != "2.0":
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream message is malformed",
)
if "method" in message and "id" not in message:
if on_notification is None:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream notification is unexpected",
)
on_notification(message)
elif "id" in message and "method" not in message:
self._accept_response(message)
else:
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream message direction is invalid",
)
finally:
with self._lock:
if self._reader_processing is processing:
self._reader_processing = (connection, False)
except Exception as exc:
self._disconnect(connection, exc)
def _accept_response(self, message: dict[str, object]) -> None:
request_id = message.get("id")
if isinstance(request_id, str):
with self._lock:
abandoned = request_id in self._abandoned
if abandoned:
self._abandoned.pop(request_id)
if abandoned:
if "error" not in message and not _is_json_object(message.get("result")):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream response result is malformed",
)
return
if self._validate_response is not None:
self._validate_response(message)
if not isinstance(request_id, str):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream response identity is malformed",
)
with self._lock:
pending = self._pending.get(request_id)
if pending is None or pending.ready.is_set():
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream response request is unknown",
)
if "error" in message:
pending.failure = LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream request was declined",
)
pending.ready.set()
return
result_value = message.get("result")
if not _is_json_object(result_value):
raise LiveTransportError(
LiveTransportFailureKind.REJECTED,
"native control stream response result is malformed",
)
response: dict[str, object] = copy.deepcopy(dict(result_value))
self._completed[request_id] = _CompletedResponse(
pending.fingerprint,
copy.deepcopy(response),
)
while len(self._completed) > _DEFAULT_MAXIMUM_CACHED_RESPONSES:
self._completed.popitem(last=False)
# Publish the response atomically with its completed-cache entry.
# An update owner may close immediately after the last response
# byte. EOF must not overtake that accepted response and turn an
# already-committed handoff into an ambiguous client failure.
pending.response = response
pending.ready.set()
def _remember_ambiguous(self, request_id: str) -> None:
with self._lock:
self._ambiguous[request_id] = None
while len(self._ambiguous) > _DEFAULT_MAXIMUM_CACHED_RESPONSES:
self._ambiguous.popitem(last=False)
def _disconnect(self, connection: _Connection, failure: BaseException | None = None) -> None:
with self._lock:
if self._connection is not connection:
return
self._connection = None
on_disconnect, self._connection_on_disconnect = (
self._connection_on_disconnect,
None,
)
self._reader = None
self._completed.clear()
self._abandoned.clear()
# A response accepted before EOF remains authoritative. Only
# requests without a terminal response become ambiguous when the
# owner disconnects.
pending = tuple(
request for request in self._pending.values() if not request.ready.is_set()
)
for request in pending:
self._ambiguous[request.request_id] = None
while len(self._ambiguous) > _DEFAULT_MAXIMUM_CACHED_RESPONSES:
self._ambiguous.popitem(last=False)
if on_disconnect is not None:
try:
on_disconnect()
except Exception:
pass
_close_connection(connection)
if isinstance(failure, LiveTransportError):
error = failure
elif failure is not None:
error = LiveTransportError(
LiveTransportFailureKind.REJECTED, "native control stream disconnected"
)
error.__cause__ = failure
else:
error = LiveTransportError(
LiveTransportFailureKind.UNAVAILABLE,
"native control stream disconnected",
)
for request in pending:
request.failure = error
request.ready.set()
SHA-256: f1817f5685b3a2d5721faedad501c1a7ae862deb7fe21602ec9702476b69266b