Skip to content

Streaming & Async IO

Spikard treats streaming and real-time protocols as first-class citizens so the same APIs work for HTTP, WebSocket, and SSE flows.

Capabilities

  • WebSockets: bidirectional handlers with typed messages and backpressure-aware send/receive loops.
  • Server-Sent Events: push event streams with graceful client disconnect handling.
  • Chunked responses: stream files or long-running computations without buffering everything in memory.

Concurrency Model

  • Built on Tokio with cooperative scheduling
  • Binding bridges expose async primitives that map to the language runtime (e.g., Python event loop thread, Node async iterators)
  • Cancellation and shutdown signals propagate through middleware and handlers

Validation for Streaming

  • Envelope/message schemas validated per message when configured
  • Connection-level auth and capability checks enforced in middleware before the handler executes

Streaming Response

from spikard import SseEvent, sse

@sse("/events")
async def events():
    for i in range(3):
        yield SseEvent(data={"tick": i})
import { Spikard, StreamingResponse } from "spikard";

const app = new Spikard();

async function* makeStream() {
  for (let i = 0; i < 3; i++) {
    yield JSON.stringify({ tick: i }) + "\n";
    await new Promise((resolve) => setTimeout(resolve, 100));
  }
}

app.addRoute(
  { method: "GET", path: "/stream", handler_name: "stream", is_async: true },
  async () =>
    new StreamingResponse(makeStream(), {
      statusCode: 200,
      headers: { "Content-Type": "application/x-ndjson" },
    }),
);
require "spikard"
require "json"

app = Spikard::App.new

app.get "/stream" do |_params, _query, _body|
  Enumerator.new do |y|
    3.times do |i|
      y << JSON.dump({ tick: i }) + "\n"
      sleep 0.1
    end
  end
end
<?php

declare(strict_types=1);

use Spikard\App;
use Spikard\Attributes\Get;
use Spikard\Config\ServerConfig;
use Spikard\Http\StreamingResponse;

$app = new App(new ServerConfig(port: 8000));

final class StreamController
{
    #[Get('/stream')]
    public function stream(): StreamingResponse
    {
        $generator = function (): Generator {
            for ($i = 0; $i < 10; $i++) {
                yield json_encode(['chunk' => $i]) . "\n";
                usleep(100000); // 100ms delay
            }
        };

        return new StreamingResponse(
            $generator(),
            headers: ['Content-Type' => 'application/x-ndjson']
        );
    }
}

$app = $app->registerController(new StreamController());
use spikard::prelude::*;
use tokio_stream::StreamExt;

app.route(get("/stream"), |_ctx: Context| async move {
    let stream = tokio_stream::iter(0..3).then(|i| async move {
        serde_json::to_vec(&serde_json::json!({ "tick": i }))
    });
    Ok(StreamingBody::new(stream))
})?;

Server-Sent Events

from spikard import Spikard, SseEvent, sse

app = Spikard()

@sse("/events")
async def events():
    for i in range(3):
        yield SseEvent(data={"tick": i})
import { Spikard, StreamingResponse } from "spikard";

const app = new Spikard();

async function* sseStream() {
  for (let i = 0; i < 3; i++) {
    yield `data: ${JSON.stringify({ tick: i })}\n\n`;
  }
}

app.addRoute(
  { method: "GET", path: "/events", handler_name: "events", is_async: true },
  async () =>
    new StreamingResponse(sseStream(), {
      statusCode: 200,
      headers: { "Content-Type": "text/event-stream" },
    }),
);
require "spikard"
require "json"

app = Spikard::App.new

app.get "/events" do |_params, _query, _body|
  Enumerator.new do |y|
    3.times do |i|
      y << "data: #{JSON.dump({ tick: i })}\n\n"
      sleep 0.1
    end
  end
end
<?php

declare(strict_types=1);

use Spikard\App;
use Spikard\Attributes\Get;
use Spikard\Config\ServerConfig;
use Spikard\Http\StreamingResponse;

final class EventsController
{
    #[Get('/events')]
    public function events(): StreamingResponse
    {
        $generator = function (): Generator {
            for ($i = 0; $i < 5; $i++) {
                $data = json_encode(['tick' => $i, 'time' => time()]);
                yield "data: {$data}\n\n";
                sleep(1);
            }
            yield "data: " . json_encode(['message' => 'done']) . "\n\n";
        };

        return StreamingResponse::sse($generator());
    }
}

$app = (new App(new ServerConfig(port: 8000)))
    ->registerController(new EventsController());
use spikard::prelude::*;
use tokio_stream::StreamExt;

app.route(get("/events"), |_ctx: Context| async move {
    let stream = tokio_stream::iter(0..3).map(|i| {
        format!("data: {}\n\n", serde_json::json!({"tick": i}))
    });
    Ok(StreamingBody::new(stream).with_header("content-type", "text/event-stream"))
})?;

WebSocket

from spikard import Spikard, websocket

app = Spikard()

@websocket("/ws")
async def echo(message: dict) -> dict | None:
    return {"echo": message}
import { Spikard } from "spikard";

const app = new Spikard();

app.addRoute({ method: "WS", path: "/ws", handler_name: "ws", is_async: true }, async (socket) => {
  for await (const message of socket) {
    await socket.send({ echo: message });
  }
});
require "spikard"

app = Spikard::App.new

class ChatHandler < Spikard::WebSocketHandler
  def handle_message(message)
    # Echo JSON message back
    message
  end
end

app.websocket("/chat") { ChatHandler.new }
<?php

declare(strict_types=1);

use Spikard\App;
use Spikard\Config\ServerConfig;
use Spikard\Handlers\WebSocketHandlerInterface;

final class ChatHandler implements WebSocketHandlerInterface
{
    public function onConnect(): void
    {
        error_log('Client connected');
    }

    public function onMessage(string $message): void
    {
        $data = json_decode($message, true);
        error_log('Received: ' . json_encode($data));
    }

    public function onClose(int $code, ?string $reason = null): void
    {
        error_log("Client disconnected: {$code}" . ($reason ? " ({$reason})" : ''));
    }
}

$app = (new App(new ServerConfig(port: 8000)))
    ->addWebSocket('/ws', new ChatHandler());

$app->run();
use spikard::prelude::*;
use futures::StreamExt;

app.websocket("/ws", |mut socket| async move {
    while let Some(msg) = socket.next().await {
        let text = msg.unwrap_or_default();
        socket.send(text).await.ok();
    }
});

More details and decisions live in ADR 0006.

Edit this page on GitHub