Skip to content

Commit a3ee459

Browse files
linfeng.wanglinfeng.wang
authored andcommitted
Delete stopped session topics automatically
1 parent fb56c1c commit a3ee459

2 files changed

Lines changed: 191 additions & 1 deletion

File tree

‎src/telegram_codex_controller/console.py‎

Lines changed: 105 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -72,6 +72,7 @@ async def start(self, application: Application) -> None:
7272
self._application = application
7373
self._prune_expired_records()
7474
self._restore_reply_routes()
75+
self._queue_existing_stopped_session_cleanups()
7576
if self._stuck_task is None:
7677
self._stuck_task = asyncio.create_task(self._stuck_watcher())
7778
if self._pin_task is None:
@@ -176,7 +177,11 @@ async def update_status(self, spec: SessionStatusSpec) -> None:
176177
status_text = self._render_status_text(record)
177178
last_text = record.get("status_text")
178179
record["_pending_status_text"] = status_text
179-
if self._session_has_pending_work(record):
180+
if spec.state == "stopped":
181+
self._dirty_session_routes.discard(spec.route.label)
182+
self._pending_status_due_monotonic.pop(spec.route.label, None)
183+
self._queue_stopped_session_cleanup(record)
184+
elif self._session_has_pending_work(record):
180185
self._mark_session_dirty(
181186
spec.route.label,
182187
record,
@@ -354,6 +359,7 @@ def clear_session(self, route: ReplyRoute) -> None:
354359
self._topic_bump_reason_by_route.pop(route.label, None)
355360
self._topic_bump_queue_seen.discard(route.label)
356361
self._last_status_render_monotonic.pop(route.label, None)
362+
self._mark_index_dirty(immediate=True, urgent=True)
357363
self._persist()
358364

359365
def lookup_route_by_message(self, chat_id: int, message_id: int) -> ReplyRoute | None:
@@ -953,6 +959,10 @@ async def _flush_once(self) -> None:
953959
if handled:
954960
return
955961

962+
handled = await self._flush_stopped_session_cleanup()
963+
if handled:
964+
return
965+
956966
handled = await self._flush_stale_topic_cleanup()
957967
if handled:
958968
return
@@ -1397,6 +1407,100 @@ def _mark_session_dirty(
13971407
self._pending_status_due_monotonic[route_label] = due
13981408
self._dirty_session_routes.add(route_label)
13991409

1410+
def _queue_existing_stopped_session_cleanups(self) -> None:
1411+
for record in self._state.get("sessions", {}).values():
1412+
if record.get("state") == "stopped":
1413+
self._queue_stopped_session_cleanup(record)
1414+
1415+
def _queue_stopped_session_cleanup(self, record: dict[str, Any]) -> None:
1416+
if record.get("_pending_topic_delete"):
1417+
return
1418+
route_kind = record.get("route_kind")
1419+
route_target = record.get("route_target")
1420+
if not route_kind or not route_target:
1421+
return
1422+
record["_pending_topic_delete"] = {
1423+
"route_kind": route_kind,
1424+
"route_target": route_target,
1425+
"topic_id": record.get("topic_id"),
1426+
"status_thread_id": record.get("status_thread_id"),
1427+
"topic_managed": bool(record.get("topic_managed", True)),
1428+
"chat_id": record.get("status_chat_id") or self._console_chat_id(),
1429+
"status_message_id": record.get("status_message_id"),
1430+
"topic_bump_message_id": record.get("topic_bump_message_id"),
1431+
"last_event_message_id": record.get("last_event_message_id"),
1432+
"label": record.get("label") or route_target,
1433+
}
1434+
1435+
async def _flush_stopped_session_cleanup(self) -> bool:
1436+
if self._application is None:
1437+
return False
1438+
for route_label, record in list(self._state.get("sessions", {}).items()):
1439+
pending = record.get("_pending_topic_delete")
1440+
if not pending or record.get("state") != "stopped":
1441+
continue
1442+
1443+
route = ReplyRoute(kind=pending["route_kind"], target=pending["route_target"])
1444+
chat_id = int(pending.get("chat_id") or self._console_chat_id())
1445+
status_message_id = pending.get("status_message_id")
1446+
topic_bump_message_id = pending.get("topic_bump_message_id")
1447+
event_message_id = pending.get("last_event_message_id")
1448+
topic_ids = []
1449+
for candidate in [pending.get("topic_id"), pending.get("status_thread_id")]:
1450+
if candidate and candidate not in topic_ids:
1451+
topic_ids.append(candidate)
1452+
topic_managed = bool(pending.get("topic_managed", True))
1453+
now = time.monotonic()
1454+
1455+
try:
1456+
for message_id in [status_message_id, topic_bump_message_id, event_message_id]:
1457+
if not message_id:
1458+
continue
1459+
with contextlib.suppress(BadRequest):
1460+
await self._application.bot.delete_message(
1461+
chat_id=chat_id,
1462+
message_id=int(message_id),
1463+
)
1464+
if topic_ids and self.forum_enabled():
1465+
if topic_managed:
1466+
for topic_id in topic_ids:
1467+
with contextlib.suppress(BadRequest):
1468+
await self._application.bot.delete_forum_topic(
1469+
chat_id=chat_id,
1470+
message_thread_id=int(topic_id),
1471+
)
1472+
else:
1473+
for topic_id in topic_ids:
1474+
with contextlib.suppress(BadRequest):
1475+
await self._application.bot.edit_forum_topic(
1476+
chat_id=chat_id,
1477+
message_thread_id=int(topic_id),
1478+
name=_archived_topic_title(
1479+
{
1480+
"label": pending.get("label") or pending.get("route_target") or "?",
1481+
"topic_title": record.get("topic_title") or "",
1482+
}
1483+
),
1484+
)
1485+
with contextlib.suppress(BadRequest):
1486+
await self._application.bot.close_forum_topic(
1487+
chat_id=chat_id,
1488+
message_thread_id=int(topic_id),
1489+
)
1490+
self._record_write(now, state="stopped", urgent=True)
1491+
self.clear_session(route)
1492+
LOG.info("Deleted stopped session topic for %s", route_label)
1493+
return True
1494+
except RetryAfter as exc:
1495+
self._enter_write_backoff(float(exc.retry_after) + 1)
1496+
LOG.info("Stopped session cleanup backed off for %.1fs", float(exc.retry_after))
1497+
return False
1498+
except TelegramError:
1499+
LOG.exception("Failed to delete stopped session topic for %s", route_label)
1500+
self.clear_session(route)
1501+
return False
1502+
return False
1503+
14001504
def _queue_stale_topic_cleanup(self, record: dict[str, Any], *, new_topic_id: int | None) -> None:
14011505
old_topic_id = record.get("topic_id")
14021506
if old_topic_id is None or old_topic_id == new_topic_id:

‎tests/test_console.py‎

Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,9 @@ async def edit_forum_topic(self, **kwargs: object) -> None:
3333
async def close_forum_topic(self, **kwargs: object) -> None:
3434
self.calls.append(("close_forum_topic", kwargs))
3535

36+
async def delete_forum_topic(self, **kwargs: object) -> None:
37+
self.calls.append(("delete_forum_topic", kwargs))
38+
3639

3740
def make_settings(root: Path) -> Settings:
3841
sidecar = root / "sidecar" / "runner.mjs"
@@ -389,6 +392,89 @@ def test_flush_stale_topic_cleanup_archives_previous_topic(self) -> None:
389392
],
390393
)
391394

395+
def test_stopped_transition_queues_topic_delete_cleanup(self) -> None:
396+
with tempfile.TemporaryDirectory() as temp_dir:
397+
settings = make_settings(Path(temp_dir))
398+
manager = TelegramConsoleManager(settings)
399+
manager._application = FakeApplication(bot=FakeBot())
400+
manager.set_forum_enabled(True)
401+
route = ReplyRoute("mirror", "ttys039")
402+
manager._state["sessions"][route.label] = {
403+
"route_kind": "mirror",
404+
"route_target": "ttys039",
405+
"kind": "mirror",
406+
"label": "ttys039",
407+
"title": "Terminal",
408+
"state": "running",
409+
"step": "mirroring",
410+
"summary": "still working",
411+
"updated_at": "2026-03-15T21:00:00+00:00",
412+
"status_message_id": 11,
413+
"status_chat_id": -1001,
414+
"topic_id": 42,
415+
"topic_title": "🟢 ttys039",
416+
"topic_bump_message_id": 12,
417+
}
418+
419+
asyncio.run(
420+
manager.update_status(
421+
SessionStatusSpec(
422+
route=route,
423+
kind="mirror",
424+
label="ttys039",
425+
title="Terminal",
426+
state="stopped",
427+
step="session ended",
428+
summary="Mirror source disappeared.",
429+
)
430+
)
431+
)
432+
433+
pending = manager._state["sessions"][route.label].get("_pending_topic_delete")
434+
self.assertIsNotNone(pending)
435+
self.assertEqual(pending["topic_id"], 42)
436+
self.assertEqual(pending["status_message_id"], 11)
437+
438+
def test_flush_stopped_session_cleanup_deletes_managed_topic_and_clears_session(self) -> None:
439+
with tempfile.TemporaryDirectory() as temp_dir:
440+
settings = make_settings(Path(temp_dir))
441+
bot = FakeBot()
442+
manager = TelegramConsoleManager(settings)
443+
manager._application = FakeApplication(bot=bot)
444+
manager.set_forum_enabled(True)
445+
manager._state["sessions"]["mirror:ttys039"] = {
446+
"route_kind": "mirror",
447+
"route_target": "ttys039",
448+
"label": "ttys039",
449+
"kind": "mirror",
450+
"state": "stopped",
451+
"_pending_topic_delete": {
452+
"route_kind": "mirror",
453+
"route_target": "ttys039",
454+
"topic_id": 42,
455+
"topic_managed": True,
456+
"chat_id": -1001,
457+
"status_message_id": 11,
458+
"topic_bump_message_id": 12,
459+
"last_event_message_id": 13,
460+
"label": "ttys039",
461+
},
462+
}
463+
464+
handled = asyncio.run(manager._flush_stopped_session_cleanup())
465+
466+
self.assertTrue(handled)
467+
self.assertNotIn("mirror:ttys039", manager._state["sessions"])
468+
self.assertEqual(
469+
bot.calls,
470+
[
471+
("delete_message", {"chat_id": -1001, "message_id": 11}),
472+
("delete_message", {"chat_id": -1001, "message_id": 12}),
473+
("delete_message", {"chat_id": -1001, "message_id": 13}),
474+
("delete_forum_topic", {"chat_id": -1001, "message_thread_id": 42}),
475+
],
476+
)
477+
392478

393479
if __name__ == "__main__":
394480
unittest.main()

0 commit comments

Comments
 (0)