"""Observe configuration publications and recover from a slow subscription.

Run ``uv run python examples/events.py``. The client starts with no sources,
so this example needs no network service or credentials. Empty groups illustrate
event ordering; their snapshots correctly remain unavailable for validation.
"""

import argparse
from dataclasses import replace

import anyio

from rpkiparrot import Client
from rpkiparrot.config import ClientConfig, Limits, SourceGroupConfig
from rpkiparrot.errors import ResyncRequiredError
from rpkiparrot.models import SnapshotEvent


async def main() -> None:
    """Watch an atomic initial point, process an event, and explicitly resync."""
    config = ClientConfig(limits=Limits(subscription_queue_size=1))
    async with Client(config) as client, anyio.create_task_group() as tasks:
        await tasks.start(client.run)
        async with client.watch() as subscription:
            assert subscription.initial_status is not None
            assert subscription.initial_status.snapshot_id == subscription.initial.id
            initial = subscription.initial
            config = replace(config, groups=(SourceGroupConfig(id="primary", priority=0),))
            first = await client.apply_config(config, expected_revision=initial.config_revision)
            with anyio.fail_after(5):
                event = await anext(subscription)
            assert isinstance(event, SnapshotEvent)
            assert event.before_id == initial.id
            assert event.after_id == first.snapshot_id

            # Deliberately lag behind a queue of one, while the publisher keeps
            # working. A lost event is an explicit terminal resync condition.
            second = await client.apply_config(
                replace(config, groups=()), expected_revision=first.config_revision
            )
            third = await client.apply_config(config, expected_revision=second.config_revision)
            try:
                await anext(subscription)
            except ResyncRequiredError:
                pass
            else:
                raise AssertionError("a slow subscription must request resynchronization")

        async with client.watch() as fresh:
            assert fresh.initial.id == third.snapshot_id
        await client.aclose()
    print("Observed atomic configuration events and recovered from explicit ResyncRequired.")


if __name__ == "__main__":
    parser = argparse.ArgumentParser(description=__doc__)
    parser.add_argument("--backend", choices=("asyncio", "trio"), default="asyncio")
    anyio.run(main, backend=parser.parse_args().backend)
