mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 20:27:56 +00:00
WIP: lane C network hygiene implementation (git bounded, keepalive, telegram backoff)
This commit is contained in:
parent
4a5505899e
commit
a23e12b13e
13 changed files with 297 additions and 56 deletions
|
|
@ -1501,6 +1501,14 @@ detector, while explicit task deadline, budget, cancellation and the absolute ce
|
|||
remain independent hard axes. A stale terminal from an earlier retry or task attempt
|
||||
cannot clear the current row.
|
||||
|
||||
Every remote httpx LLM client (the cached OpenAI-compatible clients and the
|
||||
no-proxy per-call clients) is built on one shared transport factory that sets
|
||||
platform-guarded TCP keepalive socket options, so a NAT/VPN mapping silently
|
||||
dropped during a long silent reasoning stretch is detected by kernel probes
|
||||
within minutes instead of hanging until the read timeout. Disclosed residual:
|
||||
the native Anthropic `requests` session (and the GigaChat library client) do
|
||||
not carry these socket options.
|
||||
|
||||
The configured-session startup/recovery receipt and every newly minted meaningful wake
|
||||
also carry one host-rendered `coordination_context`: the complete parent-authored advisory
|
||||
`delegation_budget.intent_note`, explicit-deadline time remaining, known/partial/unknown
|
||||
|
|
|
|||
|
|
@ -52,6 +52,13 @@ PACING_INTERVAL_DEFAULT_SEC = 600
|
|||
# supervisor loop STALLED if it has not ticked within this many seconds (healthy tick
|
||||
# ~0.5s), so it only fires on a real wedge. 0 disables.
|
||||
SUPERVISOR_LIVENESS_DEADLINE_DEFAULT_SEC = 90
|
||||
# TCP keepalive tuning for long-lived remote LLM connections: a NAT/VPN mapping
|
||||
# silently dropped during a long silent reasoning stretch is detected by kernel
|
||||
# probes (idle threshold, probe interval, probe count) instead of hanging until
|
||||
# the transport read timeout. Consumed by platform_layer's socket-option builder.
|
||||
TCP_KEEPALIVE_IDLE_SEC = 60
|
||||
TCP_KEEPALIVE_INTERVAL_SEC = 60
|
||||
TCP_KEEPALIVE_PROBE_COUNT = 5
|
||||
|
||||
|
||||
def _guard_live_settings_write() -> None:
|
||||
|
|
|
|||
|
|
@ -1446,13 +1446,36 @@ class LLMClient:
|
|||
target = self._resolve_remote_target("openrouter::")
|
||||
return self._get_remote_client(target)
|
||||
|
||||
@staticmethod
|
||||
def _remote_transport(async_client: bool = False):
|
||||
"""One transport factory for every remote httpx client.
|
||||
|
||||
TCP keepalive probes let the kernel detect a NAT/VPN mapping silently
|
||||
dropped during a long silent reasoning stretch, instead of the socket
|
||||
hanging until the transport read timeout. Options live on the
|
||||
transport (httpx ignores ``socket_options`` on the Client itself).
|
||||
"""
|
||||
import httpx
|
||||
|
||||
from ouroboros.platform_layer import tcp_keepalive_socket_options
|
||||
|
||||
options = tcp_keepalive_socket_options()
|
||||
if async_client:
|
||||
return httpx.AsyncHTTPTransport(socket_options=options)
|
||||
return httpx.HTTPTransport(socket_options=options)
|
||||
|
||||
@staticmethod
|
||||
def _new_remote_client(target: Dict[str, Any]):
|
||||
from openai import OpenAI
|
||||
# DefaultHttpxClient keeps the SDK's own timeout/limits defaults while
|
||||
# letting the shared keepalive transport carry the socket options.
|
||||
from openai import DefaultHttpxClient, OpenAI
|
||||
|
||||
kwargs: Dict[str, Any] = {
|
||||
"api_key": str(target.get("api_key") or ""),
|
||||
"max_retries": 0,
|
||||
"http_client": DefaultHttpxClient(
|
||||
transport=LLMClient._remote_transport()
|
||||
),
|
||||
}
|
||||
base_url = str(target.get("base_url") or "")
|
||||
headers = dict(target.get("default_headers") or {})
|
||||
|
|
@ -1517,11 +1540,14 @@ class LLMClient:
|
|||
|
||||
client = self._async_remote_clients.get(cache_key)
|
||||
if client is None:
|
||||
from openai import AsyncOpenAI
|
||||
from openai import AsyncOpenAI, DefaultAsyncHttpxClient
|
||||
|
||||
kwargs: Dict[str, Any] = {
|
||||
"api_key": api_key,
|
||||
"max_retries": 0,
|
||||
"http_client": DefaultAsyncHttpxClient(
|
||||
transport=self._remote_transport(async_client=True)
|
||||
),
|
||||
}
|
||||
if base_url:
|
||||
kwargs["base_url"] = base_url
|
||||
|
|
@ -1551,6 +1577,7 @@ class LLMClient:
|
|||
trust_env=False,
|
||||
mounts={},
|
||||
timeout=cls._no_proxy_timeout(timeout),
|
||||
transport=cls._remote_transport(),
|
||||
)
|
||||
oa_client = OpenAI(
|
||||
api_key=str(target.get("api_key") or ""),
|
||||
|
|
@ -1570,6 +1597,7 @@ class LLMClient:
|
|||
trust_env=False,
|
||||
mounts={},
|
||||
timeout=cls._no_proxy_timeout(timeout),
|
||||
transport=cls._remote_transport(async_client=True),
|
||||
)
|
||||
oa_client = AsyncOpenAI(
|
||||
api_key=str(target.get("api_key") or ""),
|
||||
|
|
|
|||
|
|
@ -837,6 +837,49 @@ def kill_processes_referencing(marker: str) -> None:
|
|||
force_kill_pid(pid)
|
||||
|
||||
|
||||
def tcp_keepalive_socket_options() -> List[tuple]:
|
||||
"""Cross-platform TCP keepalive options for long-lived remote sockets.
|
||||
|
||||
A NAT/VPN gateway that silently drops an idle connection's mapping leaves
|
||||
the local socket half-open: without keepalive probes the process only
|
||||
learns at the (deliberately long) transport read timeout. Kernel probes
|
||||
detect the dead peer within minutes instead.
|
||||
|
||||
Every platform gets ``SO_KEEPALIVE``; the probe-tuning constants are set
|
||||
only where the platform exposes them (Linux spells the idle threshold
|
||||
``TCP_KEEPIDLE``, Darwin spells it ``TCP_KEEPALIVE``), each behind a
|
||||
``hasattr`` guard so an older interpreter still gets the safe minimum.
|
||||
"""
|
||||
import socket
|
||||
|
||||
from ouroboros.config import (
|
||||
TCP_KEEPALIVE_IDLE_SEC,
|
||||
TCP_KEEPALIVE_INTERVAL_SEC,
|
||||
TCP_KEEPALIVE_PROBE_COUNT,
|
||||
)
|
||||
|
||||
options: List[tuple] = [(socket.SOL_SOCKET, socket.SO_KEEPALIVE, 1)]
|
||||
if IS_LINUX:
|
||||
if hasattr(socket, "TCP_KEEPIDLE"):
|
||||
options.append(
|
||||
(socket.IPPROTO_TCP, socket.TCP_KEEPIDLE, TCP_KEEPALIVE_IDLE_SEC)
|
||||
)
|
||||
if hasattr(socket, "TCP_KEEPINTVL"):
|
||||
options.append(
|
||||
(socket.IPPROTO_TCP, socket.TCP_KEEPINTVL, TCP_KEEPALIVE_INTERVAL_SEC)
|
||||
)
|
||||
if hasattr(socket, "TCP_KEEPCNT"):
|
||||
options.append(
|
||||
(socket.IPPROTO_TCP, socket.TCP_KEEPCNT, TCP_KEEPALIVE_PROBE_COUNT)
|
||||
)
|
||||
elif IS_MACOS:
|
||||
if hasattr(socket, "TCP_KEEPALIVE"):
|
||||
options.append(
|
||||
(socket.IPPROTO_TCP, socket.TCP_KEEPALIVE, TCP_KEEPALIVE_IDLE_SEC)
|
||||
)
|
||||
return options
|
||||
|
||||
|
||||
def kill_process_on_port(port: int) -> None:
|
||||
"""Kill any process listening on the given TCP port."""
|
||||
try:
|
||||
|
|
|
|||
|
|
@ -5,6 +5,7 @@ from __future__ import annotations
|
|||
import json
|
||||
import logging
|
||||
import os
|
||||
import pathlib
|
||||
import re
|
||||
import time
|
||||
import urllib.error
|
||||
|
|
@ -60,8 +61,14 @@ def _gh_api(method: str, path: str, token: str, body: Optional[dict] = None,
|
|||
|
||||
|
||||
def _push_branch(repo_dir: str, branch: str) -> Tuple[bool, str]:
|
||||
# Bounded network push: a hung remote returns the ordinary (False, message)
|
||||
# shape instead of blocking the CI tool on an unbounded subprocess wait.
|
||||
from ouroboros.tools.git import _run_git_network_cmd
|
||||
|
||||
try:
|
||||
result = run_cmd(["git", "push", "-u", "origin", branch], cwd=repo_dir)
|
||||
result = _run_git_network_cmd(
|
||||
["git", "push", "-u", "origin", branch], cwd=pathlib.Path(repo_dir),
|
||||
)
|
||||
return True, result
|
||||
except Exception as e:
|
||||
return False, str(e)
|
||||
|
|
|
|||
|
|
@ -3183,6 +3183,25 @@ def _git_diff(
|
|||
return f"⚠️ GIT_ERROR: {_sanitize_git_error(str(e))}"
|
||||
|
||||
|
||||
def _run_git_network_cmd(cmd: List[str], cwd: pathlib.Path) -> str:
|
||||
"""``run_cmd``-shaped adapter for network git commands (fetch/push).
|
||||
|
||||
Routes through the shared bounded git runner (wall-clock ceiling, process
|
||||
tree kill, ``GIT_TERMINAL_PROMPT=0``, low-speed abort) with the caller's
|
||||
repository as the explicit cwd, so a dead network cannot hang the tool
|
||||
forever. Keeps ``run_cmd``'s contract: stdout on success, ``RuntimeError``
|
||||
with the same message shape on failure.
|
||||
"""
|
||||
from supervisor.update_source import _git_network_bounded
|
||||
|
||||
rc, out, err = _git_network_bounded(list(cmd[1:]), cwd=cwd)
|
||||
if rc != 0:
|
||||
raise RuntimeError(
|
||||
f"Command failed: {' '.join(cmd)}\n\nSTDOUT:\n{out}\n\nSTDERR:\n{err}"
|
||||
)
|
||||
return out
|
||||
|
||||
|
||||
def _ff_pull(repo_dir: pathlib.Path) -> str:
|
||||
try:
|
||||
branch = run_cmd(
|
||||
|
|
@ -3193,7 +3212,7 @@ def _ff_pull(repo_dir: pathlib.Path) -> str:
|
|||
if not branch or branch == "HEAD":
|
||||
return "⚠️ PULL_ERROR: Not on a named branch (detached HEAD). Cannot pull."
|
||||
try:
|
||||
run_cmd(["git", "fetch", "origin"], cwd=repo_dir)
|
||||
_run_git_network_cmd(["git", "fetch", "origin"], cwd=repo_dir)
|
||||
except Exception as e:
|
||||
return f"⚠️ PULL_ERROR: git fetch failed: {_sanitize_git_error(str(e))}"
|
||||
try:
|
||||
|
|
|
|||
|
|
@ -41,6 +41,26 @@ class TelegramTransportError(RuntimeError):
|
|||
"""Telegram could not return a trustworthy API response yet."""
|
||||
|
||||
|
||||
# Shared degraded-transport pacing contract for the poller and notifier loops:
|
||||
# the first retry waits TELEGRAM_RETRY_INITIAL_SEC and each further consecutive
|
||||
# failure doubles the wait monotonically up to TELEGRAM_RETRY_MAX_SEC; any
|
||||
# successful API round resets the wait to the initial value.
|
||||
TELEGRAM_RETRY_INITIAL_SEC = 5
|
||||
TELEGRAM_RETRY_MAX_SEC = 60
|
||||
|
||||
|
||||
def next_telegram_retry_delay(current: float) -> float:
|
||||
"""Next monotone backoff step after one more consecutive transient failure."""
|
||||
return min(max(float(current), TELEGRAM_RETRY_INITIAL_SEC) * 2, TELEGRAM_RETRY_MAX_SEC)
|
||||
|
||||
|
||||
def is_transient_telegram_error(exc: BaseException) -> bool:
|
||||
"""Whether *exc* is a typed transient Telegram failure worth retrying."""
|
||||
if isinstance(exc, TelegramTransportError):
|
||||
return True
|
||||
return isinstance(exc, TelegramRequestRejected) and exc.transient
|
||||
|
||||
|
||||
def _chunk_raw_text(text: str, limit: int = _TELEGRAM_CHUNK_RAW) -> list[str]:
|
||||
"""Split raw text into <=limit-char pieces on line, then space, boundaries."""
|
||||
if len(text) <= limit:
|
||||
|
|
|
|||
|
|
@ -3,12 +3,14 @@ from __future__ import annotations
|
|||
|
||||
import asyncio
|
||||
import json
|
||||
from typing import Any, Dict
|
||||
from typing import Any, Dict, Optional, Tuple
|
||||
|
||||
from .telegram_api import (
|
||||
TELEGRAM_RETRY_INITIAL_SEC,
|
||||
TelegramClient,
|
||||
TelegramRequestRejected,
|
||||
TelegramTransportError,
|
||||
next_telegram_retry_delay,
|
||||
)
|
||||
from .telegram_state import (
|
||||
_data_dir,
|
||||
|
|
@ -44,37 +46,50 @@ def _pinned_chat_id(settings: Dict[str, Any]) -> int:
|
|||
return 0
|
||||
|
||||
|
||||
async def _push_notification(api, chat_id: int, text: str) -> bool:
|
||||
async def _push_notification(
|
||||
api, chat_id: int, text: str,
|
||||
) -> Tuple[str, Optional[BaseException]]:
|
||||
"""Send one notification; never raises a typed Telegram failure.
|
||||
|
||||
Returns ``("sent", None)`` on delivery, ``("transient", exc)`` when the
|
||||
failure is worth retrying next cycle, and ``("skipped", exc)`` on a
|
||||
permanent rejection: the notification is consumed so one dead send cannot
|
||||
replay forever and exhaust the supervised-restart budget.
|
||||
"""
|
||||
protected = api.get_settings(["TELEGRAM_BOT_TOKEN"])
|
||||
client = TelegramClient(protected.get("TELEGRAM_BOT_TOKEN", ""))
|
||||
try:
|
||||
await client.send_message(int(chat_id), text, parse_mode="")
|
||||
return True
|
||||
return "sent", None
|
||||
except TelegramTransportError as exc:
|
||||
api.log("error", f"Telegram notify failed: {exc}")
|
||||
return False
|
||||
return "transient", exc
|
||||
except TelegramRequestRejected as exc:
|
||||
if not exc.transient:
|
||||
raise
|
||||
api.log("error", f"Telegram notify failed: {exc}")
|
||||
return False
|
||||
if exc.transient:
|
||||
api.log("error", f"Telegram notify failed: {exc}")
|
||||
return "transient", exc
|
||||
api.log("error", f"Telegram notification permanently rejected; skipping it: {exc}")
|
||||
return "skipped", exc
|
||||
|
||||
|
||||
_BUDGET_THRESHOLDS = (100, 90, 80) # checked high → low
|
||||
|
||||
|
||||
async def _check_budget_notify(api, settings: Dict[str, Any], chat_id: int, state: Dict[str, Any], lang: str) -> None:
|
||||
async def _check_budget_notify(
|
||||
api, settings: Dict[str, Any], chat_id: int, state: Dict[str, Any], lang: str,
|
||||
) -> Optional[BaseException]:
|
||||
"""Returns the transient send failure, if one occurred, for loop pacing."""
|
||||
if not _notify_enabled(settings, "TELEGRAM_NOTIFY_BUDGET"):
|
||||
return
|
||||
return None
|
||||
snapshot = await _load_runtime_state(api)
|
||||
try:
|
||||
spent = float(snapshot["spent_usd"])
|
||||
total = float(snapshot["budget_limit"])
|
||||
pct = float(snapshot["budget_pct"])
|
||||
except (KeyError, TypeError, ValueError):
|
||||
return
|
||||
return None
|
||||
if total <= 0:
|
||||
return
|
||||
return None
|
||||
crossed = 0
|
||||
for thr in _BUDGET_THRESHOLDS:
|
||||
if pct >= thr:
|
||||
|
|
@ -84,10 +99,15 @@ async def _check_budget_notify(api, settings: Dict[str, Any], chat_id: int, stat
|
|||
if crossed > notified:
|
||||
msg = (f"⚠️ Бюджет: {pct:.0f}% (${spent:.2f} / ${total:.2f})" if lang == "ru"
|
||||
else f"⚠️ Budget: {pct:.0f}% (${spent:.2f} / ${total:.2f})")
|
||||
if await _push_notification(api, chat_id, msg):
|
||||
state["budget_threshold"] = crossed
|
||||
outcome, exc = await _push_notification(api, chat_id, msg)
|
||||
if outcome == "transient":
|
||||
return exc
|
||||
# "sent" delivered it; "skipped" consumes it (permanent rejection) so
|
||||
# the same send does not replay every cycle forever.
|
||||
state["budget_threshold"] = crossed
|
||||
elif crossed < notified:
|
||||
state["budget_threshold"] = crossed # budget raised / spend reset → re-arm
|
||||
return None
|
||||
|
||||
|
||||
def _summary_ids_in_tail(api, limit: int = 200) -> list:
|
||||
|
|
@ -103,15 +123,19 @@ def _summary_ids_in_tail(api, limit: int = 200) -> list:
|
|||
return ids
|
||||
|
||||
|
||||
async def _check_tasks_notify(api, settings: Dict[str, Any], chat_id: int, state: Dict[str, Any], lang: str) -> None:
|
||||
async def _check_tasks_notify(
|
||||
api, settings: Dict[str, Any], chat_id: int, state: Dict[str, Any], lang: str,
|
||||
) -> Optional[BaseException]:
|
||||
"""Returns the last transient send failure, if any, for loop pacing."""
|
||||
if not _notify_enabled(settings, "TELEGRAM_NOTIFY_TASKS"):
|
||||
return
|
||||
return None
|
||||
summaries = _summary_ids_in_tail(api)
|
||||
if "notified_task_ids" not in state:
|
||||
# First run with task notifications on → treat the existing backlog as seen
|
||||
# so enabling the toggle doesn't blast a notification for every old task.
|
||||
state["notified_task_ids"] = [tid for tid, _ in summaries][-300:]
|
||||
return
|
||||
return None
|
||||
transient: Optional[BaseException] = None
|
||||
seen = list(state.get("notified_task_ids") or [])
|
||||
seen_set = set(seen)
|
||||
for tid, e in summaries:
|
||||
|
|
@ -134,16 +158,23 @@ async def _check_tasks_notify(api, settings: Dict[str, Any], chat_id: int, state
|
|||
tail = (" · " + " · ".join(parts)) if parts else ""
|
||||
icon = "✅" if outcome in ("", "completed", "done") else "⚠️"
|
||||
msg = (f"{icon} Задача {tid[:8]} готова{tail}" if lang == "ru" else f"{icon} Task {tid[:8]} done{tail}")
|
||||
if await _push_notification(api, chat_id, msg):
|
||||
seen.append(tid)
|
||||
seen_set.add(tid)
|
||||
send_outcome, exc = await _push_notification(api, chat_id, msg)
|
||||
if send_outcome == "transient":
|
||||
transient = exc
|
||||
continue
|
||||
# "sent" or permanent "skipped": either way this notification is done.
|
||||
seen.append(tid)
|
||||
seen_set.add(tid)
|
||||
state["notified_task_ids"] = seen[-300:]
|
||||
return transient
|
||||
|
||||
|
||||
def _make_notifier(api):
|
||||
"""Periodic, file-based proactive notifications (task done / budget threshold).
|
||||
Read-only over durable files; sends only when a pinned chat + toggle are set."""
|
||||
async def notifier() -> None:
|
||||
retry_delay = TELEGRAM_RETRY_INITIAL_SEC
|
||||
degraded_cause = ""
|
||||
while True:
|
||||
settings = _load_settings(api)
|
||||
chat_id = _pinned_chat_id(settings)
|
||||
|
|
@ -151,11 +182,25 @@ def _make_notifier(api):
|
|||
settings,
|
||||
"TELEGRAM_NOTIFY_BUDGET",
|
||||
)
|
||||
transient: Optional[BaseException] = None
|
||||
if chat_id and want:
|
||||
lang = str(settings.get("TELEGRAM_LANGUAGE") or "en").strip().lower()
|
||||
state = _load_notif_state(api)
|
||||
await _check_budget_notify(api, settings, chat_id, state, lang)
|
||||
await _check_tasks_notify(api, settings, chat_id, state, lang)
|
||||
transient = await _check_budget_notify(api, settings, chat_id, state, lang)
|
||||
transient = await _check_tasks_notify(api, settings, chat_id, state, lang) or transient
|
||||
_save_notif_state(api, state)
|
||||
if transient is not None:
|
||||
# Transition logging only: one line entering degraded, one on
|
||||
# recovery; the monotone backoff paces retries meanwhile.
|
||||
if not degraded_cause:
|
||||
degraded_cause = type(transient).__name__
|
||||
api.log("warning", f"Telegram notifier degraded ({degraded_cause}): {transient}")
|
||||
await asyncio.sleep(retry_delay)
|
||||
retry_delay = next_telegram_retry_delay(retry_delay)
|
||||
continue
|
||||
if degraded_cause:
|
||||
api.log("info", f"Telegram notifier recovered after {degraded_cause}.")
|
||||
degraded_cause = ""
|
||||
retry_delay = TELEGRAM_RETRY_INITIAL_SEC
|
||||
await asyncio.sleep(30)
|
||||
return notifier
|
||||
|
|
|
|||
|
|
@ -15,10 +15,13 @@ from starlette.responses import JSONResponse
|
|||
from ouroboros.utils import truncate_review_artifact
|
||||
|
||||
from .lib.telegram_api import (
|
||||
TELEGRAM_RETRY_INITIAL_SEC,
|
||||
TelegramClient,
|
||||
TelegramRequestRejected,
|
||||
TelegramTransportError,
|
||||
is_transient_telegram_error,
|
||||
markdown_to_telegram_html,
|
||||
next_telegram_retry_delay,
|
||||
_LOCALIZED_TEXTS,
|
||||
)
|
||||
from .lib.telegram_state import (
|
||||
|
|
@ -563,25 +566,39 @@ async def _edit_panel(api, client, chat_id: int, message_id: int, text: str, key
|
|||
await client.send_message_with_inline_keyboard(chat_id, text, keyboard)
|
||||
|
||||
|
||||
async def _validate_bot(api, client, command_mode: str, lang: str) -> None:
|
||||
"""Validate the token (getMe) and finish bridge startup bookkeeping."""
|
||||
await client.call("getMe")
|
||||
# The owner-control commands (evolve/bg/review/restart/panic) only
|
||||
# appear in full_access — they already forward via the raw-command
|
||||
# path; listing them here just makes them discoverable/tappable.
|
||||
try:
|
||||
await client.call("setMyCommands", data={"commands": json.dumps(_bot_commands(command_mode))})
|
||||
api.log("info", "Telegram bot commands configured successfully")
|
||||
except Exception as exc:
|
||||
api.log("warning", f"Failed to set Telegram bot commands: {exc}")
|
||||
|
||||
_save_bridge_status(api, "ready")
|
||||
api.log("info", f"Telegram poller started (command_mode={command_mode}, lang={lang})")
|
||||
|
||||
|
||||
async def _start_poller(api):
|
||||
try:
|
||||
protected_settings = api.get_settings(["TELEGRAM_BOT_TOKEN"])
|
||||
_settings, pinned_chat, max_updates, command_mode, lang = _poller_preferences(api)
|
||||
client = TelegramClient(protected_settings.get("TELEGRAM_BOT_TOKEN", ""))
|
||||
offset = _load_offset(api)
|
||||
await client.call("getMe")
|
||||
# The owner-control commands (evolve/bg/review/restart/panic) only
|
||||
# appear in full_access — they already forward via the raw-command
|
||||
# path; listing them here just makes them discoverable/tappable.
|
||||
try:
|
||||
await client.call("setMyCommands", data={"commands": json.dumps(_bot_commands(command_mode))})
|
||||
api.log("info", "Telegram bot commands configured successfully")
|
||||
await _validate_bot(api, client, command_mode, lang)
|
||||
validated = True
|
||||
except Exception as exc:
|
||||
api.log("warning", f"Failed to set Telegram bot commands: {exc}")
|
||||
|
||||
_save_bridge_status(api, "ready")
|
||||
api.log("info", f"Telegram poller started (command_mode={command_mode}, lang={lang})")
|
||||
return client, offset, pinned_chat, max_updates, command_mode, lang
|
||||
if not is_transient_telegram_error(exc):
|
||||
raise
|
||||
# A dead network at startup must not burn the supervised-restart
|
||||
# budget: stay alive and revalidate from the polling loop once
|
||||
# transport returns. Permanent rejections still raise below.
|
||||
validated = False
|
||||
return client, offset, pinned_chat, max_updates, command_mode, lang, validated
|
||||
except TelegramSettingsError:
|
||||
_save_bridge_status(api, "error", "settings_invalid")
|
||||
api.log("error", "Telegram settings are invalid; owner binding is closed.")
|
||||
|
|
@ -593,11 +610,20 @@ async def _start_poller(api):
|
|||
|
||||
|
||||
async def _poller(api) -> None:
|
||||
client, offset, pinned_chat, max_updates, command_mode, lang = await _start_poller(api)
|
||||
client, offset, pinned_chat, max_updates, command_mode, lang, validated = await _start_poller(api)
|
||||
|
||||
retry_delay = TELEGRAM_RETRY_INITIAL_SEC
|
||||
degraded_cause = ""
|
||||
while True:
|
||||
try:
|
||||
if not validated:
|
||||
await _validate_bot(api, client, command_mode, lang)
|
||||
validated = True
|
||||
updates = await client.get_updates(offset)
|
||||
if degraded_cause:
|
||||
api.log("info", f"Telegram poller recovered after {degraded_cause}; polling resumed.")
|
||||
degraded_cause = ""
|
||||
retry_delay = TELEGRAM_RETRY_INITIAL_SEC
|
||||
if updates:
|
||||
local_settings, pinned_chat, max_updates, command_mode, lang = _poller_preferences(api)
|
||||
|
||||
|
|
@ -857,16 +883,19 @@ async def _poller(api) -> None:
|
|||
if updates:
|
||||
_save_offset(api, offset)
|
||||
await asyncio.sleep(0.1)
|
||||
except TelegramRequestRejected as exc:
|
||||
if not exc.transient:
|
||||
except (TelegramRequestRejected, TelegramTransportError) as exc:
|
||||
if isinstance(exc, TelegramRequestRejected) and not exc.transient:
|
||||
_save_bridge_status(api, "error", "telegram_rejected")
|
||||
api.log("error", "Telegram polling was permanently rejected.")
|
||||
raise
|
||||
api.log("warning", f"Telegram poller transient error: {exc}")
|
||||
await asyncio.sleep(5)
|
||||
except TelegramTransportError as exc:
|
||||
api.log("warning", f"Telegram poller transient error: {exc}")
|
||||
await asyncio.sleep(5)
|
||||
# Transition logging only: one line entering degraded, one on
|
||||
# recovery; consecutive transient failures stay quiet while the
|
||||
# monotone backoff (reset by any successful poll) paces retries.
|
||||
if not degraded_cause:
|
||||
degraded_cause = type(exc).__name__
|
||||
api.log("warning", f"Telegram poller degraded ({degraded_cause}): {exc}")
|
||||
await asyncio.sleep(retry_delay)
|
||||
retry_delay = next_telegram_retry_delay(retry_delay)
|
||||
except TelegramSettingsError:
|
||||
_save_bridge_status(api, "error", "settings_invalid")
|
||||
api.log("error", "Telegram settings are invalid; owner binding is closed.")
|
||||
|
|
|
|||
|
|
@ -1975,18 +1975,23 @@ def _configure_credential_helper(repo_slug: str, token: str) -> None:
|
|||
|
||||
|
||||
def push_to_remote(branch: Optional[str] = None, push_tags: bool = True) -> Tuple[bool, str]:
|
||||
"""Push current branch (and optionally tags) to origin."""
|
||||
"""Push current branch (and optionally tags) to origin.
|
||||
|
||||
Network pushes ride the shared bounded git runner: a hung remote surfaces
|
||||
as the ordinary ``(False, "git push failed: ...")`` result instead of
|
||||
pinning the worker forever on an unbounded subprocess wait.
|
||||
"""
|
||||
if not _has_remote("origin"):
|
||||
return False, "No remote configured"
|
||||
|
||||
target = branch or BRANCH_DEV
|
||||
rc, out, err = git_capture(["git", "push", "-u", "origin", target])
|
||||
rc, out, err = _git_network_bounded(["push", "-u", "origin", target])
|
||||
if rc != 0:
|
||||
return False, f"git push failed: {err}"
|
||||
|
||||
result = f"Pushed {target} to origin"
|
||||
if push_tags:
|
||||
rc_t, _, err_t = git_capture(["git", "push", "origin", "--tags"])
|
||||
rc_t, _, err_t = _git_network_bounded(["push", "origin", "--tags"])
|
||||
if rc_t != 0:
|
||||
result += f" (tags push failed: {err_t})"
|
||||
else:
|
||||
|
|
|
|||
|
|
@ -3,6 +3,7 @@
|
|||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import pathlib
|
||||
import re
|
||||
import subprocess
|
||||
from typing import Callable, List, Optional, Tuple
|
||||
|
|
@ -44,18 +45,33 @@ def official_ref_has_constitution(
|
|||
|
||||
|
||||
def _git_network_bounded(
|
||||
args: List[str], *, timeout: Optional[float] = None
|
||||
args: List[str],
|
||||
*,
|
||||
timeout: Optional[float] = None,
|
||||
cwd: os.PathLike[str] | str | None = None,
|
||||
) -> Tuple[int, str, str]:
|
||||
"""Run one non-interactive git network command under the shared ceiling."""
|
||||
"""Run one non-interactive git network command under the shared ceiling.
|
||||
|
||||
``cwd`` selects the repository the command runs in; the default remains the
|
||||
managed system repository. An explicit ``cwd`` must be an existing
|
||||
directory — network commands against a caller-selected repository must
|
||||
never silently fall back to the system repository.
|
||||
"""
|
||||
from ouroboros.update_channels import get_managed_update_fetch_timeout_sec
|
||||
from supervisor import git_ops as _g
|
||||
|
||||
if cwd is None:
|
||||
repo_dir = _g.REPO_DIR
|
||||
else:
|
||||
repo_dir = pathlib.Path(cwd)
|
||||
if not repo_dir.is_dir():
|
||||
return 127, "", f"git cwd is not a directory: {repo_dir}"
|
||||
limit = float(timeout or get_managed_update_fetch_timeout_sec())
|
||||
env = {**os.environ, "LC_ALL": "C", "LANG": "C", "GIT_TERMINAL_PROMPT": "0"}
|
||||
rc, out, err = _g._run_git_process_bounded(
|
||||
["git", "-c", "http.lowSpeedLimit=1024", "-c", "http.lowSpeedTime=30", *args],
|
||||
timeout=limit,
|
||||
cwd=_g.REPO_DIR,
|
||||
cwd=repo_dir,
|
||||
env=env,
|
||||
text=True,
|
||||
)
|
||||
|
|
|
|||
|
|
@ -38,20 +38,20 @@ def test_push_to_remote_push_tags_compatibility(monkeypatch):
|
|||
monkeypatch.setattr(git_ops, "_has_remote", lambda _name: True)
|
||||
monkeypatch.setattr(
|
||||
git_ops,
|
||||
"git_capture",
|
||||
"_git_network_bounded",
|
||||
lambda command, **_kwargs: commands.append(list(command)) or (0, "", ""),
|
||||
)
|
||||
|
||||
ok, _ = git_ops.push_to_remote("feature", push_tags=False)
|
||||
assert ok is True
|
||||
assert commands == [["git", "push", "-u", "origin", "feature"]]
|
||||
assert commands == [["push", "-u", "origin", "feature"]]
|
||||
|
||||
commands.clear()
|
||||
ok, _ = git_ops.push_to_remote("feature", push_tags=True)
|
||||
assert ok is True
|
||||
assert commands == [
|
||||
["git", "push", "-u", "origin", "feature"],
|
||||
["git", "push", "origin", "--tags"],
|
||||
["push", "-u", "origin", "feature"],
|
||||
["push", "origin", "--tags"],
|
||||
]
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -573,17 +573,31 @@ def test_typed_payment_required_maps_to_no_credits():
|
|||
|
||||
def test_ephemeral_openai_client_disables_sdk_retries(monkeypatch):
|
||||
captured = {}
|
||||
http_clients = []
|
||||
|
||||
class FakeOpenAI:
|
||||
def __init__(self, **kwargs):
|
||||
captured.update(kwargs)
|
||||
|
||||
class FakeDefaultHttpxClient:
|
||||
def __init__(self, **kwargs):
|
||||
http_clients.append(kwargs)
|
||||
|
||||
import sys
|
||||
|
||||
monkeypatch.setitem(sys.modules, "openai", types.SimpleNamespace(OpenAI=FakeOpenAI))
|
||||
monkeypatch.setitem(
|
||||
sys.modules,
|
||||
"openai",
|
||||
types.SimpleNamespace(
|
||||
OpenAI=FakeOpenAI, DefaultHttpxClient=FakeDefaultHttpxClient,
|
||||
),
|
||||
)
|
||||
LLMClient._new_remote_client({
|
||||
"api_key": "key", "base_url": "https://example.test/v1", "default_headers": {"X": "Y"},
|
||||
})
|
||||
http_client = captured.pop("http_client")
|
||||
assert isinstance(http_client, FakeDefaultHttpxClient)
|
||||
assert len(http_clients) == 1 and http_clients[0].get("transport") is not None
|
||||
assert captured == {
|
||||
"api_key": "key", "base_url": "https://example.test/v1",
|
||||
"default_headers": {"X": "Y"}, "max_retries": 0,
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue