xai-org/x-algorithm · error · RuntimeError

Task generators not started

Error message

Task generators not started

What it means

PriorityTaskGenerator._poll reads self._streams, which is only populated during start(). Polling before start() (or after stop() cleared state) means the generator was used out of order, so it raises RuntimeError rather than silently yielding nothing.

Source

Thrown at grox/core/generators/task_generator.py:102

        logger.info(
            f"Initialized priority task generator with {list(zip(self._generators.keys(), [gen.__class__.__name__ for gen in self._generators.values()], self._weights.values(), strict=True))}"
        )

    async def start(self) -> None:
        logger.info("Starting priority task generators")
        await asyncio.gather(*[gen.start() for gen in self._generators.values()])
        self._streams = {label: gen.poll() for label, gen in self._generators.items()}
        logger.info("Priority task generators started")

    async def stop(self) -> None:
        logger.warning("Stopping priority task generators")
        await asyncio.gather(*[gen.stop() for gen in self._generators.values()])
        await super().stop()
        logger.warning("Priority task generators stopped")

    async def _poll(self) -> AsyncGenerator[TaskPayload | None, None]:
        if not self._streams:
            raise RuntimeError("Task generators not started")
        while self._weights:
            _weights = self._weights.copy()
            polled = False
            while _weights:
                labels = list(_weights.keys())
                weights = list(_weights.values())
                labels = random.choices(labels, weights, k=1)
                label = labels[0]
                stream = self._streams[label]
                try:
                    payload = await anext(stream)
                    if payload:
                        self._result_cache[payload.payload_id] = label
                        yield payload
                        polled = True
                        break
                    else:
                        del _weights[label]

View on GitHub (pinned to 24c60942c5)

Solutions

  1. Ensure `await generator.start()` completes before any polling/iteration begins.
  2. Don't reuse a PriorityTaskGenerator after stop(); construct a fresh instance.
  3. In async setups, await start() inside the same task/flow that creates the generator before handing it to consumers.
  4. Add a startup readiness check or event before starting consumers.

Example fix

# before
pg = PriorityTaskGenerator(gens)
async for task in pg: ...  # RuntimeError: not started

# after
pg = PriorityTaskGenerator(gens)
await pg.start()
async for task in pg: ...
Defensive patterns

Strategy: validation

Validate before calling

await pg.start()
assert pg._streams, "start() must populate streams before polling"

Type guard

async def is_started(pg) -> bool:
    return bool(getattr(pg, "_streams", None))

Try / catch

try:
    async for task in pg:
        ...
except RuntimeError as e:
    if "not started" in str(e):
        await pg.start()
    else:
        raise

Prevention

When it happens

Trigger: Calling the generator's polling/iteration API before awaiting start(); awaiting stop() and then polling again; using the generator after an aborted startup.

Common situations: Skipping the start step in tests; race where a consumer task begins polling before the startup task completes; reusing a stopped generator instance.

Related errors


AI-assisted analysis of xai-org/x-algorithm@24c60942c5 (2026-08-28). Data as JSON: /api/errors/491f5e9a2a3e8609. Report an issue: GitHub.