Skip to main content

stapel-realtime

CI coverage pypi downloads python license llms.txt

Realtime delivery substrate: the L1 library behind the Signal primitive (stapel_core.comm.signal). Ships the Channels/Redis transport for the core's signal-delivery seam, the two consumers every browser socket in the fleet is built from (EphemeralStreamConsumer for at-most-once Signal fan-out; ResumableStreamConsumer for hello/welcome/replay/live journals with seq dedup and a bounded replay window), the versioned v1 wire envelope, the canonical :<scope_type>:<scope_id>[:] stream key, a fail-closed per-stream authorize seam with the workspace-capability authorizer, revoke-to-kick, heartbeat with JWT-exp re-check, disconnect-on-overflow backpressure, the fleet close-code canon, build_websocket_application() host assembly with a port-aware origin guard, the fleet's presence oracle, and eight system checks. Presence is a TTL lease in the fleet-shared cache (stapel_core.core.fleet_cache), written by the base consumer on connect, on every heartbeat tick and on disconnect — no model, no migration — and read by two comm Functions: realtime.is_live {user_id, family?} -> {live, sessions, last_seen} and realtime.live_batch {user_ids[<=100], family?} -> {users: {id: {...}}}. That pair is the whole comm surface; there are no models, migrations, views or urls, and no HTTP route of its own, so a peer asks over the bus like any other Function. It is installed as a Django app so the checks and the Functions are registered.

Part of the Stapel framework — composable Django apps that deploy as a monolith or as microservices without changing module code.

Install

pip install stapel-realtime

At a glance

Fact Value
Version 0.2.1
Python >=3.11 (3.11, 3.12, 3.13)
Config axes 9
Usage surface 23
Extension points 6
Fleet dependencies stapel-core

Documentation

capabilities.json · llms.txt (for agents)

What this is

The delivery half of the fourth communication primitive.

Three primitives in stapel_core.comm address code. Function — "answer me now", the caller waits. Action — "this happened, the system must know": outbox, at-least-once, 0..N module subscribers. Task — "do the long work", the system waits, not the caller.

Signal is the fourth, and its addressee is a human looking at a screen:

Show this to whoever is watching right now.

There is no obligation to an observer who is not watching. When they look, they read current state over REST — the truth is in the database, and the value of a signal expires in seconds. Losing a signal is correct behaviour, and that one property is what lets this library be small: no outbox row, no retry, no history, no delivery receipt.

stapel_core.comm.signal() is the emitter — sixty lines of stdlib, free for every library in the fleet, a silent no-op with no backend configured. This package is everything on the other side of that call.

What it ships

Transport deliver(stream_key, frame) — the callable the core's STAPEL_COMM["SIGNAL_TRANSPORT"] = "channels" resolves to, registered from this package's AppConfig.ready(); plus deliver_frame() for journal fan-out and revoke() for the kick. Best-effort by contract: no layer, no subscriber, dead redis → the frame is dropped and nothing raises.
Two consumers EphemeralStreamConsumer (Signal fan-out, no seq, no history) and ResumableStreamConsumer (hello{last_seq}welcome → replay → live, deduplicated by seq, bounded replay window). Both are generalizations of stapel_chat.ChatConsumer, the one protocol the fleet had actually proven.
Wire envelope v1 {v, type, stream, payload, seq?} — the shape comm.signal() builds and this substrate forwards verbatim, published as a JSON schema (the deliberate exception to "an L1 library ships no schemas": the contract is shared by a backend consumer and a browser client written by different hands). Frame kind is structural — seq present means journal, absent means ephemeral.
Stream keys <mod>:<scope_type>:<scope_id>[:<topic>], built and validated by the core's comm.stream_key() (re-exported here, never re-implemented). The scope is in the name, so a group physically cannot cross a workspace.
Authorization A per-stream authorize() hook that is fail-closed: a consumer that does not implement it subscribes nobody. WorkspaceCapability is the canonical implementation — the same require_capability predicate HTTP uses.
Revoke → kick revoke(stream_key, user_id) sends a kick frame and closes 4410 immediately, rather than leaking until the client happens to reconnect.
Host assembly build_websocket_application() — origin guard (compared with the port) over core's G14 JWT stack over every installed module's routing manifest, discovered rather than listed.
Presence The fleet's answer to "is this person watching right now?" — a TTL lease in the shared cache, written by the base consumer on connect / heartbeat / disconnect, read over the bus as realtime.is_live and realtime.live_batch. No model, no migration.
System checks Seven, each one a production bruise turned into a manage.py check verdict.
Test harness stapel_realtime.testing.open_stream() — an envelope-aware Channels client, so a module testing its consumer does not wire the fourth WebsocketCommunicator by hand.

Quick start

# myapp/consumers.py
from stapel_realtime import EphemeralStreamConsumer, WorkspaceCapability

class RecordingsConsumer(EphemeralStreamConsumer):
    module = "recordings"
    scope_type = "ws"
    stream_key_kwarg = "workspace_id"
    authorizer = WorkspaceCapability("recordings.read")
# myapp/routing.py — the manifest the host assembly discovers
from django.urls import path
from .consumers import RecordingsConsumer

websocket_urlpatterns = [
    path("ws/recordings/<uuid:workspace_id>", RecordingsConsumer.as_asgi()),
]
# myapp/services.py — the emit side. Note what is NOT imported: a module that
# only signals depends on the core, never on this library.
from stapel_core.comm import signal, stream_key

with transaction.atomic():
    recording.status = "ready"
    recording.save()
    signal(stream_key("recordings", "ws", recording.workspace_id),
           "recording.status",
           {"recording_id": str(recording.pk), "status": recording.status})
# asgi.py — the whole host
from django.core.asgi import get_asgi_application
from stapel_realtime.asgi import build_websocket_application

application = build_websocket_application(http_application=get_asgi_application())
# settings.py
INSTALLED_APPS += ["stapel_realtime"]      # so the system checks are registered

STAPEL_COMM = {"SIGNAL_TRANSPORT": "channels"}   # opt in; the default is "none"

STAPEL_REALTIME = {
    "ALLOWED_ORIGINS": ["https://app.example.com"],   # WITH the port if non-default
}
CHANNEL_LAYERS = {
    "default": {
        "BACKEND": "channels_redis.core.RedisChannelLayer",
        "CONFIG": {"hosts": ["redis://redis:6379/0"]},
    }
}

Install: pip install 'stapel-realtime[channels,redis]' on a host that serves sockets, and [testing] on top wherever a module tests its own consumer (that extra adds daphne, which channels.testing drags in — no reason to put an ASGI server on a production host). A module that only emits needs nothing from here at all: comm.signal() lives in the core, and that is the point.

Presence: who is watching right now

Every sender of a Signal eventually needs the question the substrate is the only place able to answer. The first to need it was an incoming call: the ring is pushed to the callee's phone and rung in the tab they already have open, and with nothing to ask, the push went out unconditionally and every client carried the workaround of suppressing a banner for a call it was already ringing.

The base consumer writes a TTL lease into the fleet-shared cache on connect, on every heartbeat tick, and on disconnect. Anyone asks over the bus:

from stapel_core.comm import call

if not call("realtime.is_live", {"user_id": str(callee_id)})["live"]:
    notify(callee_id, "call.incoming", ...)   # nobody is looking; push it

# or, for a group, in one round trip (≤100 ids, every id comes back)
live = call("realtime.live_batch", {"user_ids": ids})["users"]

{"live": bool, "sessions": int, "last_seen": iso|null}sessions counts open sockets, so two tabs are two sessions and one person. An optional "family" narrows the question to one stream family (chat, video, …).

Four properties worth knowing before you gate anything on it:

  • It is a lease, not a counter. A worker killed mid-socket never runs its disconnect; one PRESENCE_TTL_S later that session simply stops counting. Which is why the TTL must stay above HEARTBEAT_Srealtime.W006.
  • It fails to "not live". No cache, dead redis, corrupt document: the answer is false and nothing raises. A caller gating a push therefore falls back to sending it, which is exactly the behaviour that existed before.
  • It is fleet-shared on purpose. The write goes through stapel_core.core.fleet_cache, not django.core.cache, because the service holding the socket is not the service asking. On a locmem cache it is per-process and realtime.W005 says so.
  • It is not a last-seen history. last_seen outlives the session by one TTL and no longer. A durable "last online" belongs to a profile row.

The rule that keeps a fifth implementation from appearing

Before this library the fleet had three independent browser sockets (chat, video, studio-dialog) plus a machine peer protocol, each with its own JWT handling, its own close codes, and — twice, independently — its own resume protocol. The boundary is drawn by who is on the other end:

A human in a browser → stapel-realtime. One of our own processes → an application-level protocol (stapel-runner-protocol), and it owes an answer to "why not a Function or a Task".

The machine protocol stays separate on merit, not inertia: a dropped task.assign frame is unacceptable where a dropped signal is correct, it needs exactly-once apply keyed by (task_id, seq), and it is deliberately transport-agnostic so it can be tested without a network.

What it does not do

Not in v1, on purpose: an SSE fallback, one multiplexed socket for many streams (the envelope reserves stream so adding it later is not a breaking change), NATS as the signal transport (that is a future value of the core's axis, for the microservice topology), client→server commands over the socket (writes go through REST/Function), and delivering Actions to the browser as-is — an anti-pattern, because a five-minute-late "typing…" retried by an outbox is worse than no delivery at all. And a presence.changed signal: presence is asked, not announced, because a fan-out on every connect and disconnect in the fleet is a great deal of traffic for a fact that costs one cache read — a module that wants to paint a green dot subscribes to the one its own domain already emits (chat.presence.changed).

License

MIT — see LICENSE.


This page is assembled by stapel-readme from docs/readme.md plus the contract artifacts in docs/. Edit the prose in docs/readme.md; the badges, facts and links above and below it are generated — do not hand-edit README.md.

Download files

Download the file for your platform. If you're not sure which to choose, learn more about installing packages.

Source Distribution

stapel_realtime-0.2.1.tar.gz (81.1 kB view details)

Uploaded Source

Built Distribution

If you're not sure about the file name format, learn more about wheel file names.

stapel_realtime-0.2.1-py3-none-any.whl (72.4 kB view details)

Uploaded Python 3

File details

Details for the file stapel_realtime-0.2.1.tar.gz.

File metadata

  • Download URL: stapel_realtime-0.2.1.tar.gz
  • Upload date:
  • Size: 81.1 kB
  • Tags: Source
  • Uploaded using Trusted Publishing? Yes
  • Uploaded via: twine/7.0.0 CPython/3.13.14

File hashes

Hashes for stapel_realtime-0.2.1.tar.gz
Algorithm Hash digest
SHA256 a153e20cb43fd603433a9a142aaa66544fce17a361879958c97a5a60e35ffb71
MD5 7d041b24760ed99e0fe85c9043c9d2af
BLAKE2b-256 7068b4461de01e7fff5127d35c8544bf89930df7794f211fcdfc7b23c41da121

See more details on using hashes here.

Provenance

The following attestation bundles were made for stapel_realtime-0.2.1.tar.gz:

Publisher: publish.yml on usestapel/stapel-realtime

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

File details

Details for the file stapel_realtime-0.2.1-py3-none-any.whl.

File metadata

File hashes

Hashes for stapel_realtime-0.2.1-py3-none-any.whl
Algorithm Hash digest
SHA256 e49948864e4d526f021aa7182a35e33a4893f754214b6b374dff8a07817e0668
MD5 d1e6b2f63d844d60807c25e04a8b3a02
BLAKE2b-256 e3453298fb429791aac68c01e033e8f20ef993362fe96850869d03e8a6d28849

See more details on using hashes here.

Provenance

The following attestation bundles were made for stapel_realtime-0.2.1-py3-none-any.whl:

Publisher: publish.yml on usestapel/stapel-realtime

Attestations: Values shown here reflect the state when the release was signed and may no longer be current.

Release history Release notifications | RSS feed

This release

0.2.1 This release

2 files

0.2.0

2 files

0.1.4

2 files

0.1.3

2 files

0.1.2

2 files

0.1.1

2 files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page