diff --git a/CHANGELOG.md b/CHANGELOG.md index e2ccb6f..2735a90 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -27,6 +27,12 @@ secure remote command runner; Claude Code becomes an optional assistant layer. free text. Admin list (`TELEGRAM_ALLOWED_USERS`) is unchanged. - **Long output arrives as a file** instead of a flood of 4000-char chunks (threshold ~8k chars), with a chunked-send fallback if the upload fails. +- **Delivery-lag handling.** Telegram's server-side send timestamp now flows + through the transport: the receipt log shows real inbound lag + (`(sent 58s ago)`), /confirm is judged by when the user *sent* it rather + than when a flaky uplink finally delivered it, and the long-poll window is + shortened (25s/35s) so a silently dead poll can't sit on incoming messages + for 75 seconds. - **Group hardening:** authorization now checks the message *sender* in group chats, not just the chat id — allow-listing a group no longer hands the command runner to every member. diff --git a/gateway.py b/gateway.py index 968e409..881a2ef 100755 --- a/gateway.py +++ b/gateway.py @@ -428,7 +428,8 @@ def launch_detached(script: str) -> bool: return True -def handle(chat_id: str, actor: str, text: str, guest: bool) -> None: +def handle(chat_id: str, actor: str, text: str, guest: bool, + sent_ts: float = 0.0) -> None: stripped = text.strip() parts = stripped.split(maxsplit=1) cmd = parts[0].lower() if parts else "" @@ -448,7 +449,10 @@ def handle(chat_id: str, actor: str, text: str, guest: bool) -> None: send(chat_id, "Nothing awaiting confirmation.") return expiry, pcmd, pargv = pending - if time.time() > expiry: + # Judge by when the user SENT /confirm (transport's server timestamp), + # not when it finally reached us — delivery lag on a flaky uplink must + # not burn the confirmation window. + if (sent_ts or time.time()) > expiry: audit("confirm-expired", actor, pcmd) send(chat_id, f"⌛ {pcmd} expired unconfirmed — send it again if you still want it.") return @@ -548,13 +552,13 @@ def handle(chat_id: str, actor: str, text: str, guest: bool) -> None: send(chat_id, reply) -def _worker(chat_id: str, actor: str, text: str, guest: bool) -> None: +def _worker(chat_id: str, actor: str, text: str, guest: bool, sent_ts: float) -> None: """Handle one message in its own thread, then release the per-chat slot.""" # Log completion + duration: the receipt line alone can't distinguish "still # running", "killed mid-flight", and "reply delayed by a Telegram stall". t0 = time.monotonic() try: - handle(chat_id, actor, text, guest) + handle(chat_id, actor, text, guest, sent_ts) except Exception as e: # noqa: BLE001 log("handler error:", e) send(chat_id, "⚠️ Something broke handling that. Logged it.") @@ -602,7 +606,13 @@ def main() -> None: send(chat_id, f"⛔ Not authorized. Your ID is `{actor}` — " "add it to TELEGRAM_ALLOWED_USERS to enable access.") continue - log(f"[{chat_id}] {text[:80]}") + # Surface inbound delivery lag (send time is Telegram's server clock): + # a message that sat queued behind a dead long-poll looks like a slow + # bot otherwise. Small offsets are clock noise, only log real lag. + sent_ts = float(upd.get("date") or 0) + lag = (time.time() - sent_ts) if sent_ts else 0.0 + log(f"[{chat_id}] {text[:80]}" + + (f" (sent {lag:.0f}s ago)" if lag > 5 else "")) # Concurrency guard: one in-flight message per chat. Drop a second one # (with a notice) rather than overlapping runs. Handle in a thread so a # long run in one chat doesn't block polling or other chats. @@ -614,7 +624,7 @@ def main() -> None: send(chat_id, "⏳ Still working on your previous message — give me a moment, " "then resend if needed.") continue - threading.Thread(target=_worker, args=(chat_id, actor, text, guest), + threading.Thread(target=_worker, args=(chat_id, actor, text, guest, sent_ts), daemon=True).start() diff --git a/transport_telegram.py b/transport_telegram.py index 45eba6e..7b5f70c 100644 --- a/transport_telegram.py +++ b/transport_telegram.py @@ -95,17 +95,25 @@ class TelegramTransport: """Long-poll getUpdates forever, yielding one dict per text message. chat_id is where to reply; from_id is who wrote it (differs from chat_id - in groups — the gateway authorizes on the sender there). Transport errors - are logged and retried here; the gateway never sees them. + in groups — the gateway authorizes on the sender there); date is Telegram's + server-side send timestamp (unix), which lets the gateway measure delivery + lag and judge time-sensitive replies (/confirm) by when the user actually + sent them. Transport errors are logged and retried here; the gateway never + sees them. + + The poll window is deliberately short-ish (25s server / 35s client): when a + long-poll connection dies silently (flaky uplink, NAT timeout), incoming + messages queue server-side until the client times out and re-polls — the + window is the worst-case delivery lag, so keep it tolerable. """ offset = 0 while True: try: - resp = self._call("getUpdates", {"offset": offset, "timeout": 50}, - timeout=70) + resp = self._call("getUpdates", {"offset": offset, "timeout": 25}, + timeout=35) except Exception as e: # noqa: BLE001 log("getUpdates error:", e) - time.sleep(5) + time.sleep(3) continue for upd in resp.get("result", []): offset = upd["update_id"] + 1 @@ -116,6 +124,7 @@ class TelegramTransport: "chat_id": str(msg["chat"]["id"]), "from_id": str((msg.get("from") or {}).get("id") or ""), "text": msg["text"], + "date": int(msg.get("date") or 0), } # -----------------------------------------------------------------------