Spec index | Previous: Node.js | Next: Go
Python binding Core public contract¶
What this chapter defines — the public type, ownership, and error contract the
zlinkPython package provides on top of Core raw messaging.
- This document defines the Core raw messaging contract the
zlinkPython package provides. - A feature not defined by this document and the public header is not part of the Python binding contract.
- Python 3.9-and-later support and the package version follow the distribution metadata.
- The current native package target is Linux x86_64; other targets are outside this contract's supported scope until a separate candidate payload and clean-consumer verification exist.
| Section | Covers |
|---|---|
| Scope | The list of public Python types that represent Core resources |
| Package surface | The boundary between the public factories and the private area |
| Byte HWM and Auto-HWM | Mapping between Python int and Core uint64_t byte HWM |
| Ownership and lifetime | Ownership/release rules for native handles, messages, and Received |
| Send/receive and no-data | How submit and no-data are represented, and how native failures are delivered |
| Receive flow state | The receive-flow state type, setter, and monitor surface |
| Error | The ZlinkError family and its result fields |
| Python version and the type package | The supported Python version and the type-check target |
| Related documents | Links to the guide, the Core spec, and internals |
Scope¶
The Python binding represents the following Core resources as Python objects and Protocols.
| Area | Public concepts |
|---|---|
| Core | Context, ContextOptions, Message, Received, RoutingId |
| Socket | PAIR, DEALER, ROUTER, STREAM, PUB, SUB, XPUB, XSUB |
| Eventing | MonitorSocket, MonitorEvent, MonitorStatus, Poller, PollEvents, Timer |
| Utility | AtomicCounter, Stopwatch, Thread, proxy, sleep |
| Result | SubmitResult, RequestResult, RecvResult, ConfigResult, and their matching errors |
A socket preserves Core's raw endpoint and message-routing semantics as-is. The Python binding does not expose Core's internal handles, FFI symbols, or native structs as public types.
Package surface¶
Users work with the factories and contract types at the zlink package
root. They do not import implementation modules directly, and _native
and _runtime are private areas. The package root has no separate domain
type or compatibility alias outside the Core raw contract.
The main factories are create_context(), create_pair_socket(),
create_dealer_socket(), create_router_socket(),
create_stream_socket(), create_pub_socket(), create_sub_socket(),
create_poller(), create_timer(), create_received(), and the
create_message family. The exact Python signatures are governed jointly
by the contract module in the same directory and the public header.
Byte HWM and Auto-HWM¶
Core owns HWM (the queue byte threshold) calculation and queue admission. Python
send_high_water_mark and receive_high_water_mark are byte-valued int
properties and accept only the non-negative uint64_t range. The binding
passes each value to Core as an exact 8-byte option, and a getter returns
Core's 64-bit value as a Python int. 0 means unlimited.
The context passes byte-valued core_hwm_memory_limit_bytes and
core_hwm_budget_bytes, plus core_hwm_profile, unchanged to Core. Core applies
the profile ratio and distributes the result across physical directional
queues exactly once. Setting a directional HWM makes that direction a manual
override and excludes it from automatic HWM recalculation.
Context also provides core_hwm_budget_snapshot() and
reset_core_hwm_budget_metrics(). Input precedence is manual Core budget,
explicit memory limit, a runtime hint only when a distinct VM hard limit is
unambiguous, then Core fallback. Setting either of the first two values disables
automatic runtime-hint detection. The binding does not combine the hint with
Core's hard limit. If an explicit input exceeds a finite hard limit Core
detected, the binding preserves the existing configuration error corresponding
to EINVAL and does not clamp the value.
Core decides backpressure when the accounted bytes retained by a pipe reach
the applied HWM. The Python binding does not recount messages and passes the
native result through the existing submit and error contract.
monitor_open(events=..., monitor_hwm_bytes=...) accepts a non-negative
uint64_t-range Python int. Zero selects the Core monitor default; a positive
value is forwarded unchanged, with no message-count alias or conversion. Planned,
applied, and deferred HWM and in-flight usage in MonitorStatus are byte-valued
Python int fields. Pending-message counts remain display diagnostics; no
slot, message-unit, size-cap, or connection-bucket property is exposed.
snd_pending_bytes and rcv_pending_bytes are separate byte values.
The Core budget snapshot projects ABI version/size, configured/runtime/resolved
memory limits, configured/effective budgets, planned/applied/manual-reserved
HWM, Core-queue/application/current/peak/provisional accounted bytes,
completion current/peak/pending and total-messaging values, monitor/instance
aggregates, application/completion queue counts,
outstanding_application_lease_count, retired_queue_count,
deferred_origin_credit_bytes, oversize/blocked/aggregate flags,
budget_generation, and measurement_epoch as Python int/boolean values.
application_accounted_bytes and those three owner-lifecycle fields are
ABI-reserved and always zero. Reset preserves current, pending, and queue-count
gauges, rebases both peaks to current, clears epoch counters, and increments
measurement_epoch. An ABI version/size mismatch is an unsupported error.
Ownership and lifetime¶
Contextowns the native context and releases it viaclose()or context-manager exit.- A socket, monitor, poller, or timer owns the native handle it creates. A call that uses the handle after a successful
close()is not allowed. Message.from_(value)creates a native message independent of the caller's value. After a successful send submit, native ownership of the message part moves to the Core send path.Receivedis a receive storage the caller creates. On a successfulrecv_into(received), the parts and routing metadata are recorded intoReceived; native parts are released onclose()or context-manager exit.- The native view
Received'spartsprovides is valid only while its owner stays open. If it must outlive that, copy the value withto_bytes()orto_bytes_list().
Received and TopicMessage preserve native parts, routing ID, reply token, topic, and
multipart framing, released through close(), context-manager exit, or storage reuse.
The common receive ownership contract defines the boundary with receive accounting.
Send/receive and no-data¶
- Send/request builder completion boundaries follow the async completion surface policy; completion joins and cancellation follow the async execution model.
- Reply and publish end with synchronous
submit(). Only a separatePublishOpprovides publish flags. - The runtime submits every part collected by a builder to the Core whole-message API once as a native array and count. Receive uses a reusable native array, capacity, and count; when the array is too small, it grows to the required size and retries the same unconsumed record.
- Caller-provided receive with
RecvFlags.DONT_WAITreturnsFalsewhen no message is available. - Direct-return control APIs such as timer and monitor return
Nonewhen no value is pending. - An actual native failure is delivered through its corresponding error type and is not hidden as no-data.
DEALER and ROUTER request/reply preserve Core routing metadata and ReplyToken.
Received.routing_id from ROUTER receive is a routing ID and is not converted into another identity
type. The single-part accessor is named single_part_or_throw(), matching the implementation and
contract tests.
Receive flow state¶
ReceiveFlowState is an IntEnum with RUNNING = 0 and PAUSED = 1.
Socket.set_receive_flow_state(state) returns None and raises ConfigError carrying
the failed native ConfigResult and errno.
State, result, and monitor projection follow the common receive-flow contract.
Error¶
A call that returns a Core result gives its matching Python error result,
code, and native_errno. An input-format error can be checked before the
call, but a native operation failure is never turned into a plain
ValueError. SubmitError, RequestError, RecvError, BindError,
ConnectError, ConfigError, CloseError, and HandlerError are all in
the ZlinkError family.
Python version and the type package¶
Public annotations use a form the Python 3.9 parser and runtime can
resolve. The package root includes py.typed. Public contract type
checking targets the Python 3.9 target pyrightconfig.json specifies, and
src/zlink/contracts.
Related documents¶
- Usage follows the Python guide.
- The reference for Core functions and layout is the repository's
core/include/zlink.hand the Core spec. - Implementation detail and callback/native lifetime explanations belong to the internals documents, not this one.
Pull completion public contract¶
Python package information follows its distribution metadata; the Core ABI version follows Core release metadata.
Python provides blocking submit_sync() and submit() returning a result object (SendSubmission/RequestSubmission: result and admitted, plus reply for a request).
The result accessors are read-only properties. Read the initial result through submission.result,
wait for admission with await submission.admitted, and wait for a request's response with await submission.reply.
Caller wait cancellation is expressed through awaitable cancellation.
Native completion IDs, user_context, and raw drain are not public APIs.
Submission results follow the common result projection;
completion joins, lifetime, and progress conditions for PollEventFlag.POLLCOMPLETION follow the
async execution model.
Poller.add_monitor(monitor: MonitorSocket, events: PollEventFlag, slot: int) -> None,
Poller.modify_monitor(monitor: MonitorSocket, events: PollEventFlag) -> None and
Poller.remove_monitor(monitor: MonitorSocket) -> None register, modify and remove a socket monitor as a
poller source (common spec "Monitor sources in Poller"); the existing add_socket/modify_socket/remove_socket
also accept a monitor. Only PollEventFlag.POLLIN is valid for a monitor; any other bit is rejected with a typed
ConfigResult.INVALID_ARGUMENT. Drain with monitor.recv(RecvFlags.DONT_WAIT) after readiness; the poll event
reports the monitor through the same slot and source kind as a socket.
Only module-private _reply_token_from_native creates a ReplyToken; public construction and
serialization are rejected. The factory fills private _owner and _value through
object.__new__(ReplyToken) and object.__setattr__. Equality and hashing use both owner identity and
the opaque value. StreamPacket is an empty reusable output. Publish preserves its existing flags and
synchronous submit semantics on a PublishOp separate from send. A token provides no raw property,
int() conversion, ordering, or close(). Concurrent recv into the same output is invalid-state.
Message references remain valid only until the next recv entry or close(). Before the first
bind/connect, the recv_mode setter accepts only RAW and PACKET and rejects UNSPECIFIED.
Public interface¶
class SendSubmission:
@property
def result(self) -> SubmitResult: ... # OK | BACKPRESSURED, submit-time snapshot
@property
def admitted(self) -> Awaitable[None]: ... # completed when result is OK
class RequestSubmission:
@property
def result(self) -> SubmitResult: ...
@property
def admitted(self) -> Awaitable[None]: ...
@property
def reply(self) -> Awaitable[list[Message]]: ... # completes after successful admission
class SendOp(Protocol):
def message(self, payload) -> "SendOp": ...
def messages(self, *payloads) -> "SendOp": ...
def submit(self) -> SendSubmission: ...
def submit_sync(self) -> None: ...
class RequestOp(Protocol):
def message(self, payload) -> "RequestOp": ...
def messages(self, *payloads) -> "RequestOp": ...
def timeout(self, timeout) -> "RequestOp": ...
def submit(self) -> RequestSubmission: ...
def submit_sync(self) -> list[Message]: ...
class ReplyOp(Protocol):
def message(self, payload) -> "ReplyOp": ...
def messages(self, *payloads) -> "ReplyOp": ...
def submit(self) -> None: ...
class PublishOp(Protocol):
def message(self, payload) -> "PublishOp": ...
def messages(self, *payloads) -> "PublishOp": ...
def flags(self, flags) -> "PublishOp": ...
def submit(self) -> None: ...
@final
class ReplyToken:
__slots__ = ("_owner", "_value")
def __new__(cls) -> NoReturn:
raise TypeError("ReplyToken is created by ROUTER request receive")
def __eq__(self, other: object) -> bool: ...
def __hash__(self) -> int: ...
def __repr__(self) -> str: return "ReplyToken()"
def __copy__(self) -> "ReplyToken": return self
def __deepcopy__(self, memo) -> "ReplyToken": return self
def __reduce_ex__(self, protocol):
raise TypeError("ReplyToken cannot be serialized")
class StreamRecvMode(IntEnum):
UNSPECIFIED = 0
RAW = 1
PACKET = 2
class StreamSocketOptions(Protocol):
@property
def recv_mode(self) -> StreamRecvMode: ...
@recv_mode.setter
def recv_mode(self, mode: StreamRecvMode) -> None: ...
class StreamPacket:
routing_id: Optional[RoutingId]
header: Optional[Message]
body: Optional[Message]
def __init__(self) -> None: ...
@property
def is_empty(self) -> bool: ...
def close(self) -> None: ...
def __enter__(self) -> "StreamPacket": ...
def __exit__(self, exc_type, exc, tb) -> None: ...
class Received:
routing_id: Optional[RoutingId]
reply_token: Optional[ReplyToken]
class StreamSocket:
def send(self, routing_id: RoutingId) -> SendOp: ...
def recv_into(
self, out: Received, *, flags: RecvFlags = RecvFlags.NONE
) -> bool: ...
def recv_packet_into(
self, out: StreamPacket, *, flags: RecvFlags = RecvFlags.NONE
) -> bool: ...
The operation-start signatures are PAIR send() -> SendOp, DEALER send() -> SendOp and
request() -> RequestOp, ROUTER send(routing_id) -> SendOp,
request(routing_id) -> RequestOp, and reply(routing_id, token) -> ReplyOp, and STREAM
send(routing_id) -> SendOp. A send factory captures the target in the builder. PUB and XPUB
publish(topic) return PublishOp. Received.send() returns a SendOp that captures the source
target, and Received.reply() returns a ReplyOp that captures the source RID and token.
The public Python surface contains no RoutedSendOp, StreamSocket.send_async() or on_packet(),
request callback argument, reply _FlaggedFluentMessageOp or flags, monitor ignore_handler or
on_event, timer on_fire, or pair/generation member.
Monitor provides recv(*, flags=RecvFlags.NONE) -> Optional[MonitorEvent], status(), and close().
Timer provides start(interval_ns:int, repeat_count:int), stop(), recv() -> Optional[int], and
close(). Monitor-event connection_id is used only for diagnostics and correlation, not as a
send/reply target or reconnect fence. The internal FFI enum mirror uses only
ZLINK_OPT_PENDING_MAX_MSGS and ZLINK_OPT_PENDING_MAX_BYTES and adds no public option property.
Implementation and contract-test verification requirements¶
Verify the following using only public Python protocols, results, exceptions, and poller events. Each item maps to one contract test.
Operations and completion
- Send/request expose only the flag-free awaitable and synchronous terminals in the Public interface section and retain request timeout.
publish(topic)returns a separatePublishOpwith publish flags and synchronous submit.- Common completion, cancellation, and poller observations follow the execution-model verification requirements.
ReplyToken and STREAM
- Public
ReplyToken()and pickle serialization fail;copy.copy()andcopy.deepcopy()return the same immutable valid token. - Only tokens with the same owner and value are equal, and reply with a token from another owner fails before the native call.
recv_packet_into()fills the output after success and leaves it empty on no-data or error. The output can be reused afterclose().
Pull eventing
- Monitor and timer recv return no-data as
Noneand expose events and fire counts without callbacks.