streams API
SDK 0.4.0. This reference describes the supported authoring interface.
StreamError
Sanitized machine-readable terminal outcome.
Fields
code: Literal['provider_error', 'transport_error', 'unreachable', 'cancelled', 'timeout', 'invalid_frame']
message: str
| Field | Description | Declared defaults and constraints |
|---|---|---|
code | — | — |
message | — | max_length=2048 |
InlineMedia
Raw media plus format metadata; never encode the bytes as JSON.
Fields
kind: Literal['inline']
data: bytes
media_type: str
codec: str | None
sample_rate: int | None
channels: int | None
timestamp_seconds: float | None
duration_seconds: float | None
| Field | Description | Declared defaults and constraints |
|---|---|---|
kind | — | default="inline" |
data | — | — |
media_type | — | max_length=256 |
codec | — | default=null |
sample_rate | — | default=null; gt=0 |
channels | — | default=null; gt=0 |
timestamp_seconds | — | default=null; ge=0 |
duration_seconds | — | default=null; ge=0 |
BlobMedia
Opaque immutable blob reference; resolution remains Fabric-owned.
Fields
kind: Literal['blob']
blob_id: str
size_bytes: int
sha256: str
media_type: str
| Field | Description | Declared defaults and constraints |
|---|---|---|
kind | — | default="blob" |
blob_id | — | — |
size_bytes | — | ge=0 |
sha256 | — | pattern="^[a-f0-9]{64}$" |
media_type | — | — |
StreamFrame
One sequenced input or output frame, aligned with Skulk™'s public contract.
StreamFrame.lifecycle_shape
Refuse malformed lifecycle frames before they reach the transport.
def lifecycle_shape(self) -> Self:
...
StreamFrame.is_terminal
Whether this frame closes its own direction.
def is_terminal(self) -> bool:
...
Fields
call_id: str
direction: Literal['caller_to_provider', 'provider_to_caller']
sequence: int
kind: Literal['started', 'chunk', 'completed', 'failed', 'cancelled']
synthetic: Literal[False]
payload: dict[str, JsonValue] | None
media: Annotated[InlineMedia | BlobMedia, Field(discriminator='kind')] | None
error: StreamError | None
| Field | Description | Declared defaults and constraints |
|---|---|---|
call_id | — | min_length=1; max_length=128 |
direction | — | — |
sequence | — | ge=0 |
kind | — | — |
synthetic | — | default=false |
payload | — | default=null |
media | — | default=null |
error | — | default=null |
StreamEnd
Child cleanup acknowledgment following exactly one output terminal.
Fields
kind: Literal['end']
call_id: str
| Field | Description | Declared defaults and constraints |
|---|---|---|
kind | — | default="end" |
call_id | — | min_length=1; max_length=128 |
StreamPacket
type StreamPacket = StreamFrame | StreamEnd
write_packet
Write one bounded header and raw attachment, awaiting socket backpressure.
async def write_packet(writer: asyncio.StreamWriter, frame: StreamPacket) -> None:
...
read_packet
Bound lengths before reading and reconstruct media without base64 expansion.
async def read_packet(reader: asyncio.StreamReader) -> StreamPacket:
...
Sequence
Fence one connection by call identity, direction, sequence and terminal.
Sequence.__init__
Set the initial sequence for one independently closed stream direction.
def __init__(self, call_id: str, direction: Literal['caller_to_provider', 'provider_to_caller'], *, next_sequence: int=0) -> None:
...
Sequence.accept
Reject gaps, duplicate terminals, wrong direction and stale call frames.
def accept(self, frame: StreamFrame) -> None:
...
StreamHandler
type StreamHandler = Callable[[StreamInvoke, AsyncIterator[StreamFrame]], AsyncGenerator[StreamFrame]]
execute_stream
Serve one admitted stream while the child's main loop answers health.
Call this in an owned task from the main IPC loop. The handler receives caller lifecycle frames; completed input is a half-close. It yields output starting at sequence one (Skulk™ owns output started) and exactly one terminal. Disconnect, invalid input, or deadline expiry cancels and closes its iterator. The bounded eight-frame ingress queue applies backpressure independently of health. The owner additionally validates all declared payload schemas.
async def execute_stream(startup: Startup, call: StreamInvoke, handler: StreamHandler) -> None:
...