stormlog.infer.open_loop

Send requests at scheduled times, bounded by an in-flight limit.

Functions

cancel_all(tasks)

Cancel tasks and wait until each has finished cancelling.

dispatch_schedule(offsets, *, mode, limiter, ...)

Start each request at its offset from now, in seconds.

Classes

Arrival(index, mode, intended_at_ns[, ...])

One scheduled request and how it got onto the wire.

Dispatch(started_at_ns, tasks)

The requests a schedule sent, and when the schedule started.

InFlightLimiter(limit)

Count outstanding requests and hold new ones at the limit.

class stormlog.infer.open_loop.Arrival(index, mode, intended_at_ns, held_for_slot=False, in_flight_at_dispatch=None)[source]

Bases: object

One scheduled request and how it got onto the wire.

Parameters:
  • index (int)

  • mode (str)

  • intended_at_ns (int)

  • held_for_slot (bool)

  • in_flight_at_dispatch (int | None)

index: int
mode: str
intended_at_ns: int
held_for_slot: bool = False
in_flight_at_dispatch: int | None = None
class stormlog.infer.open_loop.InFlightLimiter(limit)[source]

Bases: object

Count outstanding requests and hold new ones at the limit.

Parameters:

limit (int)

property full: bool
async acquire()[source]

Wait for a free slot; return the in-flight count including this one.

Return type:

int

release()[source]
Return type:

None

class stormlog.infer.open_loop.Dispatch(started_at_ns, tasks)[source]

Bases: object

The requests a schedule sent, and when the schedule started.

tasks holds the requests still running and any that failed, whose errors are raised by whoever drains them; finished requests are dropped so a long schedule holds no more than its in-flight limit.

Parameters:
  • started_at_ns (int)

  • tasks (list[Task[None]])

started_at_ns: int
tasks: list[Task[None]]
async stormlog.infer.open_loop.dispatch_schedule(offsets, *, mode, limiter, overflow, send, drop, deadline=None)[source]

Start each request at its offset from now, in seconds.

With overflow="wait" an arrival that finds every slot busy is held until one frees up; while it waits, later arrivals fall behind schedule too, which their dispatch lag shows. With "drop" it is handed to drop and never sent, with the reason. deadline is when the phase’s drain ends, in seconds from now: an arrival still unsent then is dropped rather than held any longer. Returns once every arrival has been sent or dropped; the requests themselves may still be running.

Parameters:
  • offsets (Sequence[float])

  • mode (str)

  • limiter (InFlightLimiter)

  • overflow (Literal['wait', 'drop'])

  • send (Callable[[Arrival], Awaitable[None]])

  • drop (Callable[[Arrival, str], None])

  • deadline (float | None)

Return type:

Dispatch

async stormlog.infer.open_loop.cancel_all(tasks)[source]

Cancel tasks and wait until each has finished cancelling.

Parameters:

tasks (Sequence[Task[None]])

Return type:

None