Skip to content

Keep a single EventManager per process (event loop) #2242

Description

@Mantisus

Make EventManager a single service per process (per event loop, in practice), owned by the global service_locator.

Motivation

Today it's possible to create several EventManager instances that emit the same events. BasicCrawler builds its own ServiceLocator and may put a crawler-specific event manager into it, while Snapshotter and RecoverableState always resolve the global service_locator. So a crawler-scoped manager runs a second SystemInfo sampling loop and emits PersistState that nobody listens to, while the components it was supposed to serve keep listening on the global one. Working around that is why BasicCrawler._run_crawler currently enters both managers.

The second problem is lifetime. An event manager only runs while at least one crawler is active, but several systems that depend on it are usable without a crawler. FileSystemRequestQueueClient keeps its ordering and in-progress state in RecoverableState, and the same holds for autosaved KeyValueStore values, RequestList and SitemapRequestLoader. Used standalone, they register a PersistState listener on a manager that is never started, so their state is only written on explicit teardown.

Implementation

Keep exactly one EventManager, reachable through a single access point in service_locator, resolved per running event loop so that a manager never leaks from one asyncio.run into the next. service_locator.set_event_manager stays as the way to install a custom manager, for tests or for the Apify SDK's platform event manager.

The manager starts itself on the first on or emit call and closes itself when its event loop shuts down, so it works regardless of whether a crawler is running. There is no context to enter and no async with in the calling code. All current on call sites are already inside async code, and so is every RecurringTask.start call site, so the lazy start always has a running loop to attach to.

Activity

  1. added
    t-toolingIssues with this label are in the ownership of the tooling team.
    on Sep 19, 2026
  2. added this to the 2.0 milestone on Sep 19, 2026
  3. janbuchar commented on Sep 21, 2026

    @janbuchar
    Collaborator

    Hey @Mantisus, thanks for filing this. We'll probably want to do this right in crawlee-js 4.0. So are you saying that we should simply not allow passing the event_manager option to the BasicCrawler constructor?

  4. Mantisus commented on Sep 21, 2026

    @Mantisus
    CollaboratorAuthor

    So are you saying that we should simply not allow passing the event_manager option to the BasicCrawler constructor?

    Yes, and the same goes for the other places that accept it today, SessionPool, for example. The other half is decoupling the EventManager lifecycle from the crawler. Today, it only runs while its context is entered, and in practice only the crawler does that.

    So that users don't have to start the event manager themselves, I'd propose lazy initialization: the manager starts on first use and keeps running until the asyncio event loop shuts down or the user closes it explicitly.

  5. janbuchar commented on Sep 22, 2026

    @janbuchar
    Collaborator

    Lazy initialization and async don't play well together, do you have something concrete in mind?

  6. Mantisus commented on Sep 22, 2026

    @Mantisus
    CollaboratorAuthor

    Starting it isn't a problem, since EventManager.__aenter__ never actually awaits anything. It only starts the recurring tasks, which a sync on or emit can do just as well.

    The part that needs thought is closing it. I'm considering two options.

    1. A task created with loop.create_task that waits until it's cancelled. On exit, asyncio.run cancels every pending task and then awaits them, so the task's finally block still runs with a working loop.
    class SingleEventManager(EventManager):
    
        def __init__(self, **kwargs: Any) -> None:
            super().__init__(**kwargs)
    
            self._shutdown_task: asyncio.Task[None] | None = None
    
        def on(self, *, event: Event, listener: EventListener[Any]) -> None:
            self._ensure_started()
            super().on(event=event, listener=listener)
    
        def emit(self, *, event: Event, event_data: EventData) -> None:
            self._ensure_started()
            super().emit(event=event, event_data=event_data)
    
        async def close(self) -> None:
            if (shutdown_task := self._shutdown_task) is None:
                return
    
            shutdown_task.cancel()
            await asyncio.wait([shutdown_task])
    
            if not shutdown_task.cancelled() and (exc := shutdown_task.exception()) is not None:
                raise exc
    
            await self._close()
    
        def _ensure_started(self) -> None:
            if self.active:
                return
    
            try:
                loop = asyncio.get_running_loop()
            except RuntimeError:
                raise RuntimeError(
                    'The event manager can only be used from within a running asyncio event loop.'
                ) from None
    
            self._emit_persist_state_event_rec_task.start()
    
            self._shutdown_task = loop.create_task(self._wait_for_shutdown(), name='single-event-manager-shutdown')
    
        async def _wait_for_shutdown(self) -> None:
            try:
                await asyncio.Event().wait()
            finally:
                # Stops the recurring task, emits a last `PersistState` and waits for the listeners.
                await self._close()
    1. An async generator that's closed by loop.shutdown_asyncgens(), which asyncio.run calls on exit, after all tasks are done.
    class SingleEventManager(EventManager):
    
        def __init__(self, **kwargs: Any) -> None:
            super().__init__(**kwargs)
    
            self._shutdown_hook: AsyncGenerator[None, None] | None = None
    
        def on(self, *, event: Event, listener: EventListener[Any]) -> None:
            self._ensure_started()
            super().on(event=event, listener=listener)
    
        def emit(self, *, event: Event, event_data: EventData) -> None:
            self._ensure_started()
            super().emit(event=event, event_data=event_data)
    
        async def close(self) -> None:
            if self._shutdown_hook is not None:
                await self._shutdown_hook.aclose()
    
        def _ensure_started(self) -> None:
            if self.active:
                return
    
            try:
                asyncio.get_running_loop()
            except RuntimeError:
                raise RuntimeError(
                    'The event manager can only be used from within a running asyncio event loop.'
                ) from None
    
            self._emit_persist_state_event_rec_task.start()
    
            self._hook_into_loop()
    
        def _hook_into_loop(self) -> None:
            hook = self._wait_for_shutdown()
    
            # Start the async generator from a sync method, so that the loop registers it.
            send = hook.asend(None)
            try:
                send.send(None)
            except StopIteration:
                self._shutdown_hook = hook
                return
    
            send.close()
            raise RuntimeError('The shutdown hook must not await anything before its `yield`.')
    
        async def _wait_for_shutdown(self) -> AsyncGenerator[None, None]:
            try:
                yield
            finally:
                # Stops the recurring task, emits a last `PersistState` and waits for the listeners.
                await self._close()

    The main difference is ordering. With the task, the final flush runs concurrently with the teardown of the other tasks. The generator runs after all of them have finished.

  7. janbuchar commented on Sep 24, 2026

    @janbuchar
    Collaborator

    Tying the event manager lifecycle to the lifetime of the event loop is a viable option, sure.

    Then again, in crawlee-js, we went with the owned-or-injected pattern for most crawler dependencies (either crawler constructs and tears down a default implementation, or it receives a readymade instance via constructor and the caller is responsible for the lifecycle management). This is nice for services that the user may actually want to override, because it makes the lifecycle clear.

    I guess that pretty much nobody wants to override the event manager though, so we shouldn't force people to own the lifecycle management in this case.

    Regarding the async startup, this is true for the current implementations, but what if some alternative event manager actually needs async initialization?

  8. Mantisus commented on Sep 24, 2026

    @Mantisus
    CollaboratorAuthor

    The context manager stays for the cases that need explicit lifecycle management. In the SDK the Actor owns the ApifyEventManager: it initializes it and closes it when the Actor exits. I believe, that's correct behaviour and it shouldn't change.

    The rule is then the same in both directions: whoever enters the context also closes it, and a manager nobody owns closes itself with the event loop.

    So a custom implementation that needs to await during initialization uses the context manager. And users who don't want to think about the event manager at all get the lazy initialization, with the manager working whether or not a crawler is running at that moment.

    For implementations that need the async initialization, we can add a class attribute like supports_lazy_start, so that the lazy path raises instead of leaving the manager half-initialized.

  9. B4nan commented on Oct 6, 2026

    @B4nan
    Member

    @janbuchar I guess we forgot about this one, if you think we should make this in the JS part too, let's create an issue on that side

  10. B4nan commented on Oct 6, 2026

    @B4nan
    Member

    Created the JS counterpart: apify/crawlee#4224

    Same two problems there. Passing any custom service to a crawler (even just configuration) can create a second event manager, and standalone storages never get periodic persistState. Node has no event loop shutdown hook, so the close side would go through an unref()'d interval plus a final flush in process.once('beforeExit').

    One thing to watch here: the lazy start should hold its own ref in _active_ref_count. Otherwise, when the Actor enters and exits a manager that already started lazily, the last exit clears the listeners that standalone components registered.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    t-toolingIssues with this label are in the ownership of the tooling team.

    Type

    No type

    Projects

    No projects

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions