Skip to main content

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
FieldDescriptionDeclared 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
FieldDescriptionDeclared 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
FieldDescriptionDeclared 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
FieldDescriptionDeclared 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
FieldDescriptionDeclared 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:
...