Skip to content
SubEtha Python

SubEtha Python

The Python binding (subetha-py)

subetha-py gives Python the memory-mapped primitives directly, with no C shim between the interpreter and the Rust. It ships as a wheel rather than to crates.io.

pip install subetha-ipc

The distribution is subetha-ipc because subetha on PyPI belongs to an unrelated project. The package it installs is subetha, so code says import subetha.

This page describes the surface by family. The complete surface, every name with every type, is generated from the type stub and split across two pages:

Both are produced by crates/subetha-py/tools/export_reference.py reading python/subetha/__init__.pyi, which tests/test_surface.py holds against the compiled module, so a name there is a name that ships, and python/subetha/sidecar.py. Regenerate them whenever the surface changes.

To install it and send a first message, start at SubEtha from Python .

What a call costs, and what follows

Measured on a Ryzen 9 7900X, taking the same atomic load every way it can be reached:

reached byns per operation
Rust, through the C ABI7.1
Python, this binding30.1
Python, this binding, with an argument36.8
Python, through a C shim with ctypes584.4
Python, this binding, batched a thousand at a time1.3

The two Python rows to read against each other are 30.1 and 584.4: the same operation, nineteen times cheaper. An empty Python loop costs 7.7 ns an iteration on the same machine, so about 22 ns of the 30 is the call itself.

Two shapes cross the boundary less often, and the surface is built around them:

  • A batch call carries many operations over one crossing. fetch_add_many, send_many, recv_many, insert_many, get_many, drain and their kin are all this shape.
  • A buffer crosses once and is then read with no call at all. memoryview(region) is a view over the mapped file: read whole it is 0.3 ns a byte. Walking that same view one element at a time from Python costs 33 ns an element, a hundred times the bulk figure, and none of that is the mapping. A buffer pays off when something consumes it whole, numpy included, and not when a Python loop walks it.

bench/call_shapes.py in the crate reproduces every row above on the reader’s own machine.

The conventions

  • A refusal is an answer, not a fault. A push that does not fit returns False; a pop with nothing to take returns None; a sweep that freed nothing returns 0; a sample that kept nothing returns None. Only a genuine fault raises.
  • Anything holding a resource is a context manager, and gives the resource back when its block ends, including on the way out of an exception, and again if it is dropped without a block.
  • Exceptions in place of error codes. ValueError for an argument that cannot be right, OSError for a failure underneath, plus two of the package’s own: Lagged when a subscriber fell far enough behind that what it asked for was overwritten, and Contended when a lease or a lane is held by somebody else.
  • A value that is not measured yet is None, which is a different answer from zero.
  • The package ships py.typed and a full stub, and tests/test_surface.py holds the stub against the module in both directions, so a class in one and not the other fails the suite.

The front door

Three classes that pick the shape underneath from what is described, rather than asking the caller to choose one. Reach for these first; the families below give more control when it is wanted.

ClassWhat it is
ChannelA queue between processes that can be waited on. recv answers None at once when there is nothing there; recv_for waits up to a timeout, sleeping rather than spinning, with the interpreter free so other threads run. send_for is the same on the other side.
WorkQueueWork one process owns and others take from when idle. The owner pushes and pops at the cheap end it has to itself; thieves steal from the other end.
KvMapA lookup table between processes, from one unsigned integer to another.
AdaptiveQueueThe one that picks from what happens rather than what was declared. It counts the sizes it is sent and moves between a ring and a work-stealing deque while running, without either end reconnecting. Worth it when the traffic is not known in advance; when it is, saying so is cheaper, because this one pays a counter on every send.

KvMap.insert answers whether the key was new rather than what it held, because that is what the map underneath reports, and there is no way to take a key out because it has no removal. HashMap is the one to reach for when entries have to go away.

Rings and channels

ClassWhat it is
RingThe adaptive ring: producers and consumers register for an id, payloads larger than a slot are framed, and the shape changes under the traffic. Built with stamps it marks each item with the order its sender made it in. With managed=True it also runs a shape sidecar of its own, scanning every scan_interval_us microseconds, 250 when not given, and sidecar_morphs counts the morphs that sidecar made.
SpscRingOne writer, one reader.
BroadcastRingOne writer, many readers, each seeing everything.
PubSub, SubscriberKeeps the last N and tells a slow reader what it lost, by raising Lagged.
MpscProducer, MpscConsumerMany writers, one reader. Built by mpsc_pool / mpsc_pool_open.
MpmcProducer, MpmcConsumerMany of each. Built by mpmc_grid / mpmc_grid_open.
LamportProducer, LamportConsumerThe Lamport pair. Built by lamport_pair / lamport_pair_open.

The constructors hand out both ends together because the shape is what makes them correct.

Rings that change themselves

ClassWhat it is
CapacityRingResizes without losing what is in flight: a reader keeps reading the superseded backing until it has drained. Reports the prewarmed backing, the morphs that took it, and the reads served from a superseded one.
LocaleRingMoves between process-private memory, a mapped file and named memory without its senders and readers reconnecting, carrying what it holds across the move.

Both accept stamped, and report ordering_mode as one of unordered, merge_by_stamp or merge_strict.

Both also accept managed=True, which starts a sidecar of the ring’s own scanning every scan_interval_us microseconds. Neither names a default interval, so a managed one needs it, and scan_interval_us without managed=True raises ValueError. A managed CapacityRing doubles when 85 percent full and halves when 10 percent full, within 64 to 65536 slots and no sooner than 100 ms after its last resize; sidecar_morphs and sidecar_prewarms count what it did. A managed LocaleRing moves to the locale request_locale last asked for, no sooner than 250 ms after its last move, and sidecar_migrations counts the moves. On a managed LocaleRing, migrate_to is also what its sidecar is asked to keep, so the sidecar does not move the ring back.

The sidecar

One sidecar serves the process: a scan thread per NUMA node drains the observation ring of every registered object into its stats and, for an object registered with a policy, asks that policy which tag the object should run at.

NameWhat it is
obj.observe(), obj.observe(policy)Registers the object with the process’s sidecar and returns a Registration: stats() reads what the sidecar has drained as an InstanceStats, tag the tag the object runs at, and close(), leaving a with block, or collection unregisters it. An object has one registration at a time.
AdaptiveAn adaptive object of the caller’s own: record puts an operation in its ring, and tag is what its policy last answered.
subetha.sidecarThe sidecar as a whole: scan_now() has every scan thread scan and waits for it, instance_count(), max_instances() and set_max_instances() read and move the cap of 10,000 objects, and node_count() counts the scan threads.

Twenty-six classes have observe: Arena, Atomic, BitVec, BlockedBloomFilter, BloomFilter, BroadcastRing, CountMinSketch, EpochBarrier, FenceClock, Graph, HandleTable, HashMap, Heartbeat, Histogram, HyperLogLog, LeaderElection, OwnerLease, RWLock, RateLimiter, Reservoir, Ring, Semaphore, TimePointTile, TopologyMap, Universal and VersionChain, and Adaptive does too. Every structure keeps its own layout whatever tag a policy answers, so on a structure a policy moves a tag nothing reads; a policy that decides something is one on an Adaptive, whose caller reads tag to choose how it works.

A policy is a callable taking the stats and the current tag and returning a tag, an integer from 0 to 4294967295, or None to stay. It runs on the sidecar’s scan thread and is asked after every scan that drained something new. One that raises, or returns anything else, leaves the tag where it was and is counted by policy_errors, with last_policy_error holding the exception.

Order

ClassWhat it is
OrderedReceiverDelivers a ring’s items in the order their senders made them. Built by Ring.ordered_receiver, which refuses a ring with no stamps rather than returning nothing forever. It picks its own strategy and strategy says which.
ReorderWindowThe same window over items from anywhere else: push with a stamp, take the smallest first, and the window widens itself when something arrives further out of order than it covered.

OrderedReceiver.recv answers None for two different reasons, the window still filling and the ring being empty, so drain is the shape to reach for: it takes the ring and the held-back tail in one crossing.

Shared state

Atomic (load, store, add, subtract, the bitwise operations, swap and compare-and-exchange, each taking a named memory ordering), Cell, Vec, Slab, HashMap, BTreeMap, LinkedList, Deque, Stack, Arena, Region and FrameRegion.

Region is the one that carries a buffer: memoryview(region) is the mapping itself. It is the only family that does, because the others guard each slot with a seqlock and a raw view would bypass the retry and hand back torn elements.

State with a history

ClassWhat it is
VersionChainOne value’s versions. A reader holding a version number gets the value as it stood then.
VersionedSlab, SlabPinNumbered slots each keeping their recent history, read through a pin that fixes one epoch for the whole slab. history reports each version with the epochs it was current between.
VersionedMap, MapPinAn ordered index from one unsigned integer to another, scanned through a pin that sees one unchanging view while writers carry on.
LanedMap, LaneClaim, LanedPinThat map split across several trees so several writers work at once. A writer claims a lane; reading takes none, and a scan merges every lane in key order.
EpochsThe epoch table the above share.

Two names here mean what the Rust says rather than what they sound like. void_epoch is not reclamation: it undoes the writes stamped at exactly one epoch and makes what they superseded current again, for a writer that died partway through. sweep and sweep_slot are the ones that reclaim, and a live pin holds the horizon still, so a sweep during a scan frees only what was already unreachable when the scan began.

A key belongs to one lane of a LanedMap for its whole life, because a lane is a separate tree. Writing it through another raises WrongLane naming the right one, rather than reporting a row gone that no reader has stopped seeing.

Coordination

ClassWhat it is
RWLock, HoldRead and write holds, as context managers rather than tokens. read_for and write_for give up after a timeout and answer None.
Semaphore, PermitHoldA counting semaphore across processes. acquire_for gives up after a timeout.
CondvarWhose predicate really is called from inside the wait.
LazyValueComputed once across every process.
OwnerLease, LeaseHoldOne process owns a small value and another takes it over when that one dies.
Heartbeat, EpochBarrier, LeaderElection, HolderTable, FenceClock, SharedArc, Notifier, NotifierSetThe rest of the coordination surface.

The waiting forms that take a timeout sleep rather than spin, so a long wait costs no processor, and they never wait past the deadline even if whoever holds the lock never gives it back. Only the bounded forms are here: an unbounded park sleeps until something signals it, which against a peer that releases without signaling would never wake.

A lease changes hands two ways and a caller has to know them apart. A process whose id is lower than the owner’s takes it on the spot, beaten or not, which is what settles who leads when several processes start together. A process whose id is higher waits out the grace period, counted in epochs the holders step themselves with tick_epoch; nothing steps it on its own.

Probabilistic

BloomFilter, BlockedBloomFilter, CountMinSketch, HyperLogLog, Histogram, RateLimiter, BitVec and Reservoir.

BlockedBloomFilter puts every bit for one item in a single cache line, so a lookup is one miss rather than several, and its suggest works the size and hash count out from the item count and the rate of wrong yeses that can be lived with. Reservoir keeps a bounded unbiased sample of a stream of any length; a refused value there is how the sample stays unbiased.

BitVec.toggle answers the new value, while set and clear answer the previous one.

Specialist

ClassWhat it is
HandleTableValues reached by a handle carrying both the place and how often it has been reused, so a handle to a removed value does not follow the place to its new occupant.
TimePointTileSixteen values each stamped with the version it was written at, compared against a reader’s version in a handful of vector instructions.
TowerValues reached by a path down through levels, where each level checks that it still points at the next number in the path and refuses at the first that does not.
GraphNodes and the directed edges between them, in mapped files.
TopologyMapCounts who sends to whom and reads the shape off those counts: point_to_point, broadcast_tree or all_to_all_mesh. Reading a recommendation does not publish it.
UniversalA set stored as a list while that is quicker to walk and moved to a map when it is not.
QosPolicy, QosSnapshotWhat a stream needs written down, and what follows: where the bytes should live, and whether the ordering must change.

QosPolicy.ordering is set by the caller and never inferred from the traffic, because whether a reader needs one sender’s order or every sender’s order is something only the application knows.

Reaching another machine

SensSender and SensReceiver are the two ends of a link that keeps working as the network gets worse. The link sends more than the items, so a reader rebuilds what was lost without asking again, and it changes between the sliding code and the block code as the measured loss moves, without either end reconnecting.

reader = subetha.SensReceiver(("0.0.0.0", 9000), max_item_size=1024)
writer = subetha.SensSender(("0.0.0.0", 0), ("10.0.0.5", 9000), max_item_size=1024)
writer.send_many([b"one", b"two"])

# Items do not arrive one at a time. poll answers whatever the link
# could rebuild this time round, which may be nothing.
for item in reader.poll():
    handle(item)

Two things a reader has to know. Under the block code items are grouped in eights and none of a group is delivered until enough of it has arrived, so a run shorter than a group waits. And SensSender.loss is None until the reader has reported anything, which is not the same as zero.

The bridges, and what a wheel has

TcpBridgeClient / TcpBridgeServer and QuicBridgeClient / QuicBridgeServer carry a whole ring to another host. Both are off by default, because each brings a network stack with it and a process sharing memory on one host needs none of it:

maturin build --release --features tcp-bridge,quic-bridge

A wheel therefore exports what it was built with, so a name a caller cannot find can be told from a feature left out:

if "tcp" in subetha.transports:
    server = subetha.TcpBridgeServer(incoming_ring, ("0.0.0.0", 9100))

subetha.OPTIONAL_BY_TRANSPORT names the classes each optional transport brings.

A QUIC certificate is made as two pieces of bytes by generate_self_signed_cert: the reading end holds both and proves itself with them, the sending end holds the certificate alone and checks the reading end against it. Carrying it between hosts is the caller’s to arrange.

Sensing

Eight classes that are fed measurements and answer what they worked out from them. They hold no shared memory and touch no network: they are the arithmetic a transport does on its own numbers, which is why they are worth having whether or not a SubEtha link produced the numbers.

ClassWhat it answers
LossKindWhether a loss came from noise or from a queue overflowing. The two want opposite responses, so reading one as the other is expensive.
LossBurstsWhether losses arrive alone or in runs, and how long a run lasts. The same loss rate needs different redundancy depending on the answer.
TimingJitter, spacing, and whether delay is climbing. trend_debiased takes a constant difference between two clocks out, so unsynchronized clocks do not read as a filling queue.
RoundTripShapeWhether round trips fall into two groups, which is what a radio retry looks like from outside.
PeriodicityWhether interference arrives on a beat, and when the next spike is due, so redundancy can rise before it rather than after.
CapacityThe narrowest link on the path and what is free on it, from probes sent in pairs and in trains.
ForecastWhat the next interval is likely to carry.
PathChangesWhether the route moved under the traffic, which otherwise reads as congestion.

Every one answers None rather than a number until it has seen enough to say anything. None and zero are different answers.

Values that ride beside a pointer

ClassWhat it is
TinyBloomA whole bloom filter in one machine word, small enough to sit next to a pointer and be read in the same cache line. Its state crosses as a single number, so it travels anywhere an integer does. About eight keys.
FineBloomThe same in four words, for about sixty-four keys.
ClockWall-clock time that still orders two events sharing a reading, which a bare timestamp cannot. merge is what a receiver does with a sender’s clock so the received event orders after its cause.
CausalClockOne count per participant, answering before, after, equal, or concurrent. Concurrent is the answer a timestamp can never give.

The pointer types themselves are not here. Each holds a raw pointer or a reference count into Rust memory, and Python has no such value to point at; binding them would mean handing the interpreter raw pointers.

Asyncio

subetha.aio is a small pure-Python module for coroutines. The Rust side has a reactor, an executor and a task pool, and none of it is bound because none of it can be: every entry point takes or returns a Rust future. What asyncio needs is not Rust’s executor but a way to wait without blocking its loop.

from subetha import aio

item = await aio.recv(channel, timeout=5)
answer = await aio.with_write_lock(lock, lambda: do_the_work())

Only some of the surface can be awaited soundly, and the module says which. A queue can, because what comes back is bytes. A hold cannot be handed back, because it belongs to the thread that took it, so with_permit, with_read_lock and with_write_lock run the caller’s work on that same thread instead. A SensReceiver cannot leave its thread at all, so poll asks on this one and yields between tries.

Threads

The module declares that it does not need the interpreter lock, so on a free-threaded interpreter the lock stays off when it is imported and several threads really do run inside these calls at once.

That declaration is backed by tests/test_threading.py, which runs every call that releases the interpreter under eight threads: the counters, both kinds of lock hold, the semaphore’s permits, a ring several senders share, a pinned scan running beside writers, and a buffer view held across other threads’ work. It passes on CPython 3.14 free-threaded with the lock reported off.

subetha.free_threaded says which build is installed. A wheel for a free-threaded interpreter is a separate wheel, because the stable ABI does not cover free-threading until 3.15:

maturin build --release --no-default-features

Two things the binding does that the Rust does not

Every value is boxed. Python’s object allocator aligns to sixteen bytes, HandshakeHeader is cache-line aligned, and many of the primitives embed one. A class holding such a value inline compiles, imports, and then faults inside its constructor on the first aligned store. A compile-time assertion per class enforces the boxing.

A borrowed guard is declared before the handle that keeps its target alive. Rust drops struct fields in declaration order, so a pin or a receiver written the other way round released the object it borrowed from before the guard that still had to use it. Both SlabPin and OrderedReceiver have a test that collects the owning object and then uses the borrower.