How to Handle Asynchronous Tasks in Colibri: Streaming Completions and Cancellation
Colibri processes long-running operations like chat completions as cancellable Server-Sent Event (SSE) streams, using streamChat from web/src/lib/api.ts together with AbortController for cancellation and activeRequests from web/src/lib/runtime.ts for monitoring server workload.
Handling asynchronous tasks in Colibri requires understanding its streaming-first architecture designed for real-time AI inference. The open-source JustVugg/colibri repository provides a robust client implementation that manages HTTP streaming, partial response parsing, and concurrent request tracking through standardized web APIs.
Colibri's Server-Sent Event Architecture
Colibri delivers asynchronous results via Server-Sent Events (SSE) rather than traditional request-response cycles. When you initiate a streaming chat completion, the client opens a persistent connection to the /chat/completions endpoint with stream: true, enabling the server to push tokens as they are generated. The implementation in web/src/lib/api.ts uses a ReadableStreamDefaultReader to consume the response body incrementally, decoding bytes with TextDecoder and processing them through the extractSSE function.
Streaming Chat Completions with streamChat
The primary interface for managing asynchronous tasks is the streamChat function exported from web/src/lib/api.ts. This method orchestrates the entire streaming lifecycle from connection establishment to final resolution.
Initiating Streaming Requests
To start an asynchronous task, call streamChat with your model configuration and callback handlers. The function accepts a signal parameter for cancellation support and optional cacheSlot for KV cache optimization when the server supports it.
import { streamChat } from "./lib/api";
const abort = new AbortController();
await streamChat({
baseUrl: "https://colibri.example.com/v1",
apiKey: "YOUR_API_KEY",
model: "gpt-4o-mini",
messages: [{ role: "user", content: "Explain async patterns" }],
stream: true,
signal: abort.signal,
onDelta: (text) => console.log("Token:", text),
onReasoning: (text) => console.log("Reasoning:", text)
});
Parsing SSE Chunks with extractSSE
Internally, streamChat utilizes the extractSSE helper to parse raw byte streams. This function accumulates incoming chunks in a buffer, splits them into complete SSE frames, and preserves incomplete data for subsequent iterations. Each extracted frame contains a JSON payload with delta content that triggers your onDelta callback, while the optional reasoning_content field fires the onReasoning handler when reasoning tokens are present.
Monitoring Active Asynchronous Tasks
Before initiating resource-intensive operations, Colibri provides utilities to inspect server capacity through its health endpoint.
Checking Server Health with getHealth
The getHealth function queries the server's status endpoint and returns a HealthResponse object. This response contains scheduler metadata that indicates current system load and available inference slots.
import { getHealth } from "./lib/api";
import { activeRequests, supportsCacheSlots } from "./lib/runtime";
const health = await getHealth("https://colibri.example.com/v1", "YOUR_API_KEY");
console.log("Active jobs:", activeRequests(health));
console.log("Cache supported:", supportsCacheSlots(health));
Tracking Concurrent Requests
The activeRequests function in web/src/lib/runtime.ts extracts health.scheduler.active from the health response, returning the count of currently executing streaming jobs. This allows you to implement client-side backpressure or queue management when the server approaches capacity. The supportsCacheSlots helper determines whether you can pass the cache_slot parameter to streamChat for optimized context caching.
Cancellation and Error Handling
Robust async task management requires immediate cancellation capabilities and comprehensive error recovery.
Aborting Requests with AbortController
Pass an AbortSignal from an AbortController instance to streamChat. Calling abort() terminates the HTTP connection immediately, stopping both the server-side generation and client-side processing without waiting for the stream to complete naturally.
// Cancel after 5 seconds
setTimeout(() => abort.abort(), 5000);
Handling HTTP and Streaming Errors
If the server returns a non-OK status code, Colibri's internal responseError utility parses the JSON error payload to extract meaningful diagnostic messages. Network failures and parsing errors propagate through the promise returned by streamChat, enabling standard try-catch patterns for error recovery.
Complete Working Example
This example demonstrates initiating a cancellable streaming request while monitoring server capacity:
import { streamChat, getHealth } from "./lib/api";
import { activeRequests } from "./lib/runtime";
async function runAsyncTask() {
const abort = new AbortController();
// Check server load first
const health = await getHealth("https://colibri.example.com/v1", "API_KEY");
console.log(`Server handling ${activeRequests(health)} concurrent tasks`);
try {
await streamChat({
baseUrl: "https://colibri.example.com/v1",
apiKey: "API_KEY",
model: "gpt-4o-mini",
messages: [{ role: "user", content: "Analyze this code" }],
temperature: 0.7,
maxTokens: 512,
enableThinking: true,
cacheSlot: undefined,
signal: abort.signal,
onDelta: (text) => process.stdout.write(text),
onReasoning: (text) => console.error(`[thinking] ${text}`)
});
} catch (error) {
console.error("Stream failed:", error);
}
// Cancel if user navigates away
window.addEventListener("beforeunload", () => abort.abort());
}
runAsyncTask();
Summary
- Streaming Protocol: Colibri uses Server-Sent Events delivered via
ReadableStreamDefaultReaderinweb/src/lib/api.tsfor real-time token delivery. - Task Control: Pass
AbortSignaltostreamChatto enable immediate cancellation of in-flight asynchronous tasks. - Server Monitoring: Use
activeRequestsfromweb/src/lib/runtime.tson theHealthResponsefromgetHealthto check server load before submitting heavy workloads. - Data Parsing: The
extractSSEfunction handles buffered frame extraction, while callbacks processcontentdeltas and optionalreasoning_contentstreams. - Error Handling: Non-OK responses trigger
responseErrorto extract detailed failure messages from the JSON payload.
Frequently Asked Questions
How does Colibri handle real-time streaming?
Colibri implements real-time streaming through standard HTTP Server-Sent Events. The streamChat function in web/src/lib/api.ts establishes a connection with stream: true, then uses TextDecoder and extractSSE to parse incoming byte chunks into JSON payloads, triggering onDelta callbacks for each token received.
Can I cancel an ongoing asynchronous task in Colibri?
Yes. Create an AbortController and pass its signal property to streamChat. Calling controller.abort() immediately terminates the connection, freeing server resources. This is implemented using the standard AbortSignal interface supported by the fetch API.
How do I check if the server is overloaded before sending a stream?
Query the health endpoint using getHealth, then pass the result to activeRequests from web/src/lib/runtime.ts. This returns the value of health.scheduler.active, indicating how many streaming jobs are currently executing. If this number exceeds your threshold, queue the request or select a different endpoint.
What are cache slots and when should I use them?
Cache slots refer to KV cache optimization available when supportsCacheSlots(health) returns true. When supported, you can pass a cacheSlot identifier to streamChat to enable context caching for repeated prompts, reducing latency on subsequent asynchronous tasks. If the server does not advertise this capability, the parameter is ignored.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →