refactor(gateway/platforms): group D docstring closer folds
This commit is contained in:
@@ -277,8 +277,7 @@ class BlueBubblesAdapter(BasePlatformAdapter):
|
||||
|
||||
BlueBubbles posts to the exact registered URL and its registration
|
||||
API cannot set custom headers, so this is the only way to
|
||||
authenticate inbound webhooks without disabling auth.
|
||||
"""
|
||||
authenticate inbound webhooks without disabling auth."""
|
||||
return self._webhook_register_url_with(quote(self.password, safe=""))
|
||||
|
||||
@property
|
||||
@@ -345,8 +344,7 @@ class BlueBubblesAdapter(BasePlatformAdapter):
|
||||
membership is intentionally NOT a fallback: the same contact appears in
|
||||
a 1:1 DM and any number of groups, so a participant match could leak a
|
||||
DM reply into a group thread. Return ``None`` and let the caller create
|
||||
a fresh DM via ``_create_chat_for_handle``.
|
||||
"""
|
||||
a fresh DM via ``_create_chat_for_handle``."""
|
||||
target = (target or "").strip()
|
||||
if not target:
|
||||
return None
|
||||
|
||||
@@ -58,8 +58,7 @@ def _parse_allowed_source_cidrs(raw: Any) -> list[ipaddress._BaseNetwork]:
|
||||
|
||||
When populated, requests from source IPs outside every listed CIDR are
|
||||
rejected with 403 before the body is parsed (restrict to Microsoft
|
||||
Graph's published webhook source ranges in production).
|
||||
"""
|
||||
Graph's published webhook source ranges in production)."""
|
||||
if isinstance(raw, str):
|
||||
candidates = raw.split(",")
|
||||
elif isinstance(raw, (list, tuple, set)):
|
||||
|
||||
@@ -68,8 +68,7 @@ def _guess_extension(data: bytes) -> str:
|
||||
"""Guess file extension from magic bytes.
|
||||
|
||||
WEBP is claimed before the shared audio/AV sniffer (it shares RIFF with WAVE);
|
||||
tools/audio_container.py owns the MP3-vs-ADTS-AAC sync-word disambiguation.
|
||||
"""
|
||||
tools/audio_container.py owns the MP3-vs-ADTS-AAC sync-word disambiguation."""
|
||||
for magic, ext in _MAGIC_EXTENSIONS:
|
||||
if data.startswith(magic):
|
||||
return ext
|
||||
@@ -98,8 +97,7 @@ def _remux_aac_to_m4a(aac_data: bytes) -> Optional[Tuple[bytes, str]]:
|
||||
"""Losslessly remux raw ADTS AAC (Android voice notes, rejected by most STT APIs) to .m4a.
|
||||
|
||||
Returns ``(m4a_bytes, ".m4a")`` or ``None`` when ffmpeg is missing or fails —
|
||||
callers must then pass the input through unchanged.
|
||||
"""
|
||||
callers must then pass the input through unchanged."""
|
||||
# Fall back to common Homebrew/local prefixes on macOS dev hosts.
|
||||
ffmpeg = shutil.which("ffmpeg") or next(
|
||||
(p for p in ("/opt/homebrew/bin/ffmpeg", "/usr/local/bin/ffmpeg")
|
||||
@@ -415,8 +413,7 @@ class SignalAdapter(BasePlatformAdapter):
|
||||
"""Gate on require_mention (False = drop) and strip the bot's own @mention.
|
||||
|
||||
The self-mention is stripped from every group message so the agent doesn't
|
||||
read "@+155****4567 say hello" as a directive to contact that number.
|
||||
"""
|
||||
read "@+155****4567 say hello" as a directive to contact that number."""
|
||||
account_norm = self._account_normalized
|
||||
if self.require_mention:
|
||||
mentioned_in_text = account_norm and (f"@{account_norm}" in (text or ""))
|
||||
@@ -679,8 +676,7 @@ class SignalAdapter(BasePlatformAdapter):
|
||||
|
||||
``log_failures=False`` logs failures at DEBUG (typing path: silence repeated
|
||||
NETWORK_FAILURE spam). ``raise_on_rate_limit=True`` raises ``SignalRateLimitError``
|
||||
on a 429 / RateLimitException instead of swallowing it (multi-attachment sends).
|
||||
"""
|
||||
on a 429 / RateLimitException instead of swallowing it (multi-attachment sends)."""
|
||||
if not self.client:
|
||||
logger.warning("Signal: RPC called but client not connected")
|
||||
return None
|
||||
@@ -770,8 +766,7 @@ class SignalAdapter(BasePlatformAdapter):
|
||||
"""Split converted Signal text into chunks, translating body ranges per chunk.
|
||||
|
||||
Splitting after conversion (not before) keeps styles that cross a chunk boundary
|
||||
intact instead of leaking literal Markdown markers.
|
||||
"""
|
||||
intact instead of leaking literal Markdown markers."""
|
||||
if utf16_len(plain_text) <= max_length:
|
||||
return [(plain_text, text_styles)]
|
||||
indicator_reserve = 10 # Mirrors BasePlatformAdapter.truncate_message().
|
||||
@@ -864,8 +859,7 @@ class SignalAdapter(BasePlatformAdapter):
|
||||
"""Send a typing indicator (called every ~2s by base.py's ``_keep_typing``).
|
||||
|
||||
On NETWORK_FAILURE only the first consecutive failure logs at WARNING, and after
|
||||
three failures the RPC is skipped for an exponential cooldown; success resets.
|
||||
"""
|
||||
three failures the RPC is skipped for an exponential cooldown; success resets."""
|
||||
now = time.monotonic()
|
||||
if now < self._typing_skip_until.get(chat_id, 0.0):
|
||||
return
|
||||
@@ -905,8 +899,7 @@ class SignalAdapter(BasePlatformAdapter):
|
||||
"""Send a batch of images via chunked Signal RPC calls.
|
||||
|
||||
Alt texts are dropped (one shared body per send). Bad images are skipped with a
|
||||
warning. ``human_delay`` is ignored: the rate-limit scheduler paces batches.
|
||||
"""
|
||||
warning. ``human_delay`` is ignored: the rate-limit scheduler paces batches."""
|
||||
if not images:
|
||||
return
|
||||
scheduler = get_scheduler()
|
||||
@@ -951,8 +944,7 @@ class SignalAdapter(BasePlatformAdapter):
|
||||
"""Send one attachment batch with rate-limit pacing and a single transient retry.
|
||||
|
||||
Tokens are deducted only on validated success (a None result means the server
|
||||
never accepted the batch); 429s feed the scheduler before the retry.
|
||||
"""
|
||||
never accepted the batch); 429s feed the scheduler before the retry."""
|
||||
send_timeout = _signal_send_timeout(n)
|
||||
max_attempts = SIGNAL_RATE_LIMIT_MAX_ATTEMPTS
|
||||
for attempt in range(1, max_attempts + 1):
|
||||
|
||||
@@ -28,8 +28,7 @@ def _normalize_bullet_markers(source: str) -> str:
|
||||
"""Replace Markdown bullet markers with plain Unicode bullets.
|
||||
|
||||
Signal renders ``- item`` / ``* item`` literally. Fenced code blocks are kept
|
||||
byte-for-byte; list-looking lines inside code are code, not prose bullets.
|
||||
"""
|
||||
byte-for-byte; list-looking lines inside code are code, not prose bullets."""
|
||||
parts = re.split(r"(```.*?```)", source, flags=re.DOTALL)
|
||||
for idx, part in enumerate(parts):
|
||||
if idx % 2 == 0:
|
||||
@@ -42,8 +41,7 @@ def markdown_to_signal(text: str) -> tuple[str, list[str]]:
|
||||
|
||||
Signal uses ``bodyRanges`` (signal-cli ``textStyle`` / ``textStyles`` params) in
|
||||
the form ``start:length:STYLE``, positions in UTF-16 code units.
|
||||
Supported styles: BOLD, ITALIC, STRIKETHROUGH, MONOSPACE.
|
||||
"""
|
||||
Supported styles: BOLD, ITALIC, STRIKETHROUGH, MONOSPACE."""
|
||||
text = re.sub(r"\n{3,}", "\n\n", text).strip()
|
||||
text = _normalize_bullet_markers(text)
|
||||
styles: list[tuple[int, int, str]] = []
|
||||
|
||||
@@ -56,8 +56,7 @@ def _extract_retry_after_seconds(err: Any) -> Optional[float]:
|
||||
|
||||
Sources, in order: ``error.data.response.results[*].retryAfterSeconds`` (signal-cli
|
||||
≥ v0.14.3), then "Retry after N seconds" parsed from the message (libsignal-net
|
||||
RetryLaterException wrapped as AttachmentInvalidException, structured field null).
|
||||
"""
|
||||
RetryLaterException wrapped as AttachmentInvalidException, structured field null)."""
|
||||
if isinstance(err, dict):
|
||||
results = ((err.get("data") or {}).get("response") or {}).get("results") or []
|
||||
candidates = [
|
||||
@@ -105,8 +104,7 @@ class SignalAttachmentScheduler:
|
||||
Holds up to ``capacity`` tokens (default 50 = Signal's server bucket); each
|
||||
attachment consumes one. Tokens refill at ``refill_rate``/s, calibrated from the
|
||||
server's per-token Retry-After once a 429 has been observed (default 1 token / 4s).
|
||||
``acquire(n)`` calls serialize through an ``asyncio.Lock`` — FIFO across sessions.
|
||||
"""
|
||||
``acquire(n)`` calls serialize through an ``asyncio.Lock`` — FIFO across sessions."""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
@@ -145,8 +143,7 @@ class SignalAttachmentScheduler:
|
||||
capacity; call ``report_rpc_duration()`` after the RPC to sync. The lock is
|
||||
released during ``asyncio.sleep`` so other callers interleave, and the loop
|
||||
re-checks after each sleep in case the deadline was pessimistic. Signal's
|
||||
server is ground truth and will 429 (→ requeue) if the model drifts.
|
||||
"""
|
||||
server is ground truth and will 429 (→ requeue) if the model drifts."""
|
||||
if n <= 0:
|
||||
return 0.0
|
||||
if n > self.capacity:
|
||||
@@ -178,8 +175,7 @@ class SignalAttachmentScheduler:
|
||||
|
||||
No refill is credited for the upload window: Signal's server checks the bucket
|
||||
at RPC start and resumes refill only after the response, so crediting it causes
|
||||
cumulative drift that eventually triggers 429s. Advances ``last_refill``.
|
||||
"""
|
||||
cumulative drift that eventually triggers 429s. Advances ``last_refill``."""
|
||||
if n_attachments <= 0:
|
||||
return
|
||||
async with self._lock:
|
||||
|
||||
@@ -76,8 +76,7 @@ def _hmac_str_equal(provided: str, expected: str) -> bool:
|
||||
"""Timing-safe str equality tolerant of non-ASCII.
|
||||
|
||||
``compare_digest`` raises TypeError on non-ASCII str; ``provided`` is an
|
||||
attacker-controlled header, so compare as UTF-8 bytes to fail closed.
|
||||
"""
|
||||
attacker-controlled header, so compare as UTF-8 bytes to fail closed."""
|
||||
return hmac.compare_digest(provided.encode(), expected.encode())
|
||||
|
||||
|
||||
@@ -402,8 +401,7 @@ class WebhookAdapter(BasePlatformAdapter):
|
||||
|
||||
Returns None (no prefix, or multiplexing off and the prefix names this
|
||||
gateway's own profile), the profile name (served under multiplexing), or
|
||||
``_PROFILE_REJECTED`` (unknown / not served → 404).
|
||||
"""
|
||||
``_PROFILE_REJECTED`` (unknown / not served → 404)."""
|
||||
profile = (request.match_info.get("profile") or "").strip()
|
||||
if not profile:
|
||||
return None
|
||||
@@ -642,8 +640,7 @@ class WebhookAdapter(BasePlatformAdapter):
|
||||
|
||||
``prune_sessions`` only reaps rows with ``ended_at`` set, so unclosed webhook
|
||||
sessions leak unbounded. This hook fires at the true end of the run (success,
|
||||
failure, cancellation); ``end_session()`` is first-reason-wins.
|
||||
"""
|
||||
failure, cancellation); ``end_session()`` is first-reason-wins."""
|
||||
await self._end_webhook_session(event, event.source.chat_id)
|
||||
|
||||
async def _end_webhook_session(self, event: "MessageEvent", session_chat_id: str) -> None:
|
||||
@@ -756,8 +753,7 @@ class WebhookAdapter(BasePlatformAdapter):
|
||||
def _render_prompt(self, template: str, payload: dict, event_type: str, route_name: str) -> str:
|
||||
"""Render a prompt template with dot-notation payload access (``{pull_request.title}``).
|
||||
|
||||
``{__raw__}`` dumps the whole payload as indented JSON (truncated to 4000 chars).
|
||||
"""
|
||||
``{__raw__}`` dumps the whole payload as indented JSON (truncated to 4000 chars)."""
|
||||
if not template:
|
||||
truncated = json.dumps(payload, indent=2)[:4000]
|
||||
return f"Webhook event '{event_type}' on route '{route_name}':\n\n```json\n{truncated}\n```"
|
||||
|
||||
Reference in New Issue
Block a user