Skip to main content

Pipeline

A serving ensemble; submit with enqueue(), collect with wait(), or stream() a request for per-step chunks. Made by pipeline(). Usable as a context manager (shutdown on exit).

__init__​

__init__(self, /, *args, **kwargs)

Initialize self. See help(type(self)) for accurate signature.

cancel​

cancelcancel(self, session: int) -> None

cancel(self, session: int) -> None

Cooperatively cancel an in-flight session: the step holding it stops at its next step boundary, no later step starts, and wait() or the stream then report the cancellation. A finished session is a no-op; an unknown session raises.

enqueue​

enqueueenqueue(self, request: clika_runtime._core.runtime.Request) -> int

enqueue(self, request: clika_runtime._core.runtime.Request) -> int

Submit a request; returns the session id. Non-blocking.

run​

runrun(self, request: clika_runtime._core.runtime.Request) -> clika_runtime._core.runtime.Response

run(self, request: clika_runtime._core.runtime.Request) -> clika_runtime._core.runtime.Response

enqueue() then wait() in one call.

shutdown​

shutdownshutdown(self) -> None

shutdown(self) -> None

Drain in-flight work and stop; every accepted session completes and wait() still returns its result. A later enqueue() raises.

stream​

streamstream(self, request: clika_runtime._core.runtime.Request) -> clika_runtime._core.runtime.ResponseStream

stream(self, request: clika_runtime._core.runtime.Request) -> clika_runtime._core.runtime.ResponseStream

Submit a request for streamed retrieval: returns a ResponseStream yielding a ResponseChunk per streamed step (iteration ends with the session; wait() gives the final Response). The request is submitted with streaming=True.

wait​

waitwait(self, session: int) -> clika_runtime._core.runtime.Response

wait(self, session: int) -> clika_runtime._core.runtime.Response

Block until the session completes; returns its Response (a failed session raises RuntimeError).