Skip to content
Every class, in full

Every class, in full

Every class, in full

All 94 classes the extension exports, with every method, its full signature and what it answers. Generated from the type stub by crates/subetha-py/tools/export_reference.py; the stub is held against the compiled module by tests/test_surface.py, so a name here is a name that ships.

A method answering T | None uses None for an ordinary absent answer rather than a fault, the same way the other bindings use their empty value.

Signatures come from the stub and descriptions from the Rust doc comments. help() in a REPL shows both, because PyO3 derives __text_signature__ from the #[pyo3(signature = ...)] attribute; this page exists to read the surface whole rather than a name at a time.

Each description here is the opening paragraph of the method’s doc comment. Many carry more than that, on what an answer means or what the call does not do, and help() shows all of it.

All 728 methods carry a description.

Contents

Adaptive , AdaptiveQueue , Arena , Atomic , BTreeMap , BitVec , BlockedBloomFilter , BloomFilter , BroadcastRing , Capacity , CapacityRing , CausalClock , Cell , Channel , Clock , Condvar , CountMinSketch , Deque , EpochBarrier , Epochs , FenceClock , FineBloom , Forecast , FrameRegion , Graph , HandleTable , HashMap , Heartbeat , Histogram , Hold , HolderTable , HyperLogLog , InstanceStats , KvMap , LamportConsumer , LamportProducer , LaneClaim , LanedMap , LanedPin , LazyValue , LeaderElection , LeaseHold , LinkedList , LocaleRing , LossBursts , LossKind , LruCache , MapPin , MpmcConsumer , MpmcProducer , MpscConsumer , MpscProducer , Notifier , NotifierSet , OrderedReceiver , OwnerLease , PathChanges , Periodicity , PermitHold , PubSub , QosPolicy , QosSnapshot , QuicBridgeClient , QuicBridgeServer , RWLock , RateLimiter , Region , Registration , ReorderWindow , Reservoir , Ring , RoundTripShape , Semaphore , SensReceiver , SensSender , SharedArc , Slab , SlabPin , SpscRing , Stack , Subscriber , TcpBridgeClient , TcpBridgeServer , TimePointTile , Timing , TinyBloom , TopologyMap , Tower , Universal , Vec , VersionChain , VersionedMap , VersionedSlab , WorkQueue

Adaptive

AttributeType
tagint
MethodKindWhat it does
`observe(self, policy: _PolicyNone=None) -> Registration`method
record(self, op_kind: int, latency_ticks: int=0, contended: bool=False, empty: bool=False) -> boolmethodRecord one operation of kind op_kind that took latency_ticks, in whatever unit the policy reads. Kinds 1 to 6 are counted apart, 7 and above share one count, and 0 is unspecified. contended marks an operation that took a slow path, which is what contention_rate counts, and empty one that found nothing. False when the ring did not take it: before the object is observed, or while the ring is full.
set_tag(self, tag: int) -> NonemethodMove the tag from the caller’s own code.
__init__(self) -> NonemethodA fresh object at tag 0. It records nothing until it is observed.

AdaptiveQueue

AttributeType
max_item_sizeint
orderingstr
inversionsint
shapestr
shape_generationint
traffictuple[int, float]
MethodKindWhat it does
change_shape_to(self, shape: str) -> NonemethodMove to a named shape, whatever the traffic says.
`maybe_change_shape(self) -> strNone`method
`recv(self) -> bytesNone`method
`recv_for(self, timeout: floatNone=None) -> bytesNone`
recv_many(self, max_items: int=256) -> list[bytes]methodEverything waiting, up to max_items, in one crossing.
send(self, item: bytes) -> boolmethodSend one item, answering False when it is full.
`send_for(self, item: bytes, timeout: floatNone=None) -> bool`method
send_many(self, items: Sequence[bytes]) -> intmethodSend a run of items as one batch, which is also what tells the queue the traffic comes in batches and may be worth a different shape.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Whatever shape it changed into inside the block is the one it is still in afterward, and anything queued stays queued.
`init(self, path: str, capacity: int=1024, senders: int=1, readers: int=1, ordering: str=‘per_producer’, auto_order: floatNone=None) -> None`method

Arena

AttributeType
capacity_bytesint
remaining_bytesint
used_bytesint
writablebool
MethodKindWhat it does
get(self, reference: int) -> strmethodResolve a reference to its string.
get_bytes(self, reference: int) -> bytesmethodResolve a reference to its bytes, for anything that is not text.
get_many(self, references: Sequence[int]) -> list[str]methodResolve a run of references in one crossing.
`intern(self, value: str) -> intNone`method
`intern_bytes(self, value: bytes) -> intNone`method
intern_many(self, values: Sequence[str]) -> list[int]methodIntern a run of strings in one crossing, stopping at the first the arena cannot take. Returns the references it managed.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity_bytes: int) -> ArenastaticmethodAttach to an arena that already exists, able to intern into it. Raises OSError when it is not there. capacity_bytes must be the one it was created with.
open_read_only(path: str, capacity_bytes: int) -> ArenastaticmethodAttach for reading only, so intern cannot be called and writable answers False.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. References handed out inside the block stay valid after it: they name bytes in the file, not anything this object owns.
__init__(self, path: str, capacity_bytes: int) -> NonemethodObtain the arena at path holding capacity_bytes of interned text, creating it when the file does not exist.

Atomic

MethodKindWhat it does
compare_exchange(self, expected: int, new: int, order: str='seq_cst') -> intmethodPut new in only if the value is still expected, and answer what was there either way.
fetch_add(self, value: int=1, order: str='seq_cst') -> intmethodAdd, and answer what the value was before. Adding past what a sixty-four bit number holds wraps round, as it does in Rust and in C.
fetch_add_many(self, count: int, value: int=1, order: str='seq_cst') -> intmethodAdd value count times and return the value before the run.
fetch_and(self, value: int, order: str='seq_cst') -> intmethodClear every bit that is not set in value, and answer what the value was before.
fetch_or(self, value: int, order: str='seq_cst') -> intmethodSet every bit that is set in value, and answer what the value was before.
fetch_sub(self, value: int=1, order: str='seq_cst') -> intmethodSubtract, and answer what the value was before. Taking more than the value holds wraps round to the top, which is what a counter going below zero means here.
fetch_xor(self, value: int, order: str='seq_cst') -> intmethodFlip every bit that is set in value, and answer what the value was before.
load(self, order: str='seq_cst') -> intmethodRead the value.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str) -> AtomicstaticmethodAttach to an atomic that already exists, leaving its value alone.
store(self, value: int, order: str='seq_cst') -> NonemethodWrite the value, discarding whatever was there without reading it. Use swap when the previous value matters, or compare_exchange when the write should land only if nobody else got there first.
swap(self, value: int, order: str='seq_cst') -> intmethodPut value in and answer what was there, in one step nothing else can get between.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The mapping goes when the last reference to this object goes, not at the end of the block, and the file outlives the process either way.
__init__(self, path: str, init: int=0) -> NonemethodObtain the atomic at path, creating it holding init when the file does not exist and attaching to its live value when it does.

BTreeMap

AttributeType
capacityint
key_sizeint
nodesint
value_sizeint
MethodKindWhat it does
clear(self) -> NonemethodRemove every entry, so the map is empty and its room is free again. The capacity and both widths are unchanged, and every process mapping the file sees it.
`first(self) -> tuple[bytes, bytes]None`method
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the entries without this; flushing is about surviving a machine that stops.
`get(self, key: bytes) -> bytesNone`method
`insert(self, key: bytes, value: bytes) -> bytesNone`method
insert_many(self, pairs: Sequence[tuple[bytes, bytes]]) -> intmethodInsert or replace a run of pairs, stopping at the first the map has no room for, and answer how many landed.
`last(self) -> tuple[bytes, bytes]None`method
open(path: str, capacity: int, key_size: int, value_size: int, tag: int=0) -> BTreeMapstaticmethodAttach to a map that already exists, raising OSError when it does not. The capacity, both widths and the tag must be the ones it was created with.
`remove(self, key: bytes) -> bytesNone`method
__contains__(self, key: bytes) -> boolmethodWhether the key is present, without bringing its value back over the boundary. Cheaper than get when the value is not wanted.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The entries stay in the file for whoever attaches next, and nothing is flushed: call flush where durability matters.
__init__(self, path: str, capacity: int, key_size: int, value_size: int, tag: int=0) -> NonemethodObtain the map at path holding capacity entries of key_size and value_size bytes, creating it when the file does not exist. A capacity below one is a ValueError.
__len__(self) -> intmethodHow many entries the map holds, which is not the capacity and not nodes: a node is a block of the tree and holds several entries.

BitVec

AttributeType
capacity_bitsint
MethodKindWhat it does
clear(self, index: int) -> boolmethodClear bit index and report what it was.
get(self, index: int) -> boolmethodWhether bit index is set, without changing it. An index past the capacity is an OSError.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity_bits: int) -> BitVecstaticmethodAttach to a bit vector that already exists, raising OSError when it does not. capacity_bits must be the one it was created with.
set(self, index: int) -> boolmethodSet bit index and report what it was.
set_range(self, lo: int, hi: int) -> NonemethodSet every bit from lo up to but not including hi, in one call rather than one per bit.
toggle(self, index: int) -> boolmethodFlip bit index and report what it now is.
__getitem__(self, index: int) -> boolmethodThe same answer get gives, so bits[3] reads bit three. There is no slicing and no negative indexing: an index is a bit number.
__init__(self, path: str, capacity_bits: int) -> NonemethodThe constructor asserts a capacity of at least one bit, so that is refused here first.
__len__(self) -> intmethodHow many bits the vector addresses, which is its capacity and not a count of the bits that are set. It never changes, so len on a bit vector is a constant rather than a measurement.

BlockedBloomFilter

AttributeType
blocksint
hashesint
MethodKindWhat it does
clear(self) -> NonemethodEmpty the filter.
contains(self, item: bytes) -> boolmethodAs in, spelled out.
contains_many(self, items: Sequence[bytes]) -> list[bool]methodAsk about a run of items in one crossing.
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the bits without this; flushing is about surviving a machine that stops.
insert(self, item: bytes) -> NonemethodAdd an item, setting its bits inside a single cache-line block.
insert_many(self, items: Sequence[bytes]) -> intmethodAdd a run of items in one crossing.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, bits: int, hashes: int) -> BlockedBloomFilterstaticmethodAttach to a filter another holder made, with the size and hash count it was made with.
reset(path: str, bits: int, hashes: int) -> BlockedBloomFilterstaticmethodEmpty the filter and remake it at this size, throwing away everything in it.
suggest(items: int, false_positive_rate: float) -> tuple[int, int]staticmethodThe size and hash count for holding items with no more than false_positive_rate of wrong yeses. Hand both to the constructor.
__contains__(self, item: bytes) -> boolmethodFalse means the item was definitely never added. True means it probably was.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The bits stay set for whoever attaches next, and nothing is flushed on the way out.
__init__(self, path: str, bits: int, hashes: int) -> Nonemethodbits is the size of the filter and hashes how many bits each item sets. suggest works both out from the number of items and the rate of wrong yeses that can be lived with.

BloomFilter

AttributeType
false_positive_ratefloat
n_bitsint
n_hashesint
MethodKindWhat it does
clear(self) -> NonemethodClear every bit, so the filter is empty again and false_positive_rate goes back to nothing. The size and the hash count are unchanged, and every process mapping the file sees it.
contains(self, item: bytes) -> boolmethodFalse means definitely absent, True means probably present at about false_positive_rate. The same answer in gives, named so a caller can pass it around rather than write the operator.
contains_many(self, items: Sequence[bytes]) -> list[bool]methodLook a run of items up in one crossing.
insert(self, item: bytes) -> NonemethodAdd an item.
insert_many(self, items: Sequence[bytes]) -> NonemethodInsert a run of items in one crossing.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, n_bits: int, n_hashes: int) -> BloomFilterstaticmethodAttach to a filter that already exists, raising OSError when it does not. n_bits and n_hashes must be the ones it was created with, because together they decide which bits an item touches.
suggest_config(n_items: int, false_positive_rate: float) -> tuple[int, int]staticmethodThe bits and hash count for n_items at a false-positive rate of p, so a caller sizes the filter from what it means rather than from arithmetic it has to do itself.
__contains__(self, item: bytes) -> boolmethodFalse means it is definitely absent. True means it is probably present, at about false_positive_rate.
__init__(self, path: str, n_bits: int, n_hashes: int) -> NonemethodObtain the filter at path with n_bits bits and n_hashes hash functions, creating it when the file does not exist. Either being zero is a ValueError.

BroadcastRing

AttributeType
active_consumersint
capacityint
payload_sizeint
producer_positionint
MethodKindWhat it does
`lag(self, consumer: int) -> intNone`method
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int) -> BroadcastRingstaticmethodAttach to a broadcast ring that already exists, raising OSError when it does not. capacity must be the one it was created with.
push(self, item: bytes) -> boolmethodPublish one item, which every registered consumer will see.
push_many(self, items: Sequence[bytes]) -> intmethodPublish a run of items, stopping at the first that will not fit, and answer how many landed. A count below what was handed in means the rest were not published and are still the caller’s to keep.
`recv(self, consumer: int) -> bytesNone`method
recv_many(self, consumer: int, max_items: int) -> list[bytes]methodRead up to max_items for consumer in one call.
register_consumer(self) -> intmethodTake a consumer position. Every registered consumer sees every item published after it registered.
unregister_consumer(self, consumer: int) -> NonemethodGive up a consumer position, so the producer stops holding slots for it.
wait_for_consumers(self, want: int, timeout: float=5.0) -> intmethodWait until want consumers have registered, and answer how many there are when the wait ends.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. A consumer position taken inside the block is still registered after it: give it up with unregister_consumer rather than relying on the block to do it.
__init__(self, path: str, capacity: int) -> NonemethodObtain the broadcast ring at path holding capacity slots, creating it when the file does not exist and attaching to what is already there when it does.

Capacity

AttributeType
available`float
link_capacity`float
samplestuple[int, int]
train_rate`float
MethodKindWhat it does
observe_pair(self, index: int, arrived_microseconds: float) -> NonemethodFeed the arrival of one probe of a pair, numbered within the pair, in microseconds.
observe_train(self, arrived_microseconds: float) -> NonemethodFeed the arrival of one probe of a train, in microseconds.
reset(self) -> NonemethodForget everything and start again, which is what a caller does when the path may have changed.
__init__(self, probe_bytes: int=1400) -> Nonemethodprobe_bytes is how big each probe is.

CapacityRing

AttributeType
capacityint
inversionsint
managedbool
ordering_mode`str
pin_generationint
sidecar_morphsint
sidecar_prewarmsint
stale_popsint
stampedbool
warm_capacity`int
warm_hitsint
MethodKindWhat it does
clear_warm(self) -> NonemethodDrop the prewarmed backing, giving back its memory and its file.
morph_to(self, capacity: int) -> NonemethodResize. Items already in the ring stay readable through the old backing until they have been taken, so nothing in flight is lost.
`open(path: str, capacity: int, max_producers: int=1, max_consumers: int=1, stamped: bool=False, managed: bool=False, scan_interval_us: intNone=None) -> CapacityRing`staticmethod
prewarm(self, capacity: int) -> NonemethodBuild a backing of capacity ahead of needing it, so the morph that switches to it does not pay for the allocation.
`recv(self, consumer: int) -> bytesNone`method
recv_many(self, consumer: int, max_items: int) -> list[bytes]methodTake up to max_items for consumer in one crossing, stopping early when the ring runs dry. An empty list means there was nothing, which is the same answer recv gives as None.
register_consumer(self) -> intmethodTake a consumer position, which every recv names.
register_producer(self) -> intmethodTake a producer position, which every send names. A ring built with one producer has exactly one to take, and asking past max_producers raises rather than answering.
send(self, producer: int, item: bytes) -> boolmethodSend one item from producer. False means the ring is full, which is an answer rather than a failure; an item too long for a slot raises.
send_many(self, producer: int, items: Sequence[bytes]) -> intmethodSend a run of items from one producer, stopping at the first that will not fit, and answer how many landed. A count short of what was handed in means the rest were not sent and are still the caller’s to hold.
set_ordering_mode(self, mode: str) -> NonemethodSet the ordering discipline across the current backing and every superseded one still draining, so a reader walking both applies one discipline.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. A prewarmed backing is still held afterward, and a producer or consumer position taken inside is still registered: clear_warm gives back the first, and nothing gives back the second.
`init(self, path: str, capacity: int, max_producers: int=1, max_consumers: int=1, stamped: bool=False, managed: bool=False, scan_interval_us: intNone=None) -> None`method

CausalClock

AttributeType
nodesint
countslist[int]
MethodKindWhat it does
compare(self, other: CausalClock) -> strmethodHow this stands to another: before, after, equal, or concurrent when neither caused the other.
concurrent_with(self, other: CausalClock) -> boolmethodWhether neither caused the other.
count(self, node: int) -> intmethodOne participant’s count.
happened_before(self, other: CausalClock) -> boolmethodWhether this happened before the other, which is the same as compare answering before.
merge(self, other: CausalClock) -> CausalClockmethodThe clock a receiver should hold after taking other in: the higher of each count.
tick(self, node: int) -> CausalClockmethodStep this participant’s own count, which is what it does when something happens to it.
__init__(self) -> NonemethodAll counts at zero.

Cell

AttributeType
value_sizeint
versionint
MethodKindWhat it does
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it.
get(self) -> bytesmethodRead the value, always value_size bytes.
open(path: str, value_size: int) -> CellstaticmethodAttach to a cell that already exists, raising OSError when it does not. value_size must be the one it was created with, because it is part of the layout rather than a hint.
set(self, value: bytes) -> NonemethodReplace the value and step version, so a reader watching that number sees this write happened.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Nothing is flushed on the way out: call flush where durability matters.
__init__(self, path: str, value_size: int) -> NonemethodObtain the cell at path holding value_size bytes, creating it when the file does not exist and attaching to the value already there when it does.

Channel

AttributeType
max_item_sizeint
MethodKindWhat it does
open(path: str, capacity: int=1024) -> ChannelstaticmethodAttach to a channel another holder made, with the capacity it was made with.
`recv(self) -> bytesNone`method
`recv_for(self, timeout: floatNone=None) -> bytesNone`
recv_many(self, max_items: int=256) -> list[bytes]methodEverything waiting, up to max_items, in one crossing.
send(self, item: bytes) -> boolmethodSend one item, answering False when the channel is full.
`send_for(self, item: bytes, timeout: floatNone=None) -> bool`method
send_many(self, items: Sequence[bytes]) -> intmethodSend a run of items in one crossing, answering how many went. A short answer means the channel filled.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing the channel, and let an exception through. Nothing is signaled to the other end: a reader waiting on it goes on waiting, so a sender that means to say it is done has to say so in the messages it sends.
__init__(self, path: str, capacity: int=1024, senders: int=1, readers: int=1) -> Nonemethodcapacity is the number of items in flight, rounded up to a power of two. senders and readers are how many are expected at once, which is what the shape underneath is picked from.

Clock

AttributeType
logicalint
physicalint
MethodKindWhat it does
advance(self, physical: int) -> ClockmethodThe next reading given a new physical time. The count steps rather than resetting when the physical time has not moved.
merge(self, received: Clock, physical: int) -> ClockmethodThe reading a receiver should take, given what arrived and what its own clock says. Orders the received event after whatever caused it.
now() -> ClockstaticmethodA reading taken now, in microseconds since the epoch, with the count at zero.
__init__(self, physical: int=0, logical: int=0) -> NonemethodA reading built from the two parts given, both defaulting to zero.

Condvar

AttributeType
generationint
MethodKindWhat it does
notify_all(self) -> intmethodWake every waiter.
notify_one(self) -> intmethodWake at most one waiter. The caller is responsible for having made the predicate true first, which is the contract every condition variable has.
open(path: str) -> CondvarstaticmethodAttach to a condition that already exists, raising OSError when it does not.
`wait_for(self, predicate: Callable[[], object], timeout: floatNone=None) -> bool`method
__init__(self, path: str) -> NonemethodObtain the condition at path, creating it when the file does not exist and attaching to the live one when it does.

CountMinSketch

AttributeType
depthint
total_insertsint
widthint
MethodKindWhat it does
estimate_count(self, item: bytes) -> intmethodHow often it thinks it has seen item. Never less than the true count, sometimes more.
estimate_many(self, items: Sequence[bytes]) -> list[int]methodOne estimate per item, in the order they were given, from a single crossing rather than one per item. Each answer has the same never-under, sometimes-over property as estimate_count.
insert(self, item: bytes) -> NonemethodRecord one occurrence of item.
insert_many(self, items: Sequence[bytes]) -> intmethodRecord one occurrence of each item in one crossing, and answer how many were handed in. Nothing here can refuse, so the count is always the length of the sequence.
insert_n(self, item: bytes, count: int) -> NonemethodAdd count occurrences at once rather than calling insert that many times.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, depth: int, width: int) -> CountMinSketchstaticmethodAttach to a sketch that already exists, raising OSError when it does not. depth and width must be the ones it was created with, because together they decide which cells an item touches.
reset(self) -> NonemethodZero every cell and the insert total, so the sketch is empty again. Every process mapping the file sees it.
suggest_config(epsilon: float, delta: float) -> tuple[int, int]staticmethodThe depth and width for an error of epsilon with confidence delta, so a caller sizes the sketch from what it needs.
__init__(self, path: str, depth: int, width: int) -> NonemethodObtain the sketch at path with depth rows of width cells, creating it when the file does not exist. Either being zero is a ValueError.

Deque

AttributeType
approx_lenint
capacityint
element_sizeint
MethodKindWhat it does
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the writes without this; flushing is about surviving a machine that stops.
open_as_thief(path: str, element_size: int, alignment: int=1, tag: int=0) -> DequestaticmethodAttach as a thief: this handle steals from the deque another process owns, and does not push or pop.
`pop(self) -> bytesNone`method
push(self, item: bytes) -> boolmethodPush at the owner’s end. False means it is full.
push_many(self, items: Sequence[bytes]) -> intmethodPush a run of items in one crossing, stopping at the first that will not fit, and answer how many landed. They go on the owner’s end, the end pop takes from and steal does not.
`steal(self) -> bytesNone`method
steal_many(self, max_items: int) -> list[bytes]methodSteal up to max_items from the far end in one crossing, stopping early when there is nothing left to take. An empty list means the deque was empty or the owner won every race for the last items, which a thief cannot tell apart and does not need to.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Nothing is drained on the way out: whatever is in the deque stays there for whoever attaches next, including for a thief still stealing from the other end.
__init__(self, path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> NonemethodThe capacity must be a power of two, which the ring arithmetic relies on; it is checked here so a bad one reads as an argument error rather than a refusal from deeper down.

EpochBarrier

AttributeType
arrivedint
current_epochint
live_peersint
MethodKindWhat it does
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, heartbeat: Heartbeat, grace_epochs: int=3) -> EpochBarrierstaticmethodAttach to a barrier that already exists, raising OSError when it does not. The heartbeat table and the grace period are this participant’s own, given here the same way they are at creation.
snapshot(self) -> tuple[int, int]methodThe epoch and how many have arrived at it, read together so the two cannot disagree.
`wait(self, epoch: int, timeout: floatNone=None, quorum: intNone=None) -> bool`
__init__(self, path: str, heartbeat: Heartbeat, grace_epochs: int=3) -> NonemethodObtain the barrier at path, creating it when the file does not exist.

Epochs

AttributeType
capacityint
nowint
open_ticketsint
MethodKindWhat it does
advance(self) -> intmethodTake the next epoch and return it, for a writer stamping a version as it supersedes the last.
claim_ticket(self) -> tuple[int, int]methodReserve the next epoch for a compound write, returning the slot holding the ticket and the epoch it took.
dead_tickets(self) -> list[int]methodThe epochs of tickets whose holders are gone, which is what makes a crashed reader stop holding reclamation up for ever.
free_dead_ticket(self, epoch: int) -> boolmethodFree a ticket whose holder died. True when one was freed.
open(path: str, capacity: int) -> EpochsstaticmethodAttach to an epoch table that already exists, raising OSError when it does not. capacity must be the one it was created with.
publish_ticket(self, slot: int) -> NonemethodGive a ticket back.
__init__(self, path: str, capacity: int) -> NonemethodObtain the epoch table at path with room for capacity open tickets, creating it when the file does not exist. A capacity of zero is a ValueError.

FenceClock

AttributeType
capacityint
shared_clock_usint
MethodKindWhat it does
get_local(self, slot: int) -> tuple[int, int]methodThis participant’s current reading, without moving it.
global_fence(self) -> tuple[int, int]methodThe latest reading across every live participant: every event any of them has recorded is at or before this.
merge(self, slot: int, physical_us: int, logical: int) -> tuple[int, int]methodFold a reading received from elsewhere into this participant’s clock, which is what makes the ordering hold across processes.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int) -> FenceClockstaticmethodAttach to a clock that already exists, raising OSError when it does not. capacity must be the one it was created with.
`register(self, pid: intNone=None) -> int`method
tick(self, slot: int) -> tuple[int, int]methodMove this participant’s clock on and return the new reading.
unregister(self, slot: int) -> NonemethodGive a participant slot back, so another process can take it. The readings it contributed stay in the shared clock: giving up the slot stops this participant advancing, it does not rewind what it already published.
__init__(self, path: str, capacity: int) -> NonemethodObtain the clock at path with room for capacity participants, creating it when the file does not exist. A capacity of zero is a ValueError.

FineBloom

AttributeType
suggested_capacityint
MethodKindWhat it does
contains(self, key: bytes) -> boolmethodAs in, spelled out.
contains_many(self, keys: Sequence[bytes]) -> list[bool]methodAsk about a run of keys in one crossing, answering one True or False per key in the order they were given.
insert(self, key: bytes) -> NonemethodAdd a key, setting its bits. Nothing can be taken out again, and the rate of wrong yeses climbs past suggested_capacity keys.
insert_many(self, keys: Sequence[bytes]) -> intmethodAdd a run of keys in one crossing, and answer how many were handed in. Nothing here can refuse, so the count is always the length of the sequence.
__contains__(self, key: bytes) -> boolmethodFalse means the key is definitely absent, True means it is probably present.
`init(self, keys: Sequence[bytes]None=None) -> None`method

Forecast

AttributeType
mean_ratefloat
next_ratefloat
MethodKindWhat it does
observe(self, bytes: int, seconds: float) -> NonemethodFeed one interval: how many bytes arrived and how long it was, in seconds.
__init__(self) -> NonemethodA forecast with nothing observed yet. It lives in this process and has no file behind it. Feed it whole intervals with observe, each one a byte count and how long it covered.

FrameRegion

AttributeType
block_countint
block_sizeint
MethodKindWhat it does
`allocate(self) -> intNone`method
free(self, index: int) -> NonemethodGive a block back. Freeing one nobody holds corrupts the pool, so only free what this process allocated and has finished with.
open(path: str, block_size: int, block_count: int) -> FrameRegionstaticmethodAttach to a frame region that already exists, raising OSError when it does not. The block size and count must be the ones it was created with: they are the layout, so a mismatch reads the blocks at the wrong offsets.
read_block(self, index: int, length: int) -> bytesmethodRead length bytes from a block.
take_block(self, index: int, length: int) -> bytesmethodRead a block and give it back in one call, which is the shape a consumer wants.
write_block(self, index: int, payload: bytes) -> NonemethodWrite payload into block index, checking first that it fits.
`write_new(self, payload: bytes) -> intNone`method
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The blocks stay as they were written for whoever attaches next.
__init__(self, path: str, block_size: int, block_count: int) -> NonemethodObtain the frame region at path holding block_count blocks of block_size bytes, creating it when the file does not exist.

Graph

AttributeType
edge_countint
max_edgesint
max_nodesint
node_countint
MethodKindWhat it does
add_edge(self, source: int, target: int, value: int) -> intmethodAdd an edge from one node to another, carrying value, and answer its index.
add_edges(self, edges: Sequence[tuple[int, int, int]]) -> list[int]methodAdd a run of edges in one crossing, each a source, a target and a value.
add_node(self, value: int) -> intmethodAdd a node carrying value, answering its index.
add_nodes(self, values: Sequence[int]) -> list[int]methodAdd a run of nodes in one crossing, answering their indexes.
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the nodes and edges without this; flushing is about surviving a machine that stops.
flush_async(self) -> NonemethodStart writing the mapping back to its file and return at once, without waiting for the write to land. Nothing is on disk yet when this returns; flush is the form that waits.
neighbors(self, source: int) -> list[tuple[int, int, int]]methodEverything reachable in one step from a node, as the edge index, the node it leads to and what the edge carries.
`node_value(self, node: int) -> intNone`method
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, max_nodes: int, max_edges: int) -> GraphstaticmethodAttach to a graph another holder made, with the sizes it was made with.
`out_degree(self, source: int) -> intNone`method
`remove_edge(self, source: int, edge: int) -> intNone`method
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Node and edge indices handed out inside the block still name the same nodes and edges after it, and nothing is flushed on the way out.
__init__(self, path: str, max_nodes: int, max_edges: int) -> NonemethodThe graph lives in two files beside the path given, one for the nodes and one for the edges.

HandleTable

AttributeType
max_value_bytesint
capacityint
MethodKindWhat it does
contains(self, handle: int) -> boolmethodAs in, spelled out.
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the entries without this; flushing is about surviving a machine that stops.
flush_async(self) -> NonemethodStart writing the mapping back to its file and return at once, without waiting for the write to land. Nothing is on disk yet when this returns; flush is the form that waits.
`get(self, handle: int) -> bytesNone`method
`get_many(self, handles: Sequence[int]) -> list[bytesNone]`method
insert(self, value: bytes) -> intmethodPut a value in and get its handle. Raises when the table is full.
insert_many(self, values: Sequence[bytes]) -> list[int]methodPut a run of values in, answering their handles in order. Stops at the first one that does not fit, so a short answer means the table filled.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int) -> HandleTablestaticmethodAttach to a table another holder made, with the capacity it was made with.
`remove(self, handle: int) -> bytesNone`method
reset(path: str, capacity: int) -> HandleTablestaticmethodEmpty the table and remake it at this capacity.
__contains__(self, handle: int) -> boolmethodWhether the handle is still live. A handle that has been given back answers False, and so does one from an earlier generation of the same slot, which is what the generation in a handle is for.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Handles taken inside the block are still live after it: they are named by the table, not owned by this object.
__init__(self, path: str, capacity: int) -> NonemethodHow many values the table holds at most. A value is at most 44 bytes.
__len__(self) -> intmethodHow many handles are live, which is not the capacity. A handle given back stops being counted at once.

HashMap

AttributeType
capacityint
key_sizeint
tombstonesint
value_sizeint
MethodKindWhat it does
clear(self) -> NonemethodRemove every entry and every tombstone, so both len and tombstones go to zero and the whole capacity is usable again. The widths are unchanged, and every process mapping the file sees it. This is the only thing that clears tombstones.
compare_exchange(self, key: bytes, expected: bytes, new: bytes) -> tuple[bool, bytes]methodReplace expected with new only if that is what is there. Returns whether it swapped and what was found.
`get(self, key: bytes) -> bytesNone`method
`get_many(self, keys: Sequence[bytes]) -> list[bytesNone]`method
`insert(self, key: bytes, value: bytes) -> strNone`method
insert_many(self, pairs: Sequence[tuple[bytes, bytes]]) -> intmethodInsert a run of pairs, stopping at the first the map refuses.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int, key_size: int, value_size: int) -> HashMapstaticmethodAttach to a map that already exists, raising OSError when it does not. The capacity and both widths must be the ones it was created with.
`remove(self, key: bytes) -> bytesNone`method
__contains__(self, key: bytes) -> boolmethodWhether the key is present, without bringing its value back over the boundary. Cheaper than get when the value is not wanted.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The entries stay in the file for whoever attaches next.
__init__(self, path: str, capacity: int, key_size: int, value_size: int) -> NonemethodObtain the map at path holding capacity entries of key_size and value_size bytes, creating it when the file does not exist.
__len__(self) -> intmethodHow many entries the map holds, which is not the capacity and does not count tombstones: a removed key stops being counted here while its marker still occupies a slot. Read tombstones for those.

Heartbeat

AttributeType
capacityint
global_epochint
MethodKindWhat it does
beat(self, slot: int) -> NonemethodSay this slot is still alive. A slot that stops beating for longer than the grace period counts as dead.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int) -> HeartbeatstaticmethodAttach to a heartbeat table that already exists, raising OSError when it does not. capacity must be the one it was created with.
`register(self, pid: intNone=None) -> int`method
`snapshot(self, slot: int) -> dict[str, int]None`method
tick_global_epoch(self) -> intmethodStep the global epoch and answer the new value, not the old one.
unregister(self, slot: int) -> NonemethodGive a slot back, so another process can register into it.
__init__(self, path: str, capacity: int) -> NonemethodObtain the heartbeat table at path with capacity slots, creating it when the file does not exist. A capacity of zero is a ValueError.

Histogram

AttributeType
boundarieslist[int]
countslist[int]
n_bucketsint
total_countint
MethodKindWhat it does
count(self, bucket: int) -> intmethodHow many values fell in one bucket, by its index. A bucket past n_buckets is an OSError. Use counts when you want them all: it costs one crossing rather than one per bucket.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, boundaries: Sequence[int]) -> HistogramstaticmethodAttach to a histogram that already exists, raising OSError when it does not. boundaries must be the ones it was created with: they are the bucket edges written into the file, so a different list reads the counts against the wrong edges.
percentile(self, p: float) -> intmethodThe value at percentile p, from the bucket boundaries, so it is as precise as the buckets are.
record(self, value: int) -> intmethodRecord one value and return the bucket it fell in.
record_many(self, values: Sequence[int]) -> intmethodRecord a run of values in one crossing.
__init__(self, path: str, boundaries: Sequence[int]) -> Nonemethodboundaries must rise and must not be empty, which is checked here so a bad one names itself.

Hold

AttributeType
heldbool
MethodKindWhat it does
release(self) -> NonemethodGive the hold back now rather than at the end of a block. Calling it twice is harmless.
__enter__(self) -> SelfmethodAnswer the same object, which is why a hold is normally taken as with lock.write() as held: rather than bound to a name.
__exit__(self, *args: object) -> boolmethodGive the lock back, and let an exception through.

HolderTable

AttributeType
capacityint
liveint
MethodKindWhat it does
`claim(self, payload: int) -> intNone`method
open(path: str, capacity: int) -> HolderTablestaticmethodAttach to a holder table that already exists, raising OSError when it does not. capacity must be the one it was created with.
`payload(self, slot: int) -> intNone`method
publish(self, slot: int, payload: int) -> NonemethodPut a payload in a slot this caller already holds.
release(self, slot: int) -> NonemethodGive a slot back, so the next claim or reserve can hand it to somebody else. A process that exits without releasing leaves its slot held, and nothing here reclaims it.
`reserve(self) -> intNone`method
__init__(self, path: str, capacity: int) -> NonemethodObtain the holder table at path with capacity slots, creating it when the file does not exist. A capacity of zero is a ValueError, because a table with no slots can serve nobody.

HyperLogLog

AttributeType
n_registersint
precisionint
MethodKindWhat it does
estimate(self) -> intmethodHow many distinct items it estimates it has seen. An estimate, not a count: that is the bargain the structure makes.
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the registers without this; flushing is about surviving a machine that stops.
insert(self, item: bytes) -> NonemethodRecord having seen item.
insert_many(self, items: Sequence[bytes]) -> intmethodInsert a run of items in one crossing, which is what this structure is usually fed: a stream rather than single items.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, precision: int=14) -> HyperLogLogstaticmethodAttach to a counter that already exists, raising OSError when it does not. precision must be the one it was created with, because it sets how many registers the file holds.
reset(self) -> NonemethodZero every register, so the counter is empty again and estimate answers nothing seen. Every process mapping the file sees it.
__init__(self, path: str, precision: int=14) -> Nonemethodprecision sets the trade between memory and accuracy: 4 is the smallest the format allows and 16 the largest.

InstanceStats

AttributeType
N_OP_KINDSint
MAX_TRACKED_THREADS_PER_KINDint
contention_opsint
last_drain_usint
migrations_triggeredint
op_kind_countslist[int]
ops_observedint
per_op_kind_distinct_countlist[int]
per_op_kind_distinct_threadslist[list[int]]
total_latency_ticksint
MethodKindWhat it does
average_latency_ticks(self) -> intmethodtotal_latency_ticks over ops_observed, and 0 before any observation.
contention_rate(self) -> floatmethodcontention_ops over ops_observed, and 0.0 before any observation.
distinct_threads_for(self, kind: int) -> intmethodDistinct producer threads seen for kind.
is_multi_thread_for(self, kind: int) -> boolmethodWhether kind was seen from two producer threads or more.
op_kind_total(self) -> intmethodThe per-kind counts summed.
ratio_of(self, kind: int, total_kinds: Sequence[int]) -> floatmethodThe count of kind over the counts of total_kinds summed, and 0.0 while those are all zero.

KvMap

MethodKindWhat it does
`get(self, key: int) -> intNone`method
`get_many(self, keys: Sequence[int]) -> list[intNone]`method
insert(self, key: int, value: int) -> boolmethodPut an entry in, answering True when the key was not there before and False when this replaced what it held.
insert_many(self, entries: Sequence[tuple[int, int]]) -> list[bool]methodPut a run of entries in, in one crossing, answering True for each key that was not there before.
__contains__(self, key: int) -> boolmethodWhether the key has a value. Reads the value and throws it away, so it costs what get costs; call get when you want the value as well rather than asking twice.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The entries stay in the file for whoever attaches next.
__init__(self, path: str, capacity: int=1024, readers: int=1, writers: int=1) -> NonemethodObtain the map at path, creating it when the file does not exist. Keys and values are both sixty-four bit integers.
__len__(self) -> intmethodHow many keys the map holds.

LamportConsumer

AttributeType
capacityint
MethodKindWhat it does
`pop(self) -> bytesNone`method
pop_many(self, max_items: int) -> list[bytes]methodTake up to max_items in one crossing, stopping early when the ring runs dry. An empty list means there was nothing, which is the same answer pop gives as None.

LamportProducer

AttributeType
capacityint
payload_sizeint
MethodKindWhat it does
push(self, item: bytes) -> boolmethodPush one item. False means the ring is full and the consumer has not caught up, which is an answer rather than a failure; an item longer than payload_size raises.
push_buffer(self, data: bytes, item_len: int) -> intmethodPush items cut out of one buffer, item_len bytes each, and answer how many landed.
push_many(self, items: Sequence[bytes]) -> intmethodPush a run of items in one crossing, stopping at the first that will not fit, and answer how many landed. A count short of what was handed in leaves the rest with the caller.

LaneClaim

AttributeType
heldbool
indexint
MethodKindWhat it does
`insert(self, key: int, value: int) -> intNone`method
`insert_at(self, key: int, value: int, born: int) -> intNone`method
`insert_many(self, entries: Sequence[tuple[int, int]]) -> list[intNone]`method
release(self) -> NonemethodGive the lane back now rather than at the end of a block. Calling it twice is harmless.
`remove(self, key: int) -> intNone`method
`remove_at(self, key: int, died: int) -> intNone`method
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodGive the lane back, and let an exception through.

LanedMap

AttributeType
held_lanesint
lanesint
MethodKindWhat it does
claim_lane(self) -> LaneClaimmethodClaim a free lane to write keys that are not in the map yet. Raises Contended when every lane is held.
claim_lane_for(self, key: int) -> LaneClaimmethodClaim the lane a key already lives in, to rewrite or remove it. Raises KeyError when no lane holds the key, and Contended when its lane is held by someone else. Retry rather than writing elsewhere: elsewhere is a different tree.
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Every lane goes at once, not only the ones this process has claimed. Another process mapping the same file sees the entries without this; flushing is about surviving a machine that stops.
`get(self, key: int) -> intNone`method
`get_many(self, keys: Sequence[int]) -> list[intNone]`method
`lane_of(self, key: int) -> intNone`method
open(directory: str, lanes: int=4, nodes_per_lane: int=256, max_pins: int=16) -> LanedMapstaticmethodAttach to a laned map another holder made, with the lane count, lane size and pin count it was made with.
pin(self) -> LanedPinmethodTake a pin, fixing one epoch to scan every lane at.
reap_dead_claims(self) -> intmethodGive back the lanes of writers whose process has gone, answering how many came back. Without this a process that died holding a lane keeps it forever.
sweep(self) -> intmethodTake away every entry marked as no longer current that nothing can still reach, across every lane. Zero means nothing could be taken.
void_epoch(self, epoch: int) -> intmethodUndo every write stamped at exactly this epoch across every lane, answering how many entries were touched. For a writer that died partway through a change spanning several lanes.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. A LaneClaim or LanedPin taken inside the block is not given back here; each belongs in a with block of its own.
__init__(self, directory: str, lanes: int=4, nodes_per_lane: int=256, max_pins: int=16) -> NonemethodThe map lives in a directory of its own, one file per lane plus the shared epochs and the claims. nodes_per_lane is how many entries each lane holds and max_pins how many readers may scan at once across all of them.
__len__(self) -> intmethodEntries across every lane, counting the ones marked as no longer current.

LanedPin

AttributeType
epochint
heldbool
MethodKindWhat it does
`get(self, key: int) -> intNone`method
`get_many(self, keys: Sequence[int]) -> list[intNone]`method
release(self) -> NonemethodGive the pin back now rather than at the end of a block.
`scan(self, low: intNone=None, high: intNone=None, limit: int=1024) -> list[tuple[int, int]]`
`scan_from(self, low: intNone=None, high: intNone=None, limit: int=1024) -> tuple[list[tuple[int, int]], int
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodGive the pin back, and let an exception through.

LazyValue

AttributeType
readybool
value_sizeint
MethodKindWhat it does
`claim(self, pid: intNone=None) -> bool`method
`get(self) -> bytesNone`method
open(path: str, value_bytes: int) -> LazyValuestaticmethodAttach to a lazy value that already exists, raising OSError when it does not. value_bytes must be the size it was created with. Attaching says nothing about whether anything has been published yet: ready answers that.
`publish(self, value: bytes, pid: intNone=None) -> bool`method
wait(self, timeout: float=30.0) -> bytesmethodWait for whoever claimed it to publish, up to timeout seconds. The interpreter is detached while waiting.
__init__(self, path: str, value_bytes: int) -> NonemethodObtain the lazy value at path holding value_bytes bytes, creating it when the file does not exist. Zero bytes is a ValueError.

LeaderElection

AttributeType
global_epochint
leader`int
termint
MethodKindWhat it does
`am_i_leader(self, pid: intNone=None) -> bool`method
`beat(self, pid: intNone=None) -> bool`method
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str) -> LeaderElectionstaticmethodAttach to an election that already exists, raising OSError when it does not.
`step_down(self, pid: intNone=None) -> bool`method
tick_epoch(self) -> intmethodStep the global epoch and answer the new value, not the old one.
`try_claim(self, pid: intNone=None, grace_epochs: int=3) -> bool`method
__init__(self, path: str) -> NonemethodObtain the election at path, creating it when the file does not exist. Creating one claims nothing: try_claim is what takes the role.

LeaseHold

AttributeType
heldbool
MethodKindWhat it does
beat(self) -> boolmethodSay this process is still here, so the claim does not lapse during a long block.
`read(self) -> bytesNone`method
release(self) -> NonemethodGive the lease back now rather than at the end of a block. Calling it twice is harmless.
write(self, value: bytes) -> boolmethodWrite the value under the lease this hold has.
__enter__(self) -> SelfmethodAnswer the same object, which is why a hold is normally taken as with lease.acquire() as held: rather than bound to a name.
__exit__(self, *args: object) -> boolmethodGive the lease back, and let an exception through.

LinkedList

AttributeType
capacityint
element_sizeint
MethodKindWhat it does
get(self, index: int) -> bytesmethodRead the node at index.
open(path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> LinkedListstaticmethodAttach to a list that already exists, raising OSError when it does not. Every layout argument describes the file rather than asking anything of it, so all of them must match.
`pop_back(self) -> bytesNone`method
`pop_front(self) -> bytesNone`method
push_back(self, value: bytes) -> intmethodAdd at the back and return the index of the node holding it.
push_back_many(self, values: Sequence[bytes]) -> list[int]methodAdd a run of values at the back in one crossing, and answer the node index of each, in order.
push_front(self, value: bytes) -> intmethodAdd at the front and return the index of the node holding it.
remove(self, index: int) -> bytesmethodUnlink the node at index and return what it held.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The nodes stay linked for whoever attaches next, and a node index taken inside the block is still valid after it.
__init__(self, path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> NonemethodObtain the list at path with room for capacity nodes of element_size bytes, creating it when the file does not exist.
__len__(self) -> intmethodHow many nodes the list holds, which is not the capacity. An empty list is falsy, so if not list: works.

LocaleRing

AttributeType
inversionsint
localestr
locale_generationint
managedbool
ordering_mode`str
sidecar_migrationsint
stampedbool
MethodKindWhat it does
migrate_to(self, locale: str) -> NonemethodMove the ring to another locale, carrying what is already in it. On a stamped ring the transfer keeps the order every sender saw; on an unstamped one the drain can interleave senders, as a shape change can. On a managed ring the move is also what its sidecar is asked to keep, so the sidecar holds the ring there.
`open(path: str, capacity: int, max_producers: int=1, max_consumers: int=1, stamped: bool=False, managed: bool=False, scan_interval_us: intNone=None) -> LocaleRing`staticmethod
`recv(self, consumer: int) -> bytesNone`method
recv_many(self, consumer: int, max_items: int) -> list[bytes]methodTake up to max_items for consumer in one crossing, stopping early when the ring runs dry. An empty list means there was nothing, which is the same answer recv gives as None.
register_consumer(self) -> intmethodTake a consumer position, which every recv names. Registered on all three backings at once, like a producer, so a migration does not lose it.
register_producer(self) -> intmethodRegister on all three backings at once, so the registration is there whichever locale is live.
request_locale(self, locale: str) -> NonemethodAsk a managed ring’s sidecar to move the ring to locale, one of anon, file or shmfs. The sidecar moves it on a later scan, once 250 ms have passed since its last move. A strict ring has no sidecar to ask and raises; migrate_to moves it.
send(self, producer: int, item: bytes) -> boolmethodSend one item from producer into whichever locale is live. False means the ring is full, which is an answer rather than a failure; an item too long for a slot raises.
send_many(self, producer: int, items: Sequence[bytes]) -> intmethodSend a run of items from one producer, stopping at the first that will not fit, and answer how many landed. A count short of what was handed in means the rest were not sent and are still the caller’s to hold.
set_ordering_mode(self, mode: str) -> NonemethodSet the ordering discipline on all three backings, so a migration does not change the discipline underneath a reader.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The ring stays in whichever locale it was migrated to, and a position taken inside the block is still registered.
`init(self, path: str, capacity: int, max_producers: int=1, max_consumers: int=1, stamped: bool=False, managed: bool=False, scan_interval_us: intNone=None) -> None`method

LossBursts

AttributeType
mean_run_length`float
samplesint
steady_loss`float
transition_rates`tuple[float, float]
MethodKindWhat it does
observe(self, lost: bool) -> NonemethodFeed one item: True when it was lost.
observe_many(self, losses: Sequence[bool]) -> intmethodFeed a run of items in one crossing.
__init__(self) -> NonemethodA model with nothing observed yet. It lives in this process and has no file behind it. Feed it with observe or observe_many, one call per item, saying whether that item was lost.

LossKind

AttributeType
congestion_sharefloat
delay_spreadfloat
MethodKindWhat it does
classify(self, gap: int, spacing_microseconds: float) -> strmethodWhat a loss of gap items with this spacing was: wireless or congestion.
observe_delay(self, microseconds: float) -> NonemethodFeed a one-way delay, in microseconds.
observe_spacing(self, microseconds: float) -> NonemethodFeed the spacing between two arrivals, in microseconds.
__init__(self) -> NonemethodA sensor with nothing observed yet. It lives in this process and has no file behind it, so two processes each keep their own. Feed it with observe_spacing and observe_delay before asking it anything.

LruCache

AttributeType
capacityint
key_sizeint
value_sizeint
MethodKindWhat it does
`get(self, key: bytes) -> bytesNone`method
`get_and_touch(self, key: bytes) -> bytesNone`method
`get_many(self, keys: Sequence[bytes]) -> list[bytesNone]`method
open(path: str, capacity: int, key_size: int, value_size: int) -> LruCachestaticmethodAttach to a cache that already exists, raising OSError when it does not. The capacity and both widths must be the ones it was created with.
put(self, key: bytes, value: bytes) -> boolmethodPut a key in at the most recent end, evicting the least recently used first if the cache is full. True means the key was already present and its value was replaced.
put_many(self, pairs: Sequence[tuple[bytes, bytes]]) -> intmethodPut a run of pairs in, and answer how many.
`remove(self, key: bytes) -> bytesNone`method
touch(self, key: bytes) -> boolmethodCount a key as used without reading it.
__contains__(self, key: bytes) -> boolmethodWhether the key is present, without counting as use. Asking does not save an entry from eviction; touch is what does that.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The entries stay in the file for whoever attaches next, which is the point of a cache that outlives the process.
__init__(self, path: str, capacity: int, key_size: int, value_size: int) -> NonemethodObtain the cache at path holding capacity entries of key_size and value_size bytes, creating it when the file does not exist. A capacity of zero is a ValueError.
__len__(self) -> intmethodHow many entries the cache holds, which is at most the capacity because a full cache evicts rather than growing.

MapPin

AttributeType
epochint
heldbool
MethodKindWhat it does
`get(self, key: int) -> intNone`method
`get_many(self, keys: Sequence[int]) -> list[intNone]`method
release(self) -> NonemethodGive the pin back now rather than at the end of a block. Calling it twice is harmless.
`scan(self, low: intNone=None, high: intNone=None, limit: int=1024) -> list[tuple[int, int]]`
`scan_from(self, low: intNone=None, high: intNone=None, limit: int=1024) -> tuple[list[tuple[int, int]], int
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodGive the pin back, and let an exception through.

MpmcConsumer

AttributeType
approx_lenint
ringsint
MethodKindWhat it does
`pop(self) -> bytesNone`method
pop_many(self, max_items: int) -> list[bytes]methodTake up to max_items in one crossing from this consumer’s own rings, stopping early once they are empty. An empty list means there was nothing, which is the same answer pop gives as None.

MpmcProducer

AttributeType
capacityint
MethodKindWhat it does
push(self, item: bytes) -> boolmethodPush one item into this producer’s own ring.
push_many(self, items: Sequence[bytes]) -> intmethodPush a run of items in one crossing, stopping at the first that will not fit, and answer how many landed. A count short of what was handed in leaves the rest with the caller.

MpscConsumer

AttributeType
approx_lenint
producersint
MethodKindWhat it does
`pop(self) -> bytesNone`method
pop_many(self, max_items: int) -> list[bytes]methodTake up to max_items in one crossing, stopping early once every ring is empty. An empty list means there was nothing, which is the same answer pop gives as None.

MpscProducer

AttributeType
capacityint
payload_sizeint
MethodKindWhat it does
push(self, item: bytes) -> boolmethodPush one item into this producer’s own ring.
push_buffer(self, data: bytes, item_len: int) -> intmethodPush items cut out of one buffer, item_len bytes each, and answer how many landed.
push_many(self, items: Sequence[bytes]) -> intmethodPush a run of items in one crossing, stopping at the first that will not fit, and answer how many landed. A count short of what was handed in leaves the rest with the caller.

Notifier

AttributeType
indexint
is_signaledbool
nativeint
MethodKindWhat it does
drain(self) -> NonemethodClear a pending signal, so the next wait blocks rather than returning at once on a signal already consumed.
`wait(self, timeout: floatNone=None) -> bool`method

NotifierSet

AttributeType
attachedint
MethodKindWhat it does
attach(self) -> NotifiermethodAttach a notifier of this process’s own.
signal(self) -> intmethodWake every attached notifier, and say how many were signaled.
__init__(self, path: str) -> NonemethodObtain the notifier set at path, creating it when the file does not exist. Attaching does not by itself give this process a notifier: call attach for that.

OrderedReceiver

AttributeType
correctionsint
strategystr
MethodKindWhat it does
drain(self, max_items: int=4096) -> list[tuple[bytes, int]]methodEverything the ring holds now plus everything the window still holds back, in order, in one crossing.
`flush(self) -> tuple[bytes, int]None`method
flush_all(self) -> list[tuple[bytes, int]]methodEverything still held, in order, for a caller that would rather end a stream in one crossing than a loop of them.
`recv(self) -> tuple[bytes, int]None`method

OwnerLease

AttributeType
max_value_bytesint
owner`int
termint
MethodKindWhat it does
`beat(self, pid: intNone=None) -> bool`method
flush(self) -> NonemethodPut the lease’s file on the disk, and wait for it.
flush_async(self) -> NonemethodAsk for the lease’s file to reach the disk without waiting.
`held_by_me(self, pid: intNone=None) -> bool`method
`hold(self, grace_epochs: int=0, pid: intNone=None) -> LeaseHold`method
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str) -> OwnerLeasestaticmethodAttach to a lease that already exists, failing if it does not.
`read(self, pid: intNone=None) -> bytesNone`
`release(self, pid: intNone=None) -> bool`method
`reset(path: str, value: bytesNone=None) -> OwnerLease`staticmethod
tick_epoch(self) -> intmethodStep the epoch every holder measures the grace period in, and give back the epoch this reached. Nothing steps it on its own.
`try_acquire(self, grace_epochs: int=0, pid: intNone=None) -> bool`method
`write(self, value: bytes, pid: intNone=None) -> bool`method
`init(self, path: str, value: bytesNone=None) -> None`method

PathChanges

AttributeType
last`tuple[int, int, int]
marked_sharefloat
route_movementfloat
MethodKindWhat it does
observe(self, ttl: int, congestion_mark: int, hops: int) -> NonemethodFeed one item: its remaining time to live, its congestion marking, and how many hops it took.
__init__(self) -> NonemethodA sensor with nothing observed yet. It lives in this process and has no file behind it. Feed it with observe, one call per item received, and it infers from the spread of what it is told.

Periodicity

AttributeType
period`tuple[float, float]
seconds_to_next`float
MethodKindWhat it does
observe(self, delay_microseconds: float, at_microseconds: int) -> NonemethodFeed one delay, in microseconds, and when it was taken, also in microseconds.
__init__(self) -> NonemethodA sensor with nothing observed yet. It lives in this process and has no file behind it. Each observe needs both the delay and when it was taken, because a beat can only be found against time.

PermitHold

AttributeType
heldbool
MethodKindWhat it does
release(self) -> NonemethodGive the permit back now rather than at the end of a block. Calling it twice is harmless, and held says whether it is still out. A genuine failure is raised rather than dropped: releasing more permits than the semaphore allows is a fault in the caller’s bookkeeping, not a condition to ignore.
__enter__(self) -> SelfmethodAnswer the same object, which is why a permit is normally taken as with semaphore.acquire() as permit: rather than bound to a name.
__exit__(self, *args: object) -> boolmethodGive the permit back. A refusal here is reported rather than dropped: releasing more permits than the semaphore allows is a real fault in the caller’s bookkeeping.

PubSub

AttributeType
capacityint
headint
payload_sizeint
MethodKindWhat it does
open(path: str, capacity: int) -> PubSubstaticmethodAttach to a pubsub ring that already exists, raising OSError when it does not. capacity must be the one it was created with. Attaching does not subscribe: take a subscription of your own.
publish(self, item: bytes) -> intmethodPublish one item and return the position it landed at. The publisher never waits for a subscriber.
`publish_many(self, items: Sequence[bytes]) -> intNone`method
`read_at(self, position: int) -> bytesNone`method
subscribe(self) -> SubscribermethodA subscriber starting where the publisher is now, so it sees what follows and nothing that came before.
subscribe_from(self, position: int) -> SubscribermethodA subscriber starting at position, for one replaying from a place it recorded earlier.
__init__(self, path: str, capacity: int) -> NonemethodObtain the pubsub ring at path holding capacity items, creating it when the file does not exist.

QosPolicy

AttributeType
durabilitystr
reliabilitystr
keep_last`int
max_latencyfloat
orderingstr
MethodKindWhat it does
persistent_log() -> QosPolicystaticmethodThe settings for a stream that must survive the process: a mapped file, senders waiting, everything kept.
reliable_pubsub() -> QosPolicystaticmethodThe settings for a stream nothing may fall out of: named memory other processes can reach, senders waiting rather than dropping.
snapshot(self) -> QosSnapshotmethodEvery setting read together, so a decision is made against one consistent set rather than five separate reads.
streaming() -> QosPolicystaticmethodThe settings for a stream that would rather lose an item than hold its sender up: in-process memory, dropping when full, the last thousand or so items, a tenth of a second.
`init(self, durability: str=‘volatile’, reliability: str=‘best_effort’, keep_last: intNone=1024, max_latency: float=0.1) -> None`method

QosSnapshot

AttributeType
durabilitystr
keep_last`int
max_latencyfloat
orderingstr
reliabilitystr
MethodKindWhat it does
`recommends_locale_change(self, current: str) -> strNone`method
`recommends_ordering_change(self, current: str) -> strNone`method

QuicBridgeClient

MethodKindWhat it does
run(self, items: int) -> NonemethodConnect and ship items of them, waiting until the reading end has acknowledged the last of them.
`init(self, ring: Ring, server: tuple[str, int], cert: bytes, server_name: str, local: tuple[str, int]None=None) -> None`method

QuicBridgeServer

AttributeType
local_addrtuple[str, int]
MethodKindWhat it does
accept_one(self) -> intmethodTake one connection, read it to its end, and answer how many items arrived.
__init__(self, ring: Ring, local: tuple[str, int], cert: bytes, key: bytes) -> Nonemethodring is the ring to put arriving items into, local the address to listen on, and cert and key the certificate and key this end proves itself with. A port of zero lets the system pick one, which local_addr then reports.

RWLock

AttributeType
readersint
MethodKindWhat it does
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str) -> RWLockstaticmethodAttach to a lock that already exists, raising OSError when it does not. Attaching takes nothing: read and write are what take a hold.
read(self) -> HoldmethodTake the read hold, waiting for it. Other readers may hold it at the same time; a writer may not.
`read_for(self, timeout: float) -> HoldNone`method
`try_read(self) -> HoldNone`method
`try_write(self) -> HoldNone`method
write(self) -> HoldmethodTake the write hold, waiting for it. Nobody else holds it while this does.
`write_for(self, timeout: float) -> HoldNone`method
__init__(self, path: str) -> NonemethodObtain the lock at path, creating it when the file does not exist and attaching to the live one when it does.

RateLimiter

AttributeType
availableint
capacityint
refill_per_secondint
MethodKindWhat it does
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the token count without this.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int, refill_per_second: int) -> RateLimiterstaticmethodAttach to a limiter that already exists, raising OSError when it does not. The capacity and refill rate must be the ones it was created with.
reset(self) -> NonemethodRefill the bucket to full, which is the opposite of what the name suggests: this hands out capacity rather than taking it away.
try_acquire(self, n: int=1) -> boolmethodTake n tokens if they are there. False means they were not, which is an answer rather than a failure.
__init__(self, path: str, capacity: int, refill_per_second: int) -> NonemethodObtain the limiter at path holding at most capacity tokens and refilling refill_per_second of them each second, creating it when the file does not exist. Either being zero is a ValueError.

Region

AttributeType
capacityint
MethodKindWhat it does
allocate(self, value: bytes) -> intmethodTake a slot and write value into it, returning its index.
get(self, index: int) -> bytesmethodRead slot index.
set(self, index: int, value: bytes) -> NonemethodWrite slot index.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. A memoryview taken inside the block stays valid after it, because the view holds its own reference to the region and the export count is what keeps the mapping under it.
__init__(self, path: str, capacity: int, slot_size: int, alignment: int=1, tag: int=0) -> NonemethodObtain the region at path holding capacity slots of slot_size bytes, creating it when the file does not exist.
__len__(self) -> intmethodHow many slots are allocated.

Registration

AttributeType
closedbool
idint
last_policy_error`BaseException
policy_errorsint
tagint
MethodKindWhat it does
close(self) -> NonemethodUnregister from the sidecar. Returns once no scan is inside the registration, and the object can then be observed again.
stats(self) -> InstanceStatsmethodWhat the sidecar has drained from the object so far, as its last scan left it. subetha.sidecar.scan_now() drains what is waiting.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodUnregister as the block ends, and let an exception through.

ReorderWindow

AttributeType
correctionsint
windowint
MethodKindWhat it does
`flush(self) -> tuple[bytes, int]None`method
flush_all(self) -> list[tuple[bytes, int]]methodEverything still held, in order, in one crossing.
push(self, stamp: int, payload: bytes) -> NonemethodHold one item, with the stamp its sender gave it.
push_many(self, items: Sequence[tuple[int, bytes]]) -> intmethodHold a run of items in one crossing. Each is a stamp and a payload.
`take(self) -> tuple[bytes, int]None`method
widen_to(self, window: int) -> NonemethodWiden the window to at least this, which is what a caller does when the number of senders grows.
__init__(self, floor: int=8, cap: int=1024) -> Nonemethodfloor is the window to start at and cap the widest it may grow to. A window at least as wide as the number of senders puts every item back in order.
__len__(self) -> intmethodHow many items are held in the window waiting for a gap to fill, which is not how many have passed through. An empty window is falsy, so if not window: works.

Reservoir

AttributeType
max_value_bytesint
capacityint
total_seenint
MethodKindWhat it does
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the sample without this; flushing is about surviving a machine that stops.
flush_async(self) -> NonemethodStart writing the mapping back to its file and return at once, without waiting for the write to land. Nothing is on disk yet when this returns; flush is the form that waits.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int) -> ReservoirstaticmethodAttach to a reservoir another holder made. The capacity must be the one it was made with.
`record(self, value: bytes) -> intNone`method
record_many(self, values: Sequence[bytes]) -> intmethodOffer a run of values in one crossing, and say how many were kept.
reset(self) -> NonemethodEmpty the sample and forget how much has been offered.
snapshot(self) -> list[bytes]methodEverything in the sample right now, in one crossing.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The sample and the total seen both stay for whoever attaches next; reset is what empties a reservoir.
__init__(self, path: str, capacity: int) -> NonemethodHow many items the sample holds at most. A value is at most 52 bytes.
__len__(self) -> intmethodHow many values are held in the sample right now, which is at most the capacity and is not total_seen. A reservoir that has been offered a million values still holds only its capacity.

Ring

AttributeType
approx_lenint
capacityint
managedbool
max_consumersint
max_producersint
morph_refusalsint
shapestr
sidecar_morphsint
stampedbool
stamps`str
total_capacityint
MethodKindWhat it does
is_empty(self) -> boolmethodWhether the ring looked empty when asked. Read without stopping the producers, so it is a sighting rather than a promise: a send can land the instant after. Treat None from recv as the real answer, and use this for reporting.
`observe(self, policy: _PolicyNone=None) -> Registration`method
`open(path: str, capacity: int, max_producers: int=1, max_consumers: int=1, stamps: strNone=None, managed: bool=False, scan_interval_us: intNone=None) -> Ring`
ordered_receiver(self, consumer: int) -> OrderedReceivermethodA reader that hands items back in the order their senders made them, rather than the order they happened to arrive in.
`recv(self, consumer: int) -> bytesNone`method
`recv_frame(self, consumer: int) -> bytesNone`method
recv_many(self, consumer: int, max_items: int) -> list[bytes]methodReceive up to max_items in one crossing.
register_consumer(self) -> intmethodTake a consumer id. Every receive names one.
register_producer(self) -> intmethodTake a producer id. Every send names one.
send(self, producer: int, item: bytes) -> boolmethodSend one item as producer. False means the ring was full.
send_buffer(self, producer: int, data: bytes, item_len: int) -> intmethodSend items packed end to end, item_len bytes each, with no Python object built per item.
send_frame(self, producer: int, payload: bytes) -> boolmethodSend a payload longer than a slot, carried in frames.
send_many(self, producer: int, items: Sequence[bytes]) -> intmethodSend a run of items, stopping at the first refusal.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Producer and consumer positions taken inside the block are still registered after it, and items still in the ring stay there for whoever attaches next.
`init(self, path: str, capacity: int, max_producers: int=1, max_consumers: int=1, stamps: strNone=None, managed: bool=False, scan_interval_us: intNone=None) -> None`

RoundTripShape

AttributeType
samplesint
two_groups`float
wireless_confidencefloat
MethodKindWhat it does
observe(self, microseconds: float) -> NonemethodFeed one round trip, in microseconds.
observe_many(self, microseconds: Sequence[float]) -> intmethodFeed a run of round trips in one crossing.
__init__(self) -> NonemethodA shape with nothing observed yet. It lives in this process and has no file behind it. Feed it with observe or observe_many, in microseconds, before asking it anything.

Semaphore

AttributeType
availableint
max_permitsint
waitersint
MethodKindWhat it does
acquire(self) -> PermitHoldmethodTake a permit, waiting for one. The interpreter is detached while this waits.
`acquire_for(self, timeout: float) -> PermitHoldNone`method
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, max_permits: int) -> SemaphorestaticmethodAttach to a semaphore that already exists, raising OSError when it does not. max_permits must be the one it was created with, and attaching takes no permit: acquire is what does.
`try_acquire(self) -> PermitHoldNone`method
`init(self, path: str, initial: int, max_permits: intNone=None) -> None`method

SensReceiver

AttributeType
alivebool
codestr
local_addrtuple[str, int]
max_item_sizeint
send_failuresint
switchesint
MethodKindWhat it does
poll(self) -> list[bytes]methodDrive the link and answer every item it could rebuild this time round. An empty answer is ordinary and means nothing was ready, not that anything is wrong.
poll_from(self) -> list[tuple[int, bytes]]methodAs poll, and also which sender each item came from, for a reader taking from several at once.
`init(self, local: tuple[str, int], max_item_size: int, code: strNone=None) -> None`method

SensSender

AttributeType
codestr
datagramstuple[int, int]
local_addrtuple[str, int]
loss`float
max_item_sizeint
switchesint
MethodKindWhat it does
send(self, item: bytes) -> NonemethodSend one item, of anything up to max_item_size bytes.
send_many(self, items: Sequence[bytes]) -> intmethodSend a run of items in one crossing, which is the shape to reach for: the work per item is small enough that the crossing would otherwise be most of the cost.
`init(self, local: tuple[str, int], peer: tuple[str, int], max_item_size: int, code: strNone=None) -> None`method

SharedArc

AttributeType
holdersint
value_sizeint
MethodKindWhat it does
get(self) -> bytesmethodThe whole value.
open(path: str, value_bytes: int, max_holders: int=16, keep_on_last: bool=False) -> SharedArcstaticmethodAttach to a shared value that already exists, counting this process as another holder.
read_at(self, offset: int, length: int) -> bytesmethodPart of the value, for a caller that wants a field rather than the whole of a large record.
write_at(self, offset: int, value: bytes) -> NonemethodOverwrite as many bytes as value carries, starting at offset. A range running past the end of the value is an OSError rather than a short write. The partner of read_at, for a caller that wants one field of a large record rather than all of it.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without letting go of the value, and let an exception through.
__init__(self, path: str, value: bytes, max_holders: int=16, keep_on_last: bool=False) -> NonemethodThe constructor asserts at least one holder, so that is refused here first. keep_on_last decides what happens when the last holder lets go: keep the file for a later process, or unlink it.

Slab

AttributeType
capacityint
element_sizeint
writablebool
MethodKindWhat it does
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the slots without this; flushing is about surviving a machine that stops.
get(self, index: int) -> bytesmethodRead slot index, always element_size bytes. A slot nobody has written to reads as zeros rather than raising, because it exists from creation. An index past the capacity is an OSError.
open(path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> SlabstaticmethodAttach to a slab that already exists, able to write into it. Raises OSError when it is not there, and every layout argument must be the one it was created with.
open_read_only(path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> SlabstaticmethodAttach for reading only, so set and write_range raise and writable answers False.
read_range(self, start: int, count: int) -> bytesmethodRead count slots from start packed end to end, one crossing and one object, each slot read through its seqlock.
set(self, index: int, value: bytes) -> NonemethodOverwrite slot index, stepping its seqlock either side so a reader racing this retries instead of seeing half of each value. An index past the capacity, or a slab opened read-only, is an OSError.
slot_version(self, index: int) -> intmethodThe seqlock counter for a slot. Even means nobody is writing; the same value twice means no write landed between the two reads.
write_range(self, start: int, data: bytes) -> intmethodWrite slots packed end to end from start.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Nothing is flushed on the way out: call flush where durability matters.
__init__(self, path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> NonemethodObtain the slab at path with capacity slots of element_size bytes, creating it when the file does not exist.
__len__(self) -> intmethodHow many slots the slab has, which is its capacity and not a count of the ones written to. Every slot exists from creation, so this never changes and a slab is never empty.

SlabPin

AttributeType
epochint
heldbool
MethodKindWhat it does
`get(self, slot: int) -> bytesNone`method
`get_many(self, slots: Sequence[int]) -> list[bytesNone]`method
release(self) -> NonemethodGive the pin back now rather than at the end of a block. Calling it twice is harmless.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodGive the pin back, and let an exception through.

SpscRing

AttributeType
capacityint
payload_sizeint
MethodKindWhat it does
open(path: str, capacity: int) -> SpscRingstaticmethodAttach to a ring that already exists. capacity must be the one it was created with.
`pop(self) -> bytesNone`method
pop_buffer(self, max_items: int) -> tuple[bytes, int]methodPop up to max_items into one buffer packed end to end, and return how many were written. Nothing is allocated per item.
pop_many(self, max_items: int) -> list[bytes]methodPop up to max_items, stopping when the ring runs empty.
push(self, item: bytes) -> boolmethodPush one item. False means the ring was full, which is an answer rather than a failure; anything else raises.
push_buffer(self, data: bytes, item_len: int) -> intmethodPush items packed end to end in one buffer, item_len bytes each. No Python object is built per item, which is what makes this the fastest way in: the caller’s own array goes straight across.
push_many(self, items: Sequence[bytes]) -> intmethodPush a run of items, stopping at the first the ring refuses. Returns how many went in, so the caller keeps the rest.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Nothing is drained: items still in the ring stay there for whoever attaches next, because the file outlives the process.
__init__(self, path: str, capacity: int) -> NonemethodObtain the ring at path holding capacity slots, creating it when the file does not exist.

Stack

AttributeType
approx_lenint
capacityint
element_sizeint
MethodKindWhat it does
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the writes without this; flushing is about surviving a machine that stops.
is_empty(self) -> boolmethodWhether the stack looked empty when asked. Read without stopping anyone else, so it is a sighting rather than a promise: a push can land the instant after. Treat None from pop as the real answer, and use this for reporting.
open(path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> StackstaticmethodAttach to a stack that already exists, raising OSError when it does not. Every layout argument describes the file rather than asking anything of it, so all of them must match what it was created with.
`peek(self) -> bytesNone`method
`pop(self) -> bytesNone`method
pop_many(self, max_items: int) -> list[bytes]methodTake up to max_items in one crossing, stopping early when the stack runs dry, newest first. An empty list means there was nothing, and nothing here raises.
push(self, item: bytes) -> boolmethodPush one item. False means the stack is full.
push_many(self, items: Sequence[bytes]) -> intmethodPush a run of items in one crossing, stopping at the first that will not fit, and answer how many landed. They go on in the order given, so the last one handed in is the first one pop returns.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Nothing is popped on the way out: whatever is on the stack stays there for whoever attaches next.
__init__(self, path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> NonemethodObtain the stack at path with room for capacity items of element_size bytes, creating it when the file does not exist. A capacity below one is a ValueError, and so is a layout the alignment cannot satisfy.

Subscriber

AttributeType
positionint
MethodKindWhat it does
lag(self) -> intmethodHow far behind the publisher this subscriber is.
`next(self) -> bytesNone`method
next_many(self, max_items: int) -> list[bytes]methodUp to max_items in one crossing, stopping when caught up.

TcpBridgeClient

MethodKindWhat it does
run(self, items: int) -> NonemethodConnect and ship items of them, waiting until all have gone.
__init__(self, ring: Ring, server: tuple[str, int]) -> Nonemethodring is the ring to take items from and server the address of the reading end, as host and port.

TcpBridgeServer

AttributeType
local_addrtuple[str, int]
MethodKindWhat it does
accept_one(self) -> intmethodTake one connection, read it to its end, and answer how many items arrived. Waits until the sending end has finished.
__init__(self, ring: Ring, local: tuple[str, int]) -> Nonemethodring is the ring to put arriving items into and local the address to listen on, as host and port. A port of zero lets the system pick one, which local_addr then reports.

TimePointTile

AttributeType
lanesint
max_value_bytesint
fullbool
MethodKindWhat it does
`at(self, lane: int) -> tuple[int, bytes]None`method
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the points without this; flushing is about surviving a machine that stops.
flush_async(self) -> NonemethodStart writing the mapping back to its file and return at once, without waiting for the write to land. Nothing is on disk yet when this returns; flush is the form that waits.
insert(self, version: int, value: bytes) -> intmethodWrite a value at a version and get the place it went in. Raises when all sixteen places are taken.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str) -> TimePointTilestaticmethodAttach to a tile another holder made.
remove(self, lane: int) -> NonemethodEmpty one place.
reset(path: str) -> TimePointTilestaticmethodEmpty the tile and remake it.
visible(self, version: int) -> list[tuple[int, bytes]]methodEverything a reader at version can see, as version and value pairs, in one crossing.
visible_count(self, version: int) -> intmethodHow many places a reader at version can see.
visible_mask(self, version: int) -> intmethodWhich places a reader at version can see, as sixteen bits with the lowest standing for the first place.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The points recorded stay for whoever attaches next, and nothing is flushed on the way out.
__init__(self, path: str) -> NonemethodA value is at most 52 bytes.
__len__(self) -> intmethodHow many of the tile’s sixteen places are taken. full is the same question asked the other way round.

Timing

AttributeType
clock_skewfloat
jitterfloat
samplesint
spacingfloat
trendfloat
trend_debiasedfloat
MethodKindWhat it does
observe(self, sent: int, received: int) -> NonemethodFeed one item’s send and receive times, in microseconds. The two clocks need not agree with each other.
__init__(self, window: int=64) -> Nonemethodwindow is how many recent items the answers are taken over.

TinyBloom

AttributeType
suggested_capacityint
bitsint
set_bitsint
MethodKindWhat it does
contains(self, key: bytes) -> boolmethodAs in, spelled out.
contains_many(self, keys: Sequence[bytes]) -> list[bool]methodAsk about a run of keys in one crossing.
false_positive_rate(keys: int) -> floatstaticmethodThe share of wrong yeses to expect once keys keys are in, between zero and one.
from_bits(bits: int) -> TinyBloomstaticmethodRebuild a filter from the number bits gave.
insert(self, key: bytes) -> NonemethodAdd a key, setting its bits. Nothing can be taken out again, and the filter is only sixty-four bits wide, so the rate of wrong yeses climbs quickly past suggested_capacity keys.
insert_many(self, keys: Sequence[bytes]) -> intmethodAdd a run of keys in one crossing.
__contains__(self, key: bytes) -> boolmethodFalse means the key is definitely absent, True means it is probably present. A filter this small says yes wrongly often once it holds more than a handful of keys.
`init(self, keys: Sequence[bytes]None=None) -> None`method

TopologyMap

AttributeType
broadcast_rootint
busiest_receivertuple[int, int]
busiest_sendertuple[int, int]
participantsint
recommendation_epochint
total_sendsint
MethodKindWhat it does
fan_in(self, receiver: int) -> intmethodHow many different places send to this one.
fan_out(self, sender: int) -> intmethodHow many different places this one sends to.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, participants: int) -> TopologyMapstaticmethodAttach to a topology another holder made, with the participant count it was made with.
publish_recommendation(self) -> strmethodWork the shape out and write it down, so every process reads the same one. Answers what was published.
published_recommendation(self) -> strmethodThe shape last published, which may not be what the counts suggest now.
recommend(self) -> strmethodThe shape the counts suggest, one of point_to_point, broadcast_tree or all_to_all_mesh. Reading this does not publish it.
record_many(self, sends: Sequence[tuple[int, int]]) -> intmethodRecord a run of sends in one crossing.
record_send(self, sender: int, receiver: int) -> intmethodRecord one send, answering how many have gone that way.
`reset(path: str, participants: int, fan_out_threshold: intNone=None, fan_in_threshold: intNone=None) -> TopologyMap`
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Everything recorded stays counted for whoever attaches next; reset is what clears the observations.
`init(self, path: str, participants: int, fan_out_threshold: intNone=None, fan_in_threshold: intNone=None) -> None`

Tower

AttributeType
depthint
value_sizeint
MethodKindWhat it does
append(self, value: bytes) -> list[int]methodStore a value, taking a fresh place on the top level, and answer the path that reaches it.
append_many(self, values: Sequence[bytes]) -> list[list[int]]methodStore a run of values in one crossing, answering a path for each.
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Every level goes, not just the bottom one. Another process mapping the same file sees the values without this; flushing is about surviving a machine that stops.
get(self, path: Sequence[int]) -> bytesmethodThe value a path reaches. Raises when a level along the way no longer agrees with the path, naming that level.
get_many(self, paths: Sequence[Sequence[int]]) -> list[bytes]methodSeveral paths in one crossing.
insert_at_top(self, top: int, value: bytes) -> list[int]methodStore a value under a named place on the top level, rather than a fresh one, and answer the path that reaches it.
open(path: str, capacity: int, value_size: int, levels: Sequence[tuple[str, int]]) -> TowerstaticmethodAttach to a tower another holder made, with the shape it was made with.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Paths handed out inside the block still resolve to the same values after it, and nothing is flushed on the way out.
__init__(self, path: str, capacity: int, value_size: int, levels: Sequence[tuple[str, int]]) -> Nonemethodpath is where the bottom level lives and levels the levels above it, each a file and how many places it holds. The levels are given top first. With no levels the tower is one deep, which is a plain region reached by a path of one number.
__len__(self) -> intmethodValues stored at the bottom level.

Universal

AttributeType
generationint
migrationsint
op_countstuple[int, int]
strategystr
MethodKindWhat it does
clear(self) -> NonemethodTake every value out, so the set is empty again. The storage strategy it had migrated to is kept rather than reset, and migrations goes on counting from where it was.
contains(self, value: int) -> boolmethodAs in, spelled out.
contains_many(self, values: Sequence[int]) -> list[bool]methodAsk about a run of values in one crossing.
insert(self, value: int) -> NonemethodAdd a value to the set.
insert_many(self, values: Sequence[int]) -> intmethodAdd a run of values in one crossing.
migrate_to(self, strategy: str) -> NonemethodMove the set to list or map by hand. Moving it to where it already is does nothing.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int) -> UniversalstaticmethodAttach to a set another holder made, with the capacity it was made with.
reset(path: str, capacity: int) -> UniversalstaticmethodEmpty the set and remake it at this capacity.
snapshot(self) -> list[int]methodEverything in the set, in one crossing.
__contains__(self, value: int) -> boolmethodWhether the value is in the set. Exact, not probabilistic: unlike the filters, a True here is a fact.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Whatever strategy it migrated to inside the block is the one it is still in afterward.
__init__(self, path: str, capacity: int) -> NonemethodThe set lives in files beside the path given, one per way of storing it.
__len__(self) -> intmethodHow many values the set holds. It can raise, unlike most len implementations, because reading the count means reading the mapping and that can fail.

Vec

AttributeType
capacityint
element_sizeint
writablebool
MethodKindWhat it does
clear(self) -> NonemethodDrop every element, so the length goes to zero. The capacity and the file are left as they are, and the space is reused by the next push.
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees a write without this; flushing is about surviving a machine that stops.
`get(self, index: int) -> bytesNone`method
open(path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> VecstaticmethodAttach to a vec that already exists, raising OSError when it does not. Every layout argument must be the one it was created with: they describe the file rather than asking anything of it.
`pop(self) -> bytesNone`method
`push(self, value: bytes) -> intNone`method
push_many(self, values: Sequence[bytes]) -> intmethodAppend a run of elements, stopping at the first refusal, and say how many landed.
read_range(self, start: int, count: int) -> bytesmethodRead count elements from start, packed end to end into one object. Stops at the end of what is live.
set(self, index: int, value: bytes) -> NonemethodOverwrite element index, stepping its seqlock either side so a reader racing this retries instead of seeing half of each value.
write_range(self, start: int, data: bytes) -> intmethodWrite elements packed end to end into consecutive slots from start, and return how many landed.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. The elements stay where they are for whoever attaches next; clear is what empties a vec, not the end of a block.
__init__(self, path: str, capacity: int, element_size: int, alignment: int=1, tag: int=0) -> NonemethodObtain the vec at path holding up to capacity elements of element_size bytes, creating it when the file does not exist and attaching to what is already there when it does.
__len__(self) -> intmethodHow many elements are live, which is not the capacity. A vec that has never been pushed to is empty, and bool(vec) is False.

VersionChain

AttributeType
max_value_bytesint
capacityint
current`tuple[int, bytes]
MethodKindWhat it does
clear(self) -> NonemethodThrow away every version.
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the versions without this; flushing is about surviving a machine that stops.
flush_async(self) -> NonemethodStart writing the mapping back to its file and return at once, without waiting for the write to land. Nothing is on disk yet when this returns; flush is the form that waits.
`observe(self, policy: _PolicyNone=None) -> Registration`method
open(path: str, capacity: int) -> VersionChainstaticmethodAttach to a chain another holder made, with the capacity it was made with.
push(self, version: int, value: bytes) -> NonemethodWrite a new version. The version number must be above the one already at the front, which is what keeps the history in order.
`read_at(self, version: int) -> bytesNone`method
reset(path: str, capacity: int) -> VersionChainstaticmethodEmpty the chain and remake it at this capacity.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Every version written inside the block is still in the chain after it, and nothing is reclaimed on the way out.
__init__(self, path: str, capacity: int) -> NonemethodHow many versions the chain holds at most. A value is at most 44 bytes.
__len__(self) -> intmethodHow many versions the chain holds, which is not the capacity and not the number of distinct values: one value written three times is three versions until something reclaims the older two.

VersionedMap

AttributeType
capacityint
MethodKindWhat it does
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the entries without this; flushing is about surviving a machine that stops.
`get(self, key: int) -> intNone`method
`get_many(self, keys: Sequence[int]) -> list[intNone]`method
`insert(self, key: int, value: int) -> intNone`method
`insert_at(self, key: int, value: int, born: int) -> intNone`method
`insert_many(self, entries: Sequence[tuple[int, int]]) -> list[intNone]`method
open(path: str, capacity: int, epochs_path: str, max_pins: int=16) -> VersionedMapstaticmethodAttach to a map another holder made, with the capacity and pin count it was made with.
pin(self) -> MapPinmethodTake a pin, fixing one epoch to scan the whole map at. Give it back when the block ends, or reclaiming cannot move past it.
`remove(self, key: int) -> intNone`method
`remove_at(self, key: int, died: int) -> intNone`method
sweep(self) -> intmethodTake away every entry marked as no longer current that nothing can still reach, answering how many went.
void_epoch(self, epoch: int) -> intmethodUndo every write stamped at exactly this epoch, answering how many entries were touched.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. A MapPin taken inside the block is not given back here: pin it in its own with block, or the epoch it holds keeps sweep from reclaiming anything.
__init__(self, path: str, capacity: int, epochs_path: str, max_pins: int=16) -> Nonemethodcapacity is how many entries the map holds, counting the ones marked as no longer current, and max_pins how many readers may scan at once. The epochs live in their own file, which other structures may share.
__len__(self) -> intmethodEntries the map holds, counting the ones marked as no longer current.

VersionedSlab

AttributeType
depthint
max_value_bytesint
capacityint
MethodKindWhat it does
flush(self) -> NonemethodAsk the operating system to write the mapping back to its file, and wait for it. Another process mapping the same file sees the versions without this; flushing is about surviving a machine that stops.
`get(self, slot: int) -> bytesNone`method
`history(self, slot: int) -> list[tuple[bytes, int, intNone]]`method
open(path: str, capacity: int, epochs_path: str, max_pins: int=16) -> VersionedSlabstaticmethodAttach to a slab another holder made, with the capacity and pin count it was made with.
pin(self) -> SlabPinmethodTake a pin, fixing one epoch to read the whole slab at. Give it back when the block ends, or reclaiming cannot move past it.
`retire(self, slot: int) -> bytesNone`method
`retire_at(self, slot: int, died: int) -> bytesNone`method
set(self, slot: int, value: bytes) -> NonemethodWrite a slot, at the next epoch.
set_at(self, slot: int, value: bytes, born: int) -> NonemethodWrite a slot at a named epoch, for a caller stepping the epochs itself.
sweep_slot(self, slot: int) -> intmethodThrow away every version of a slot that nothing can still see, answering how many went.
void_epoch(self, epoch: int) -> intmethodUndo every write stamped at exactly this epoch, across the whole slab, answering how many versions were touched.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. A SlabPin taken inside the block is not given back here: pin it in its own with block, or the epoch it holds keeps reclamation waiting.
__init__(self, path: str, capacity: int, epochs_path: str, max_pins: int=16) -> Nonemethodcapacity is how many slots there are and max_pins how many readers may hold a pin at once. The epochs live in their own file, which other structures may share. A value is at most 44 bytes.

WorkQueue

AttributeType
max_item_sizeint
MethodKindWhat it does
`pop(self) -> bytesNone`method
push(self, item: bytes) -> boolmethodAdd work, at the owner’s end.
push_many(self, items: Sequence[bytes]) -> intmethodAdd a run of work in one crossing, answering how many went in.
`steal(self) -> bytesNone`method
steal_from(path: str) -> WorkQueuestaticmethodAttach to a queue somebody else owns, to take from it.
steal_many(self, max_items: int=64) -> list[bytes]methodTake up to max_items by stealing, in one crossing.
__enter__(self) -> SelfmethodAnswer the same object, so a with block can give it a name.
__exit__(self, *args: object) -> boolmethodLeave the block without closing anything, and let an exception through. Work still in the queue stays there for whoever attaches next, and no worker is told to stop.
__init__(self, path: str, capacity: int=1024, thieves: int=1) -> Nonemethodcapacity is the number of items in flight, rounded up to a power of two. thieves is how many are expected to take from it, which is what the shape underneath is picked from.