Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
29 changes: 29 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,35 @@ codes documented in `AGENT.md` are the public API.

### Fixed

- **Gateway jobs work again, and incoming updates reach the daemon at all.**
Four faults stacked up, and each one alone was enough to stop every job:

- The daemon attached its update handler before the account's client
existed, so no incoming update reached the event bus. Jobs, webhooks and
`watch` saw the daemon's own sends and nothing else. Every job stopped
with `matched=0 skipped=0`.
- The bus hands a job the raw TL update, but filters and actions read a
Telethon event (`chat_id`, `is_private`, `reply()`). The job now builds
that event the way a Telethon handler would, including
`NewMessage(incoming=True)`, so an auto-reply never answers the
account's own messages.
- `filter_chat_id` skipped `@name` refs, expecting the gateway to resolve
them, and nothing did, so `chat_id: "@channel"` never matched. Jobs now
resolve them at start, and retry once a minute if the account was
offline.
- A restarted daemon started no jobs until `tlgr job reload`. It now
starts them at boot.

- **Accounts no longer go "degraded" on every start** with `AttributeError:
'MessageBox' object attribute 'apply_difference' is read-only`. Telethon's
`MessageBox` declares `__slots__`, so the too-long hook now installs on a
per-client subclass, and the account no longer burns a reconnect or runs
without the hook.

- **`job list` reports `enabled` and `running` even when they are at their
defaults.** Before, an enabled job that was not running showed `-` in both
columns, the same as a disabled one.

- **`chat list` no longer loses a block of dialogs at a page boundary.** The
dialog walk paired each row with its top message by message id alone, but
only private chats share one id space: every channel and supergroup numbers
Expand Down
11 changes: 11 additions & 0 deletions tests/fake_telethon.py
Original file line number Diff line number Diff line change
Expand Up @@ -1429,6 +1429,17 @@ def __init__(self, world: World | None = None, session: Any = None, **_: Any) ->
self._handlers: list[tuple[Any, Any]] = []
self.disconnected: asyncio.Future[None] = asyncio.get_event_loop().create_future()
self._requests = 0
# What Telethon's events and messages read off a client: `_set_client`
# looks entities up here, `build()` takes `_self_id`, and
# `Message.text` unparses through `parse_mode` (None means plain text).
self.parse_mode = None
from telethon._updates import EntityCache

self._mb_entity_cache = EntityCache(self_id=self.world.me.id, self_bot=False)

@property
def _self_id(self) -> int:
return int(self.world.me.id)

# -- connection --------------------------------------------------------

Expand Down
70 changes: 70 additions & 0 deletions tests/test_daemon_jobs.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,70 @@
"""The daemon starts the jobs in `jobs.yaml` when it boots.

Until 2.0.1 it did not: `reload_jobs()` ran only on `tlgr job reload`, so a
restarted daemon (launchd, an upgrade, a crash) ran no jobs at all and said
nothing about it.
"""

from __future__ import annotations

import asyncio
from pathlib import Path

from fake_telethon import fake_client_factory

JOBS = """\
jobs:
- name: archive
account: {account}
actions:
- forward:
to: ["@archive"]
- name: parked
account: {account}
enabled: false
actions:
- reply: "parked"
"""


async def _until(predicate, timeout: float = 5.0) -> None:
deadline = asyncio.get_running_loop().time() + timeout
while not predicate():
if asyncio.get_running_loop().time() > deadline:
raise AssertionError("condition never became true")
await asyncio.sleep(0.02)


class TestJobsAtBoot:
async def test_enabled_jobs_are_running_once_the_daemon_is_up(
self, tlgr_home: Path, stub_account: str, world
):
from tlgr.daemon.app import Daemon

(tlgr_home / "jobs.yaml").write_text(JOBS.format(account=stub_account))
daemon = Daemon(tlgr_home, client_factory=fake_client_factory(world))
runner = asyncio.create_task(daemon.run())
try:
await _until(lambda: any(job.get("running") for job in daemon.list_jobs()))
jobs = {job["name"]: job for job in daemon.list_jobs()}
assert jobs["archive"]["running"] is True
assert "parked" not in jobs
finally:
daemon.request_shutdown()
await asyncio.wait_for(runner, timeout=10)

async def test_a_broken_jobs_file_does_not_stop_the_daemon(
self, tlgr_home: Path, stub_account: str, world
):
from tlgr.daemon.app import Daemon

(tlgr_home / "jobs.yaml").write_text("jobs: [unclosed")
daemon = Daemon(tlgr_home, client_factory=fake_client_factory(world))
runner = asyncio.create_task(daemon.run())
try:
await _until(lambda: daemon.ready)
await _until(lambda: daemon._jobs_task is not None and daemon._jobs_task.done())
assert daemon.list_jobs() == []
finally:
daemon.request_shutdown()
await asyncio.wait_for(runner, timeout=10)
63 changes: 63 additions & 0 deletions tests/test_daemon_updates.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
"""Incoming updates reach the daemon's event bus.

Until 2.0.1 none did: the daemon attached its update handler straight after
`AccountSession.start()`, which only schedules the supervisor, so it always
found `session.client is None` and returned. Jobs, webhooks and `watch` saw
the daemon's own sends and nothing else, and every job stopped with
`matched=0 skipped=0`.
"""

from __future__ import annotations

import asyncio

from telethon.tl import types

from tlgr.daemon.session import SessionState


async def _online(daemon, alias: str):
session = await daemon.sessions.ensure(alias)
deadline = asyncio.get_running_loop().time() + 2.0
while session.state != SessionState.ONLINE:
if asyncio.get_running_loop().time() > deadline:
raise AssertionError(f"{alias} never came online: {session.state}")
await asyncio.sleep(0.01)
return session


def _incoming(message_id: int = 7) -> types.UpdateNewChannelMessage:
message = types.Message(
id=message_id,
peer_id=types.PeerChannel(channel_id=1234),
date=None,
message="hello",
)
return types.UpdateNewChannelMessage(message=message, pts=1, pts_count=1)


class TestIncomingUpdates:
async def test_an_incoming_message_is_numbered_on_the_bus(self, daemon, stub_account):
session = await _online(daemon, stub_account)
before = daemon.bus.latest_seq(stub_account)

await session.client.feed(_incoming())

assert daemon.bus.latest_seq(stub_account) == before + 1
[envelope], _gap = daemon.bus.replay(stub_account, before)
assert envelope.type == "message_new"
assert envelope.chat_id == -1000000001234

async def test_the_handler_is_attached_once_across_reconnects(self, daemon, stub_account):
"""The client outlives a reconnect, so a second handler would double every event."""
session = await _online(daemon, stub_account)
handlers = len(session.client._handlers)

await session.client.disconnect()
deadline = asyncio.get_running_loop().time() + 5.0
while session.reconnects == 0 or session.state != SessionState.ONLINE:
if asyncio.get_running_loop().time() > deadline:
raise AssertionError("the session never reconnected")
await asyncio.sleep(0.01)

assert len(session.client._handlers) == handlers
Loading
Loading