library: the harness keeps the subscribers, and an id belongs to whoever opened it

Three services had each written the same thing and got it three different
ways: input polled the process list to notice a dead subscriber, and only
when someone else subscribed; the power service never noticed at all; the
device manager noticed drivers but not subscribers. The harness owns the
table now, driven by the events a protocol declares — it registers on the
reserved verb, frames each event once, posts to everyone interested without
waiting on any of them, and reclaims a slot when the kernel says its owner
died. Interest masks moved to the envelope, so a subscriber that wants only
mice asks the same way everywhere.

Two consequences the plan had not foreseen. The device manager now hears a
supervised child's death twice, once as its supervisor and once as a
subscriber, so restart backoff counted every crash twice and gave up after
half as many; it retires the id before counting. And the kernel's published
exit table had eight slots for what is now six subscriptions in a plain
boot, so it holds sixteen.

The other half is a hole the design named early and left standing: a
backend handed out a small integer and then honoured it from anyone. A
process that guessed a file's node id read another client's file; a display
layer had no owner at all, so any client could reconfigure or destroy any
layer; a USB device token was never checked against the client that opened
it. Each is now bound to the task that opened it, and a wrong owner gets
exactly what an unknown id gets — the refusal must not become the oracle
the identical answers elsewhere were designed to remove. Closing a file
changed with it: it used to succeed unconditionally, which would have told
a caller which ids existed.

Suite 111/111, with a new case in which one process holds a file and a
layer, hands both ids to a second process, and finds them untouched after
that process has tried everything with them.
This commit is contained in:
Daniel Samson
2026-08-01 09:05:26 +01:00
parent 2719b93530
commit 1b1c587c14
29 changed files with 1072 additions and 335 deletions
+32 -70
View File
@@ -60,15 +60,14 @@ const pwrbtn_bit: u16 = 1 << 8;
const sci_en_bit: u32 = 1 << 0;
const slp_en: u32 = 1 << 13;
// The `.power` subscribers: endpoints handed over as capabilities, each
// receiving events as buffered messages. Dropped on a failed send. The
// subscriber's task id is kept too — a shutdown request is honored only from a
// subscriber (init subscribes; a stray process does not), the soft gate that
// stands in for "only the system supervisor may power off" without hardcoding
// a pid the kernel's idle tasks would have taken.
const maximum_subscribers = 8;
var subscribers: [maximum_subscribers]?ipc.Handle = .{null} ** maximum_subscribers;
var subscriber_tasks: [maximum_subscribers]u32 = .{0} ** maximum_subscribers;
/// The `.power` subscribers, kept by the service harness (P4c): endpoints handed
/// over as capabilities, each receiving events as buffered messages, each swept
/// when its task dies. The table remembers which task subscribed, which is what
/// the shutdown gate below asks — a shutdown request is honored only from a
/// subscriber (init subscribes; a stray process does not), the soft gate that
/// stands in for "only the system supervisor may power off" without hardcoding a
/// pid the kernel's idle tasks would have taken.
const Subscriptions = service.Subscribers(power_protocol.Protocol, void);
// Pass-1 registration record (see main): what pass 2 reports.
const Registered = struct { hid: [8]u8 = .{0} ** 8, hid_len: usize = 0, device_id: u64 = 0, resource_count: u64 = 0 };
@@ -201,6 +200,7 @@ pub fn main(init: process.Init) void {
.init = onInit,
.on_message = onMessage,
.on_notification = onNotification,
.subscribers = Subscriptions.hooks,
});
}
@@ -407,13 +407,13 @@ fn publishNotify(node: *aml.Node, code: u64) void {
const notice = power_protocol.Notice{ .code = @truncate(code), .hid = hid };
std.log.info("power: notify {s} code {d}", .{ hid[0..7], code });
if (std.mem.eql(u8, hid[0..7], "PNP0C0A")) {
publish(.battery, notice);
Subscriptions.publish(.battery, 0, notice);
} else if (std.mem.eql(u8, hid[0..7], "ACPI0003")) {
publish(.ac, notice);
Subscriptions.publish(.ac, 0, notice);
} else if (std.mem.eql(u8, hid[0..7], "PNP0C0D")) {
publish(.lid, notice);
Subscriptions.publish(.lid, 0, notice);
} else {
publish(.notify, notice);
Subscriptions.publish(.notify, 0, notice);
}
}
@@ -425,28 +425,10 @@ fn writeHex2(out: []u8, n: u32) void {
}
fn publishButton() void {
publish(.power_button, .{});
}
/// Push one event to every subscriber. The kind is the packet's operation, so
/// this is framed once, outside the loop — every subscriber gets identical
/// bytes. A subscriber whose endpoint stops accepting (it died) is dropped on
/// the failed send, so a dead one can never stall the rest.
fn publish(comptime kind: power_protocol.Event, notice: power_protocol.Notice) void {
var packet: [envelope.post_maximum]u8 = undefined;
const framed = power_protocol.Protocol.encodeEvent(kind, 0, notice, &packet) orelse return;
for (&subscribers) |*slot| {
if (slot.*) |handle| {
if (!ipc.send(handle, framed)) slot.* = null;
}
}
}
fn isSubscriber(task: u32) bool {
for (&subscribers, 0..) |*slot, si| {
if (slot.* != null and subscriber_tasks[si] == task) return true;
}
return false;
// The kind is the packet's operation, so the harness frames it once and pushes
// the same bytes to every subscriber. There is no class here: a power event
// goes to everyone who asked for power events.
Subscriptions.publish(.power_button, 0, .{});
}
/// Enter S5 (soft off): write SLP_TYP|SLP_EN to the PM1 control register(s).
@@ -468,57 +450,37 @@ fn enterS5() void {
// --- harness callbacks --------------------------------------------------------
fn onNotification(badge: u64) void {
// The only notification the service binds is the SCI (an IRQ badge).
_ = badge;
// Two kinds of notification reach this loop now. The SCI is the one this
// service binds; the published process exits are the harness's, which it has
// already used to sweep the subscriber table before calling here. Everything
// that is not a bare IRQ badge must therefore be ignored — treating a death
// as an interrupt would clear PM1 status the firmware never set.
if (badge & (ipc.notify_exit_bit | ipc.notify_timer_bit | ipc.notify_message_bit | ipc.notify_signal_bit) != 0) return;
onSci();
}
/// The generated power dispatch. One provider per system, so the handler context
/// is empty and the subscriber table stays in this file's globals.
const Serve = power_protocol.Protocol.Provider(void);
const Invocation = envelope.Invocation;
const Answer = envelope.Answer;
/// Set by `onSubscribe` when the subscriber table has taken the capability the
/// call carried, and read by `onMessage`, where the turn's `Arrival` lives.
var capability_claimed = false;
/// The power contract: the reserved `subscribe` (the subscriber's endpoint as
/// the call's capability) and `shutdown` (subscribers only). Device discovery
/// uses a different endpoint — the device manager's — so nothing here handles a
/// tree report.
/// the call's capability, answered by the harness) and `shutdown` (subscribers
/// only). Device discovery uses a different endpoint — the device manager's — so
/// nothing here handles a tree report.
fn onMessage(message: []const u8, reply: []u8, sender: u32, arrived: *ipc.Arrival) usize {
capability_claimed = false;
const written = Serve.dispatch({}, handlers, message, sender, arrived.peek(), reply);
if (capability_claimed) _ = arrived.take();
return written;
return Subscriptions.dispatch({}, handlers, message, sender, arrived, reply);
}
const handlers = Serve.Handlers{ .shutdown = onShutdown, .subscribe = onSubscribe };
/// The subscriber's endpoint is claimed only when a slot takes it; a full table
/// refuses and the turn closes what arrived.
fn onSubscribe(_: void, invocation: Invocation(void), _: Answer(void)) isize {
const endpoint = invocation.capability orelse return -envelope.EPROTO;
for (&subscribers, 0..) |*slot, index| {
if (slot.* == null) {
slot.* = endpoint;
subscriber_tasks[index] = invocation.sender;
capability_claimed = true;
return 0;
}
}
return -envelope.ENOSPC;
}
/// `subscribe` and `unsubscribe` are absent on purpose: the harness answers both.
const handlers = Subscriptions.Handlers{ .shutdown = onShutdown };
/// Honored only from a power subscriber — init, which has already run the stop
/// sequence over everything else. The power service is mechanism (write S5);
/// deciding *when* to shut down and stopping the rest of the system first is
/// init's policy. The badge is the whole gate: it is kernel-stamped, so nothing
/// in the packet can claim to be init.
/// in the packet can claim to be init. The subscriber table moved into the
/// harness; the question it answers has not changed.
fn onShutdown(_: void, invocation: Invocation(void), _: Answer(void)) isize {
if (!isSubscriber(invocation.sender)) return -envelope.EPERM;
if (!Subscriptions.has(invocation.sender)) return -envelope.EPERM;
enterS5();
return 0;
}
@@ -27,9 +27,11 @@ const device_manager_protocol = @import("device-manager-protocol");
const envelope = @import("envelope");
const registry = @import("device-registry");
/// The generated device-manager dispatch. One manager per system, so the handler
/// context is empty and the tables stay in this file's globals.
const Serve = device_manager_protocol.Protocol.Provider(void);
/// The generated device-manager dispatch, plus the subscriber machinery the
/// harness owns (P4c): the watcher table, the reserved `subscribe` verb, the
/// exit sweep, and the fan-out. One manager per system, so the handler context is
/// empty and the tables stay in this file's globals.
const Serve = service.Subscribers(device_manager_protocol.Protocol, void);
const Invocation = envelope.Invocation;
const Answer = envelope.Answer;
@@ -156,31 +158,6 @@ var test_scanout_killed = false;
var test_kill_pid: u32 = 0;
var test_kill_due_ns: u64 = 0;
/// The application subscribers (M18.3, the input-service pattern): endpoints
/// handed over as capabilities, each receiving every child add/remove as a
/// buffered message. A subscriber whose endpoint stops accepting (it died) is
/// dropped on the failed send.
const maximum_subscribers = 8;
var subscribers: [maximum_subscribers]?ipc.Handle = .{null} ** maximum_subscribers;
/// Push one event to every subscriber: the same struct a bus driver *called*
/// with, framed as an event instead — one encoding, both directions, told apart
/// by the packet's verb rather than by anything inside it. A subscriber whose
/// endpoint stops accepting (it died) is dropped on the failed send.
fn publish(
comptime event: device_manager_protocol.Event,
target: u64,
payload: device_manager_protocol.Protocol.PayloadOf(event),
) void {
var packet: [envelope.post_maximum]u8 = undefined;
const framed = device_manager_protocol.Protocol.encodeEvent(event, target, payload, &packet) orelse return;
for (&subscribers) |*slot| {
if (slot.*) |handle| {
if (!ipc.send(handle, framed)) slot.* = null; // dead subscriber
}
}
}
/// The manager's mirror of what bus drivers report (docs/device-manager.md "the
/// tree"): the children, keyed by (parent, bus address), each remembering which
/// driver instance reported it — that is what death-pruning sweeps by.
@@ -225,7 +202,7 @@ fn pruneChildrenOf(reporter: u32) void {
if (child.used and child.reporter == reporter) {
std.log.info("child removed (device {d} port {d})", .{ child.parent, child.bus_address });
child.used = false;
publish(.child_removed, 0, .{ .parent = child.parent, .bus_address = child.bus_address });
Serve.publish(.child_removed, 0, .{ .parent = child.parent, .bus_address = child.bus_address });
}
}
}
@@ -240,7 +217,11 @@ fn childCountOf(reporter: u32) u32 {
return n;
}
/// The driver entry a live process id belongs to. Zero is not a process id here:
/// it is what `onDriverExit` writes back to retire an id it has already acted on,
/// so a second notification for the same death matches nothing.
fn driverByProcess(process_id: u32) ?*Driver {
if (process_id == 0) return null;
for (&drivers) |*driver| {
if (driver.used and driver.process_id == process_id) return driver;
}
@@ -308,8 +289,17 @@ fn spawnDriver(driver: *Driver) void {
/// is the whole restart decision: a clean exit meant to stop; anything else
/// restarts with backoff until the crash-loop cap.
fn onDriverExit(driver: *Driver) void {
pruneChildrenOf(driver.process_id);
const reason = process.exitReason(driver.process_id) orelse .fault;
const dead = driver.process_id;
// One death, two notifications: the manager is this driver's supervisor (its
// spawn named this endpoint) *and*, since P4c put the watcher table in the
// harness, a subscriber to published exits. Both badges carry the same id, and
// the ring delivers them separately — so the id is retired here, before any
// decision is taken, and the second notification finds no driver to act on.
// Without this the backoff would count one death twice and the crash-loop cap
// would fire at half the deaths it names.
driver.process_id = 0;
pruneChildrenOf(dead);
const reason = process.exitReason(dead) orelse .fault;
if (reason == .exited) {
driver.state = .stopped;
std.log.info("{s} exited cleanly; not restarting", .{driver.name()});
@@ -410,27 +400,17 @@ fn initialise(endpoint: ipc.Handle) bool {
return true;
}
/// Set by `onSubscribe` when the subscriber table has taken ownership of the
/// capability the call carried, and read by `onMessage`, 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* travels back out this way. One turn, one
/// handler, one thread: there is nothing here to race.
var capability_claimed = false;
fn onMessage(message: []const u8, reply: []u8, sender: u32, arrived: *ipc.Arrival) usize {
capability_claimed = false;
const written = Serve.dispatch({}, handlers, message, sender, arrived.peek(), reply);
if (capability_claimed) _ = arrived.take();
return written;
return Serve.dispatch({}, handlers, message, sender, arrived, reply);
}
/// `subscribe` and `unsubscribe` are absent on purpose: the harness answers both,
/// and its table is what `publish` fans out over.
const handlers = Serve.Handlers{
.hello = onHello,
.child_added = onChildAdded,
.child_removed = onChildRemoved,
.enumerate = onEnumerate,
.subscribe = onSubscribe,
};
/// The handshake. The device this driver was assigned is the packet's target.
@@ -470,7 +450,7 @@ fn onChildAdded(_: void, invocation: Invocation(device_manager_protocol.ChildAdd
if (driverByProcess(sender)) |driver| {
if (!addChild(report.parent, report.bus_address, report.identity, device_id, sender)) status = -envelope.ENOSPC;
std.log.info("child added (device {d} port {d}, identity {d}) by {s}", .{ report.parent, report.bus_address, report.identity, driver.name() });
if (status == 0) publish(.child_added, device_id, report);
if (status == 0) Serve.publish(.child_added, device_id, report);
// Matching from reports (M19.3), now data-driven via the /system/configuration/devices.csv
// registry: a registered child gets the most-specific driver its identity
// matches, once — re-reports after a bus restart dedupe on the registered
@@ -556,24 +536,11 @@ fn onEnumerate(_: void, _: Invocation(void), answer: Answer(void)) isize {
return @intCast(written);
}
/// The reserved `subscribe` verb: an application's endpoint arrived as the call's
/// capability. The table taking a slot is what claims it; a full table refuses
/// and lets the turn close it, so a subscribe storm cannot spend the handle table.
fn onSubscribe(_: void, invocation: Invocation(void), _: Answer(void)) isize {
const endpoint = invocation.capability orelse return -envelope.EPROTO;
for (&subscribers) |*slot| {
if (slot.* == null) {
slot.* = endpoint;
capability_claimed = true; // the table holds it from here
return 0;
}
}
return -envelope.ENOSPC;
}
fn onNotification(badge: u64) void {
if (badge & ipc.notify_exit_bit != 0) {
const dead: u32 = @intCast(badge & ~(ipc.notify_badge_bit | ipc.notify_exit_bit));
// The harness has already swept the watcher table for this death; what is
// left is the manager's own concern, its supervised drivers.
if (driverByProcess(dead)) |driver| onDriverExit(driver);
return;
}
@@ -592,5 +559,6 @@ pub fn main(init: process.Init) void {
.init = initialise,
.on_message = onMessage,
.on_notification = onNotification,
.subscribers = Serve.hooks,
});
}
+3 -3
View File
@@ -10,9 +10,9 @@ pub fn build(b: *std.Build) void {
.name = "display",
.root_source_file = b.path("display.zig"),
.imports = &.{
"channel", "display-client", "display-protocol", "driver", "envelope", "input-client",
"ipc", "logging", "memory", "scanout-protocol", "service",
"thread", "time",
"channel", "display-client", "display-protocol", "driver", "envelope", "input-client",
"ipc", "logging", "memory", "process", "scanout-protocol",
"service", "thread", "time",
},
.threaded = true, // real atomics/TLS (docs/threading.md)
});
+70 -19
View File
@@ -19,6 +19,7 @@ const std = @import("std");
const channel = @import("channel");
const ipc = @import("ipc");
const input = @import("input-client");
const process = @import("process");
const Thread = @import("thread").Thread;
const service = @import("service");
const time = @import("time");
@@ -74,8 +75,22 @@ var background: u32 = 0;
/// small copies rather than one huge bounding box (see compositor.DamageList).
const maximum_layers = 16;
/// A layer belongs to the client that created it. `owner` is the kernel-stamped
/// badge of that task, and `service_owned` (0, an id no task wears) marks the
/// compositor's own layers — the cursor sprite and the self-check pair — which
/// this file creates by direct call rather than over the protocol.
///
/// Layer ids are slots in a sixteen-entry table: small, dense, and guessable, so
/// before this field any client could configure, draw into, or destroy any
/// other's layer — including the cursor. The check lives in the protocol handlers
/// (docs/os-development/protocol-namespace.md: handles validated against the
/// badge); the internal helpers stay unscoped precisely so the compositor can
/// still drive its own.
const service_owned: u32 = 0;
const Layer = struct {
used: bool = false,
owner: u32 = service_owned,
x: i32 = 0,
y: i32 = 0,
z: u32 = 0,
@@ -189,7 +204,8 @@ fn layerAt(id: u32) ?*Layer {
return &layers[id];
}
fn createLayer(x: i32, y: i32, w: u32, h: u32, z: u32, visible: bool) ?u32 {
/// Create a layer for `owner` — `service_owned` for the compositor's own.
fn createLayer(owner: u32, x: i32, y: i32, w: u32, h: u32, z: u32, visible: bool) ?u32 {
if (w == 0 or h == 0) return null;
const slot = freeLayer() orelse return null;
const len = @as(usize, w) * h * 4;
@@ -197,6 +213,7 @@ fn createLayer(x: i32, y: i32, w: u32, h: u32, z: u32, visible: bool) ?u32 {
if (memory.mmapFailed(base)) return null;
layers[slot] = .{
.used = true,
.owner = owner,
.x = x,
.y = y,
.z = z,
@@ -420,8 +437,8 @@ fn selfCheck() void {
const format = backend.info().format;
const red = display_protocol.pack(format, 0xC0, 0x20, 0x20);
const green = display_protocol.pack(format, 0x20, 0xC0, 0x20);
const bottom = createLayer(100, 100, 80, 80, 0, true) orelse return fail_check("create");
const top = createLayer(140, 140, 80, 80, 1, true) orelse return fail_check("create");
const bottom = createLayer(service_owned, 100, 100, 80, 80, 0, true) orelse return fail_check("create");
const top = createLayer(service_owned, 140, 140, 80, 80, 1, true) orelse return fail_check("create");
_ = fillLayer(bottom, Rect.init(0, 0, 80, 80), red);
_ = fillLayer(top, Rect.init(0, 0, 80, 80), green);
present();
@@ -586,7 +603,7 @@ fn startCursorTracking() void {
const mode = backend.info();
cursor_origin_x = @divTrunc(@as(i32, @intCast(mode.width)), 2);
cursor_origin_y = @divTrunc(@as(i32, @intCast(mode.height)), 2);
const id = createLayer(cursor_origin_x, cursor_origin_y, cursor_size, cursor_size, cursor_z, true) orelse {
const id = createLayer(service_owned, cursor_origin_x, cursor_origin_y, cursor_size, cursor_size, cursor_z, true) orelse {
_ = logging.write("display: could not create cursor layer\n");
return;
};
@@ -603,6 +620,9 @@ fn startCursorTracking() void {
fn initialise(endpoint: ipc.Handle) bool {
service_endpoint = endpoint;
// Layers are per-client state, so the compositor needs deaths: a client that
// crashes leaves its surfaces on screen and its slots spent otherwise.
_ = process.subscribeExits(endpoint);
// Pick the scanout backend (GOP today). It logs the reason on failure.
backend = backend_mod.select() orelse return false;
@@ -635,9 +655,35 @@ fn initialise(endpoint: ipc.Handle) bool {
// id a u32: a value that does not fit is not a layer of ours, and `layerAt`
// refuses it the same way an out-of-range one is refused.
fn targetLayer(target: u64) ?u32 {
if (target > std.math.maxInt(u32)) return null;
return @intCast(target);
/// The layer a packet addresses, **for the task that sent it**: null unless the
/// target names a used slot this sender created. A layer that is somebody else's
/// is refused exactly as one that never existed, so a client cannot use the
/// refusal to learn which ids are live (P3's refusal-equals-absence, applied to
/// ids rather than names).
fn targetLayer(target: u64, sender: u32) ?u32 {
if (target > std.math.maxInt(u32)) return null; // a layer id is a u32
const id: u32 = @intCast(target);
const layer = layerAt(id) orelse return null;
if (layer.owner != sender) return null;
return id;
}
/// Destroy every layer a dead client left behind — its surface is pages nobody
/// will ever draw into again, and its slot is one of sixteen. The published
/// exit events are the notice, the same sweep idiom the FAT server uses for open
/// files and the harness uses for subscribers.
fn releaseLayersOf(dead: u32) void {
var released: u32 = 0;
for (&layers, 0..) |*layer, id| {
if (layer.used and layer.owner == dead) {
_ = destroyLayer(@intCast(id));
released += 1;
}
}
if (released != 0) {
std.log.info("released {d} layer(s) for dead client {d}", .{ released, dead });
schedulePresent(); // the screen still shows what they painted
}
}
fn onInfo(_: void, _: Invocation(void), answer: Answer(display_protocol.Info)) isize {
@@ -648,37 +694,37 @@ fn onInfo(_: void, _: Invocation(void), answer: Answer(display_protocol.Info)) i
fn onCreateLayer(_: void, invocation: Invocation(display_protocol.CreateLayer), answer: Answer(display_protocol.Created)) isize {
const request = invocation.request;
const slot = createLayer(request.x, request.y, request.width, request.height, request.z, request.visible != 0) orelse return refused;
const slot = createLayer(invocation.sender, request.x, request.y, request.width, request.height, request.z, request.visible != 0) orelse return refused;
answer.set(.{ .layer = slot });
return 0;
}
fn onConfigureLayer(_: void, invocation: Invocation(display_protocol.ConfigureLayer), _: Answer(void)) isize {
const id = targetLayer(invocation.target) orelse return refused;
const id = targetLayer(invocation.target, invocation.sender) orelse return refused;
const request = invocation.request;
return if (configureLayer(id, request.x, request.y, request.z, request.visible != 0)) 0 else refused;
}
fn onDestroyLayer(_: void, invocation: Invocation(void), _: Answer(void)) isize {
const id = targetLayer(invocation.target) orelse return refused;
const id = targetLayer(invocation.target, invocation.sender) orelse return refused;
return if (destroyLayer(id)) 0 else refused;
}
fn onFillRect(_: void, invocation: Invocation(display_protocol.FillRect), _: Answer(void)) isize {
const id = targetLayer(invocation.target) orelse return refused;
const id = targetLayer(invocation.target, invocation.sender) orelse return refused;
const request = invocation.request;
const local = Rect.init(request.x, request.y, @intCast(request.width), @intCast(request.height));
return if (fillLayer(id, local, request.colour)) 0 else refused;
}
fn onBlitTile(_: void, invocation: Invocation(display_protocol.BlitTile), _: Answer(void)) isize {
const id = targetLayer(invocation.target) orelse return refused;
const id = targetLayer(invocation.target, invocation.sender) orelse return refused;
const request = invocation.request;
return if (blitLayer(id, request.x, request.y, request.width, request.height, invocation.tail)) 0 else refused;
}
fn onDamage(_: void, invocation: Invocation(display_protocol.Damage), _: Answer(void)) isize {
const id = targetLayer(invocation.target) orelse return refused;
const id = targetLayer(invocation.target, invocation.sender) orelse return refused;
const l = layerAt(id) orelse return refused;
const request = invocation.request;
const screen = Rect{ .x = l.x + request.x, .y = l.y + request.y, .w = @intCast(request.width), .h = @intCast(request.height) };
@@ -733,13 +779,18 @@ fn onMessage(message: []const u8, reply: []u8, sender: u32, arrived: *ipc.Arriva
return Serve.dispatch({}, handlers, message, sender, arrived.peek(), reply);
}
/// Two notification sources reach the compositor, and one coalesced badge can carry
/// both, so each bit is handled independently. A **message-notification** is a poke from
/// the mouse-listener thread (a buffered self-`ipc.send`, `notify_message_bit`): fold the
/// newest cursor position into the scene. A **timer** (`notify_timer_bit`) is the frame
/// clock — or the deferred first native present after `attach_scanout` — either way,
/// present the accumulated damage.
/// Three notification sources reach the compositor, and one coalesced badge can carry
/// more than one, so each bit is handled independently. A **message-notification** is a
/// poke from the mouse-listener thread (a buffered self-`ipc.send`, `notify_message_bit`):
/// fold the newest cursor position into the scene. A **timer** (`notify_timer_bit`) is the
/// frame clock — or the deferred first native present after `attach_scanout` — either way,
/// present the accumulated damage. A **published exit** (`notify_exit_bit`) is a client
/// gone: release the layers it left.
fn onNotification(badge: u64) void {
if (badge & ipc.notify_exit_bit != 0) {
releaseLayersOf(@intCast(badge & ~(ipc.notify_badge_bit | ipc.notify_exit_bit)));
return;
}
if (badge & ipc.notify_message_bit != 0) renderCursor();
if (badge & ipc.notify_timer_bit != 0) frameTick();
}
+33 -12
View File
@@ -66,8 +66,9 @@ var ipc_block: IpcBlock = undefined;
var device_dirty: bool = false;
var filesystem: engine.FileSystem = undefined;
// Open handles the VFS holds against this backend: each maps a node id to a
// resolved engine node.
// Open handles clients hold against this backend: each maps a node id to a
// resolved engine node, and to the client that opened it. `owner` is the
// kernel-stamped badge of the opening task — the only source identity there is.
const OpenNode = struct { used: bool = false, node: engine.Node = undefined, owner: u32 = 0 };
var open_nodes = [_]OpenNode{.{}} ** 32;
@@ -78,16 +79,31 @@ fn allocOpen() ?usize {
return null;
}
fn openAt(id: u64) ?*OpenNode {
/// The open node `id` names **for `owner`** — null unless the id is in range, in
/// use, and this client's own. Node ids are small integers drawn from a table of
/// thirty-two, so they are trivially guessable; before this check every client
/// honoured every other client's ids, which is the hole
/// docs/os-development/protocol-namespace.md names ("handles must be scoped per
/// client — validated against the badge"). Nothing else about them changed: they
/// are still per-session, still swept when their owner dies.
///
/// The owner is a *task*, not a process, because the badge is: a threaded client
/// reads and writes a node from the thread that opened it, exactly as the exit
/// sweep already released a worker thread's handles when that thread died.
fn openFor(id: u64, owner: u32) ?*OpenNode {
if (id >= open_nodes.len) return null;
const o = &open_nodes[@intCast(id)];
return if (o.used) o else null;
if (!o.used or o.owner != owner) return null;
return o;
}
/// What a handler returns when the thing asked for is not there — a bad node id,
/// a path that does not resolve, a mutation the volume refused. One errno for all
/// of them, because a filesystem's failures are all "no such thing" as far as the
/// file API can act on them.
/// a node that is someone else's, a path that does not resolve, a mutation the
/// volume refused. One errno for all of them, because a filesystem's failures are
/// all "no such thing" as far as the file API can act on them — and because
/// *someone else's* must be indistinguishable from *nobody's*, or the refusal
/// would itself tell a prober which ids are live (the same discipline the
/// protocol namespace's refused open follows).
const refused: isize = -envelope.ENOENT;
/// How often to look for a block device while none is mounted. Storage arriving
@@ -230,14 +246,14 @@ fn onOpen(_: void, invocation: Invocation(vfs_protocol.Open), answer: Answer(vfs
}
fn onRead(_: void, invocation: Invocation(vfs_protocol.Read), answer: Answer(void)) isize {
const o = openAt(invocation.target) orelse return refused;
const o = openFor(invocation.target, invocation.sender) orelse return refused;
const into = answer.tail();
const want = @min(@as(usize, invocation.request.len), into.len);
return @intCast(filesystem.readFile(o.node, @intCast(invocation.request.offset), into[0..want]));
}
fn onWrite(_: void, invocation: Invocation(vfs_protocol.Write), answer: Answer(vfs_protocol.Written)) isize {
const o = openAt(invocation.target) orelse return refused;
const o = openFor(invocation.target, invocation.sender) orelse return refused;
const data = invocation.tail[0..@min(invocation.tail.len, invocation.request.len)];
const n = filesystem.writeFile(&o.node, @intCast(invocation.request.offset), data);
answer.set(.{ .count = @intCast(n) });
@@ -245,7 +261,7 @@ fn onWrite(_: void, invocation: Invocation(vfs_protocol.Write), answer: Answer(v
}
fn onStatus(_: void, invocation: Invocation(void), answer: Answer(vfs_protocol.FileStatus)) isize {
const o = openAt(invocation.target) orelse return refused;
const o = openFor(invocation.target, invocation.sender) orelse return refused;
const kind: vfs_protocol.NodeKind = if (o.node.is_directory) .directory else .regular;
answer.set(.{ .size = o.node.size, .kind = @intFromEnum(kind), .mtime = o.node.mtime });
return 0;
@@ -255,7 +271,7 @@ fn onStatus(_: void, invocation: Invocation(void), answer: Answer(vfs_protocol.F
/// cursor past the last child — is an entry with no name, which is how the
/// protocol spells it now that the reply's length always counts the fixed part.
fn onReaddir(_: void, invocation: Invocation(vfs_protocol.Readdir), answer: Answer(vfs_protocol.DirectoryEntry)) isize {
const o = openAt(invocation.target) orelse return refused;
const o = openFor(invocation.target, invocation.sender) orelse return refused;
if (!o.node.is_directory) {
answer.set(.{});
return 0;
@@ -272,8 +288,13 @@ fn onReaddir(_: void, invocation: Invocation(vfs_protocol.Readdir), answer: Answ
return @intCast(name_len);
}
/// Closing is an operation on a node like any other, so it is scoped like any
/// other: a client may release its own handles and nobody else's. An id that is
/// not the caller's — free, out of range, or another client's — is refused
/// identically, so a close cannot be used to ask which ids are live either.
fn onClose(_: void, invocation: Invocation(void), _: Answer(void)) isize {
if (openAt(invocation.target)) |o| o.used = false;
const o = openFor(invocation.target, invocation.sender) orelse return refused;
o.used = false;
// Durable-on-close: if any block reached the device since the last flush,
// commit its cache to stable media now (best-effort). This is what makes
// init's shutdown log flush survive a real power-off, and is the right
+1 -1
View File
@@ -9,7 +9,7 @@ pub fn build(b: *std.Build) void {
const exe = build_support.userBinary(b, .{
.name = "input",
.root_source_file = b.path("input.zig"),
.imports = &.{ "envelope", "input-protocol", "ipc", "logging", "process", "service" },
.imports = &.{ "envelope", "input-protocol", "ipc", "logging", "service" },
});
b.installArtifact(exe);
}
+25 -105
View File
@@ -20,133 +20,52 @@
//! handle and `ipc.send`s each event to it. That is the envelope's reserved `subscribe`
//! verb, which this protocol adopts rather than defining its own.
//!
//! P4a moved this service onto the shared harness (library/kernel/service.zig). It was the
//! P4a moved this service onto the shared harness (library/kernel/service.zig) — it was the
//! last hand-rolled receive loop in the tree, and the one service that answered neither the
//! universal ping nor a `terminate` signal — so a shutdown had to kill it. The subscriber
//! table, the fan-out, and the prune-on-subscribe below are unchanged; lifting *those* into
//! the harness is a later milestone, and doing it here would have hidden this one.
//! universal ping nor a `terminate` signal. P4c finished the job: the subscriber table, the
//! fan-out, and the dead-subscriber sweep are the harness's now
//! (`service.Subscribers`), so what is left here is what is actually about input — which
//! device class an event belongs to, and which classes a subscriber asked for. The sweep
//! that replaced the old prune is the one idiom the system uses everywhere: published
//! process-exit notifications, not a poll of the process list on every subscribe.
const envelope = @import("envelope");
const ipc = @import("ipc");
const process = @import("process");
const service = @import("service");
const logging = @import("logging");
const input_protocol = @import("input-protocol");
const envelope = @import("envelope");
/// The generated input dispatch. One fan-out point per process, so the handler
/// context is empty and the subscriber table stays in this file's globals.
const Serve = input_protocol.Protocol.Provider(void);
/// The subscriber side of the input contract: the table, the reserved `subscribe`
/// verb, the exit sweep, and the fan-out. One fan-out point per process, so the
/// handler context is empty.
const Subscriptions = service.Subscribers(input_protocol.Protocol, void);
const Invocation = envelope.Invocation;
const Answer = envelope.Answer;
/// One registered subscriber: the endpoint we push events to (a capability it handed us at
/// subscribe time) and the task id that owns it (the subscribe call's badge), so a slot
/// left behind by a subscriber that exited can be reclaimed.
const Subscriber = struct {
used: bool = false,
endpoint: ipc.Handle = 0,
task_id: u32 = 0,
/// Which device classes this subscriber wants (an OR of input_protocol.device_*). An event
/// is delivered only if its device's bit is set here.
device_mask: u32 = 0,
};
var subscribers = [_]Subscriber{.{}} ** 8;
/// Drop any subscriber whose owning process is no longer alive, so its slot (and the
/// endpoint reference it holds) can be reused. Cheap and only run on subscribe — the async
/// `send` to a dead subscriber's orphaned endpoint is harmless (it just fills a queue no
/// one drains), so this is housekeeping, not correctness.
fn pruneDeadSubscribers() void {
var table: [32]process.ProcessDescriptor = undefined;
const total = process.processes(&table);
const count = @min(total, table.len);
for (&subscribers) |*sub| {
if (!sub.used) continue;
var alive = false;
for (table[0..count]) |descriptor| {
if (descriptor.id == sub.task_id) {
alive = true;
break;
}
}
// The slot owns the endpoint capability it was handed, so reclaiming the
// slot closes it — otherwise a process that subscribes and dies costs a
// handle-table slot that never comes back.
if (!alive) {
_ = ipc.close(sub.endpoint);
sub.* = .{};
}
}
}
/// Register `endpoint` (owned by task `task_id`) to receive the device classes in
/// `device_mask`. Returns false if the subscriber table is full.
fn addSubscriber(endpoint: ipc.Handle, task_id: u32, device_mask: u32) bool {
for (&subscribers) |*sub| {
if (!sub.used) {
sub.* = .{ .used = true, .endpoint = endpoint, .task_id = task_id, .device_mask = device_mask };
return true;
}
}
return false;
}
/// Push `event` to every subscriber whose interest mask includes its device class.
/// `ipc.send` never blocks, so a slow or dead subscriber cannot stall delivery to others.
///
/// The class is the packet's operation, so the fan-out picks the event by device and the
/// packet is framed once, outside the loop — every subscriber of a class gets identical
/// bytes, which is what "one fan-out point per event domain" means on the wire.
/// Push `event` to every subscriber whose interest mask includes its device class. The
/// class is the packet's operation, so this picks *which event* to publish and the harness
/// frames it once for the whole fan-out — the mapping from device to class is the only part
/// of a broadcast that is this service's own.
fn broadcast(event: input_protocol.InputEvent) void {
const class = input_protocol.eventOfDevice(event.device) orelse return; // no class wants it
var packet: [envelope.post_maximum]u8 = undefined;
const framed = switch (class) {
.keyboard => input_protocol.Protocol.encodeEvent(.keyboard, 0, event.asKeyboard() orelse return, &packet),
.mouse => input_protocol.Protocol.encodeEvent(.mouse, 0, event.asMouse() orelse return, &packet),
.joystick => input_protocol.Protocol.encodeEvent(.joystick, 0, event.asJoystick() orelse return, &packet),
} orelse return;
const bit = input_protocol.deviceBit(event.device);
for (&subscribers) |*sub| {
if (sub.used and sub.device_mask & bit != 0) _ = ipc.send(sub.endpoint, framed);
const class = input_protocol.deviceBit(event.device);
switch (input_protocol.eventOfDevice(event.device) orelse return) { // no class wants it
.keyboard => Subscriptions.publishClass(.keyboard, 0, event.asKeyboard() orelse return, class),
.mouse => Subscriptions.publishClass(.mouse, 0, event.asMouse() orelse return, class),
.joystick => Subscriptions.publishClass(.joystick, 0, event.asJoystick() orelse return, class),
}
}
/// Set by `onSubscribe` when the subscriber table has taken ownership of the capability the
/// call carried, and read by `onMessage`, 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 capability_claimed = false;
/// The reserved `subscribe` verb: register the caller's endpoint (the call's capability) for
/// the classes in the packet's tail. Refusals simply return, and the turn closes what arrived
/// — the ownership rule the harness states (`ipc.Arrival`), unchanged by the move onto it.
fn onSubscribe(_: void, invocation: Invocation(void), _: Answer(void)) isize {
const endpoint = invocation.capability orelse return -envelope.EPROTO; // no endpoint passed
// A zero mask means "everything" (a subscriber that named no class still wants input).
const requested = input_protocol.decodeSubscribe(invocation.tail).device_mask;
const mask = if (requested == 0) input_protocol.device_all else requested;
pruneDeadSubscribers();
if (!addSubscriber(endpoint, invocation.sender, mask)) return -envelope.ENOSPC; // table full
capability_claimed = true; // the subscriber table holds it until that task dies
return 0;
}
fn onPublish(_: void, invocation: Invocation(input_protocol.InputEvent), _: Answer(void)) isize {
broadcast(invocation.request);
return 0;
}
const handlers = Serve.Handlers{ .publish = onPublish, .subscribe = onSubscribe };
/// `subscribe` and `unsubscribe` are absent on purpose: the harness answers both.
const handlers = Subscriptions.Handlers{ .publish = onPublish };
fn onMessage(message: []const u8, out: []u8, sender: u32, arrived: *ipc.Arrival) usize {
capability_claimed = false;
const written = Serve.dispatch({}, handlers, message, sender, arrived.peek(), out);
if (capability_claimed) _ = arrived.take();
return written;
return Subscriptions.dispatch({}, handlers, message, sender, arrived, out);
}
fn initialise(_: ipc.Handle) bool {
@@ -162,5 +81,6 @@ pub fn main() void {
.service = "input",
.init = initialise,
.on_message = onMessage,
.subscribers = Subscriptions.hooks,
});
}