//clika-runtime/io.clika.runtime/Pipeline
Pipeline
[common]
class Pipeline
A multi-step serving ensemble: each step runs a model under its own executor tuning, wired by name maps into one request/response surface declared by the external schema. enqueue validates the request against that schema and answers at once, with or without callbacks; await blocks for the session's response; cancel stops a session at its next step boundary (its completion reads CANCELLED); shutdown drains and a later enqueue fails with FAILED_PRECONDITION. AutoCloseable: close shuts down and releases the pipeline, and the models stay the caller's.
Pipeline.create(listOf(PipelineStep("node", model, inputMap = listOf("x" to "x"), outputMap = listOf("y" to "y"))), schema).use { pipeline ->
pipeline.enqueue(Request(mapOf("x" to x), streaming = true), RequestCallbacks(onChunk = { chunk -> chunk.use { println(it.final) } }))
}
Types
| Name | Summary |
|---|---|
| Companion | [common] object Companion |
Functions
| Name | Summary |
|---|---|
| await | [common] fun await(session: SessionId): Response Block until session completes; its response, or its failure. |
| cancel | [common] fun cancel(session: SessionId) Cooperatively cancel an in-flight session; a completed one is a no-op, an unknown one fails with NOT_FOUND. |
| close | [common] open fun close() |
| enqueue | [common] fun enqueue(request: Request): SessionId Submit a request; fails with INVALID_ARGUMENT naming the input a declared slot refuses.[common] fun enqueue(request: Request, callbacks: RequestCallbacks): SessionId Submit a request with its hooks, delivered on the engine's thread. |
| shutdown | [common] fun shutdown() Drain in-flight work and stop; every accepted session completes and await still answers it. |