Files
danos/library/kernel/service.zig
T
Daniel Samson 0faa0fd21b service: on_buffered_message for pushed events (S5)
An additive, behavior-neutral callback. A buffered async message
(Received.isMessage — a pushed event from a provider this service
subscribed to) carries a payload in the receive buffer; run() now hands
it to on_buffered_message before falling through to on_notification with
the badge, so a coalesced timer/exit riding the same wake is not lost. A
service that does not set the callback (all of them today) is unchanged:
the isMessage branch is a no-op and on_notification still runs, exactly
as before.

This is the seam the volume manager needs to consume block
medium_changed: the event's reserved op number collides with the volume
manager's own hello, so it must be decoded by hand here, never through
the protocol dispatch. Full suite neutral (the lone device-authority
miss is a known child-cleanup race that passes on rerun).
2026-08-10 04:43:27 +01:00

401 lines
21 KiB
Zig

//! The service harness (docs/process-lifecycle.md): one replyWait loop that
//! folds protocol requests, signals, and subscribed notifications into
//! callbacks — so the lifecycle contract ("answers ping, exits on terminate")
//! is satisfied by construction and a service author writes domain logic only.
//! Nothing is asynchronous inside the process: a callback runs at a point the
//! loop chose, never on a hijacked stack — the whole reason signals are
//! messages.
//!
//! One rule a service author does have to know, and it is stated on
//! `Callbacks.on_message`: **a capability that arrives belongs to the turn** —
//! the loop closes it unless the callback claims it with `take()`. Forgetting is
//! therefore safe, and keeping is explicit; the opposite arrangement quietly
//! spends a handle-table slot per request.
//!
//! The harness also owns the **subscriber side** of a protocol that declares
//! `.events` — see `Subscribers`. The table, the reserved subscribe/unsubscribe
//! verbs, the fan-out, and the dead-subscriber sweep live here rather than in
//! each provider, so every event stream in the system has identical semantics
//! (docs/os-development/protocol-namespace.md, "Wiring").
//!
//! The liveness probe: a **zero-length request is the universal ping**, answered
//! with a zero-length reply by the harness itself. No protocol's requests start
//! at length zero, so the encoding cannot collide, and there is nothing for a
//! service author to implement — a wedged service simply fails to answer, which
//! is the diagnosis (see docs/ipc.md).
const channel = @import("channel");
const envelope = @import("envelope");
const ipc = @import("ipc");
const process = @import("process");
/// The harness's handle on a provider's subscriber table, type-erased because
/// `run` is not generic over the protocol while `Subscribers` is. A service names
/// its table once, as `Callbacks.subscribers`, and the loop does the rest: it
/// subscribes to published process exits at startup and drops a dead task's
/// subscriptions before the service's own notification callback ever sees the
/// badge.
pub const SubscriberHooks = struct {
/// Ask the kernel for published exit events on this service's endpoint.
watch: *const fn (endpoint: ipc.Handle) void,
/// Drop everything task `dead` had subscribed.
forget: *const fn (dead: u32) void,
};
pub const Callbacks = struct {
/// Called once with the service's endpoint before the loop starts — the
/// place to subscribe to exit events, bind IRQs, or announce readiness.
/// Return false to abort startup (the process exits).
init: ?*const fn (endpoint: ipc.Handle) bool = null,
/// One protocol request from `sender` (a task id): write the reply into
/// `reply`, return its length. The zero-length ping never reaches this.
///
/// `arrived` is the capability the request carried (M13 cap passing — how a
/// subscriber hands over its endpoint), and it comes with **an ownership
/// rule: the turn owns it, and a handler that wants to keep it must say so
/// with `take()`.** Whatever is left when this returns, the loop closes.
/// `peek()` reads it without claiming, which is what a handler that may
/// still refuse wants — no close of its own on the refusal paths.
///
/// The rule is stated here, in the contract, because the alternative has
/// failed in practice: an implementation that simply ignored a `?ipc.Handle`
/// argument leaked a handle table slot per request, and every operation
/// except a subscribe ignores it. Thirty-two such requests — zero-length
/// pings will do, and they need no authorization — and the service can never
/// accept another capability for the rest of the boot. See `ipc.Arrival`.
on_message: *const fn (message: []const u8, reply: []u8, sender: u32, arrived: *ipc.Arrival) usize,
/// A notification that is not a signal — a subscribed exit event, a bound
/// IRQ, a timer landing. The raw badge; decode with the ipc helpers.
on_notification: ?*const fn (badge: u64) void = null,
/// A buffered async message (`Received.isMessage`): a pushed event from a
/// provider this service subscribed to, its payload in the receive buffer.
/// Unlike `on_message`, it never goes through the protocol dispatch — so an
/// event whose reserved op number collides with one of this service's own
/// verbs (a `block` `medium_changed` reaching the volume manager, whose own
/// protocol numbers `hello` the same) is decoded by hand here, not
/// mis-dispatched. Default null: the badge alone still reaches
/// `on_notification`, exactly as before this callback existed.
on_buffered_message: ?*const fn (message: []const u8) void = null,
/// The reload signal. Default: ignored.
on_reload: ?*const fn () void = null,
/// The terminate signal, called before the loop returns. The clean exit is
/// the return itself — never put *necessary* work here (iron rule 1: a kill
/// arrives with no warning; this is for graceful extras only).
on_terminate: ?*const fn () void = null,
/// The contract this service provides: a name under `/protocol`, mirroring
/// the `library/protocol/` module that defines the wire format — a program
/// imports `display-protocol` and the provider binds `"display"`
/// (docs/os-development/protocol-namespace.md). Bound at startup, before
/// `init` runs, so the service is reachable the moment it serves. A refusal
/// (not granted, or a live provider already holds the name) aborts startup.
service: ?[]const u8 = null,
/// This provider's subscriber table — `Subscribers(Protocol, Context).hooks`
/// — for a protocol that declares `.events`. Naming it here is what buys the
/// exit-notification sweep: the loop subscribes to published deaths at
/// startup and releases a dead subscriber's slot (and the endpoint capability
/// in it) when one lands.
subscribers: ?SubscriberHooks = null,
};
/// How many subscribers one provider fans out to. Bounded like every table in
/// this system; a subscribe past the end is refused with `-ENOSPC` rather than
/// silently forgetting an earlier one.
pub const subscriber_capacity = 8;
/// The interest mask that means "every event of this protocol" — what a
/// subscriber which named no class gets, and what a provider passes when the
/// event it is publishing belongs to no class.
pub const every_event: u32 = 0;
/// The subscriber side of a protocol, for a provider whose contract declares
/// `.events` (docs/os-development/protocol-namespace.md: *the harness owns the
/// machinery — the subscriber table, the dead-subscriber sweep, and the fan-out
/// loop*). Three services hand-rolled this, with three different ideas of when a
/// dead subscriber goes away — a poll of the process list on subscribe, a drop on
/// a failed send, and nothing at all. This is the one idiom.
///
/// ```zig
/// const Subscriptions = service.Subscribers(power_protocol.Protocol, void);
/// ...
/// fn onMessage(message: []const u8, reply: []u8, sender: u32, arrived: *ipc.Arrival) usize {
/// return Subscriptions.dispatch({}, handlers, message, sender, arrived, reply);
/// }
/// pub fn main() void {
/// service.run(power_protocol.message_maximum, .{
/// .service = "power",
/// .on_message = onMessage,
/// .subscribers = Subscriptions.hooks,
/// });
/// }
/// ```
///
/// What the provider still writes is its own events — `publish(.power_button, 0,
/// .{})`. Everything else happens here: registering the caller's endpoint on the
/// reserved `subscribe` verb, taking that capability out of the turn, dropping it
/// on `unsubscribe` or on the subscriber's death, and framing one packet for the
/// whole fan-out.
///
/// The table is per instantiation (a container-level `var` inside the generic
/// type), so a process providing two contracts gets two tables and neither can
/// see the other's subscribers.
pub fn Subscribers(comptime Protocol: type, comptime Context: type) type {
return struct {
/// The generated dispatch this provider answers with.
pub const Provider = Protocol.Provider(Context);
pub const Handlers = Provider.Handlers;
/// One registered subscriber: the endpoint events are pushed to (the
/// capability it handed over at subscribe time, which this slot owns),
/// the task that handed it over — the kernel-stamped badge, the only
/// source identity there is — and which classes of event it asked for.
const Slot = struct {
used: bool = false,
endpoint: ipc.Handle = 0,
task: u32 = 0,
interest: u32 = every_event,
};
var slots: [subscriber_capacity]Slot = .{Slot{}} ** subscriber_capacity;
/// Set when a slot has taken the capability the turn carried, and read
/// back in `dispatch`, which is where the turn's `Arrival` lives. The
/// generated dispatch hands a handler the raw handle rather than the
/// `Arrival` — deliberately, since a handler has no business closing the
/// turn's property — so the *claim* has to travel back out this way. One
/// turn, one handler, one thread: there is nothing here to race.
var claimed = false;
/// What `Callbacks.subscribers` is given.
pub const hooks: SubscriberHooks = .{ .watch = watchExits, .forget = forget };
fn watchExits(endpoint: ipc.Handle) void {
// Published exits, not a poll of the process list: a service must
// never depend on clients cleaning up after themselves, and it must
// not have to walk the whole table on every subscribe to find out
// either (docs/process-lifecycle.md, "Who learns of a death").
_ = process.subscribeExits(endpoint);
}
/// Release everything task `dead` had subscribed. The slot owns the
/// endpoint capability, so reclaiming the slot closes it — otherwise a
/// process that subscribes and dies costs a handle-table slot that never
/// comes back.
pub fn forget(dead: u32) void {
for (&slots) |*slot| {
if (slot.used and slot.task == dead) {
_ = ipc.close(slot.endpoint);
slot.* = .{};
}
}
}
/// Whether `task` is a subscriber — the gate for an operation a provider
/// honours from its subscribers and nobody else. The power service's
/// shutdown is the one: the badge is kernel-stamped, so nothing in a
/// packet can claim to be the subscriber that already ran the stop
/// sequence.
pub fn has(task: u32) bool {
for (&slots) |*slot| {
if (slot.used and slot.task == task) return true;
}
return false;
}
/// A provider's own handler kept this turn's capability (stored it
/// somewhere with a lifetime beyond the turn) — the same claim the
/// reserved `subscribe` makes for its slot table. Composes with it:
/// one flag, one `take()`, whoever claims first wins the turn.
pub fn claimArrival() void {
claimed = true;
}
/// Answer one received packet, with the reserved `subscribe` and
/// `unsubscribe` verbs already wired — a provider that leaves those two
/// handlers null (every provider should) gets the harness's. The turn's
/// capability is peeked, never taken, unless a slot actually kept it.
pub fn dispatch(
context: Context,
handlers: Handlers,
packet: []const u8,
sender: u32,
arrived: *ipc.Arrival,
reply: []u8,
) usize {
var wired = handlers;
if (wired.subscribe == null) wired.subscribe = onSubscribe;
if (wired.unsubscribe == null) wired.unsubscribe = onUnsubscribe;
claimed = false;
const written = Provider.dispatch(context, wired, packet, sender, arrived.peek(), reply);
if (claimed) _ = arrived.take();
return written;
}
/// Push one event to every subscriber.
pub fn publish(
comptime event: Protocol.Event,
target: u64,
payload: Protocol.PayloadOf(event),
) void {
publishClass(event, target, payload, every_event);
}
/// Push one event to the subscribers whose interest mask includes
/// `class` (a subscriber that named no class takes everything). The
/// packet is framed **once**, outside the loop, so every subscriber of a
/// class receives identical bytes; and delivery is `ipc.send`, which
/// never blocks, so one slow or dead subscriber can never stall the rest
/// — the whole reason broadcast is a provider pattern and not a kernel
/// primitive.
pub fn publishClass(
comptime event: Protocol.Event,
target: u64,
payload: Protocol.PayloadOf(event),
class: u32,
) void {
var packet: [envelope.post_maximum]u8 = undefined;
const framed = Protocol.encodeEvent(event, target, payload, &packet) orelse return;
for (&slots) |*slot| {
if (!slot.used) continue;
if (!wants(slot.*, class)) continue;
// The sweep is what normally reclaims a dead subscriber, promptly
// and with its capability closed. This is the backstop for a
// notification that never arrived: an endpoint's notify ring is
// bounded, so a burst of deaths can drop one, and a send to an
// endpoint whose owner is gone fails rather than blocking.
if (!ipc.send(slot.endpoint, framed)) {
_ = ipc.close(slot.endpoint);
slot.* = .{};
}
}
}
fn wants(slot: Slot, class: u32) bool {
if (class == every_event) return true; // the event belongs to no class
if (slot.interest == every_event) return true; // the subscriber named none
return slot.interest & class != 0;
}
/// The reserved `subscribe` verb: register the caller's endpoint (the
/// call's capability) for the classes its tail names. A refusal simply
/// returns and the turn closes what arrived — the harness's ownership
/// rule (`ipc.Arrival`), which is why a subscribe storm against a full
/// table cannot spend the handle table.
fn onSubscribe(_: Context, invocation: envelope.Invocation(void), _: envelope.Answer(void)) isize {
const endpoint = invocation.capability orelse return -envelope.EPROTO; // no endpoint passed
const interest = envelope.decodeSubscribe(invocation.tail).interest;
for (&slots) |*slot| {
if (slot.used) continue;
// Appended, not replaced: one task may hold several subscriptions
// on different endpoints (a client taking keyboard and mouse as
// two streams), and each is its own conversation.
slot.* = .{ .used = true, .endpoint = endpoint, .task = invocation.sender, .interest = interest };
claimed = true; // the table holds it until that task dies
return 0;
}
return -envelope.ENOSPC; // table full
}
/// The reserved `unsubscribe` verb: every subscription the calling task
/// holds here goes, which is exactly what its death would do. It names no
/// endpoint because the badge already names the only subscriber a caller
/// can speak for — its own.
fn onUnsubscribe(_: Context, invocation: envelope.Invocation(void), _: envelope.Answer(void)) isize {
if (!has(invocation.sender)) return -envelope.ENOENT;
forget(invocation.sender);
return 0;
}
};
}
/// Run the service: create the endpoint, bind it under the service's contract
/// name (if it has one), bind signals to it, call `init`, then serve until
/// `terminate` arrives — at which point the loop returns and main's return is
/// the clean exit the supervisor reads as `ExitReason.exited`.
/// `maximum_message` sizes the receive and reply buffers (a service passes its
/// protocol's message maximum).
/// The capability the current turn's handler nominates to ride out with its
/// reply — init's registry idiom (`pending_capability`), lifted into the
/// harness so any service can answer an establishment request with a channel
/// (communication.md "Establishment: two planes"). Consumed by the loop at the
/// very next `replyWait`, which is the reply this turn owes; null is the
/// untouched common path. The kernel shares the endpoint as a refcounted copy,
/// so the nominating service keeps its own handle.
var pending_reply_capability: ?ipc.Handle = null;
/// Called from inside an `on_message` handler: send `handle` with this turn's
/// reply. One capability per turn — the last nomination wins, matching the
/// transport (a reply carries at most one).
pub fn replyWithCapability(handle: ipc.Handle) void {
pending_reply_capability = handle;
}
pub fn run(comptime maximum_message: usize, callbacks: Callbacks) void {
const endpoint = ipc.createIpcEndpoint() orelse return;
if (callbacks.service) |name| {
if (!channel.bindPatiently(name, endpoint)) return;
}
_ = process.bindSignals(endpoint);
// Before `init`, so a subscriber that arrives the instant the name is bound
// is already covered by the sweep that will release it.
if (callbacks.subscribers) |subscribers| subscribers.watch(endpoint);
if (callbacks.init) |initialise| {
if (!initialise(endpoint)) return;
}
var reply_buffer: [maximum_message]u8 = undefined;
var reply_len: usize = 0;
var receive: [maximum_message]u8 = undefined;
while (true) {
// The reply going out is the one the just-run handler wrote, so the
// capability it nominated (if any) rides this exact replyWait and is
// reset before the next turn can see a stale one.
const reply_capability = pending_reply_capability;
pending_reply_capability = null;
const got = ipc.replyWait(endpoint, reply_buffer[0..reply_len], &receive, reply_capability);
// Whatever capability came with this turn is the turn's, and the turn
// closes it unless a callback claims it (`ipc.Arrival`). Structural
// rather than a close per branch, because the branches are exactly what
// gets forgotten: the ping's `continue` below, and every `on_message`
// that has no use for a capability — which is every operation but a
// subscribe. A `defer` in a loop body runs on `continue` and on the
// `return` that ends the loop, so this covers all four exits.
var arrived: ipc.Arrival = .{ .handle = got.cap };
defer arrived.release();
if (got.isNotification()) {
reply_len = 0; // nothing owed for a notification
if (process.signalsFrom(got.badge)) |signals| {
if (signals.has(.reload)) {
if (callbacks.on_reload) |onReload| onReload();
}
if (signals.has(.terminate)) {
if (callbacks.on_terminate) |onTerminate| onTerminate();
return; // the loop's return IS the clean exit
}
continue;
}
// A death sweeps the subscriber table first, then still reaches the
// service: a provider often has its own per-client state to release
// (open file handles, device tokens, layers) and the same badge is
// the notice for both.
if (got.isChildExit()) {
if (callbacks.subscribers) |subscribers| subscribers.forget(got.childProcessId());
}
// A buffered async message (a pushed event) carries a payload; hand it
// to the service that asked for it. The badge still reaches
// on_notification below, so a coalesced timer/exit riding the same wake
// is not lost — and a service without this callback is unchanged.
if (got.isMessage()) {
if (callbacks.on_buffered_message) |onBuffered| onBuffered(receive[0..got.len]);
}
if (callbacks.on_notification) |onNotification| onNotification(got.badge);
continue;
}
if (got.len == 0) {
reply_len = 0; // the universal ping: a zero-length reply, from the harness
continue; // any capability it carried goes out through the turn's `defer`
}
reply_len = callbacks.on_message(receive[0..got.len], &reply_buffer, got.senderTaskId(), &arrived);
}
}