Add input module: broadcast keyboard events over IPC
Programs can now subscribe to keyboard events (key_down/key_up/key_press) and drivers can broadcast them, through a new user-space input service. The delivery model is forced by danos IPC: a synchronous rendezvous holds one pending reply, so a server cannot park N subscribers blocked in a "wait for next event" call — delivery must be push. But a synchronous push has no timeout and the kernel never wakes a sender parked on a dead peer's endpoint, so one dying subscriber would hang all input. So this lands the roadmap's planned asynchronous buffered send and builds the service on it: - ipc_send (syscall 26): non-blocking post to an endpoint's bounded payload ring, delivered through reply_wait as a buffered message (notify_message_bit). A full ring drops the oldest. It can never hang on a dead/slow peer. - input-protocol + runtime.input helpers (subscribe/next, connectSource/ publish) — the first real consumer of M13 capability passing: a subscriber hands the service its own endpoint as a capability. - input service (fan-out via ipc_send, dead-subscriber pruning), a synthetic input-source, and input-test; the ps2-bus keyboard driver publishes to it. Real IRQ1 scancode decoding (which must live in the bus, the PNP0303 owner) is a documented follow-up; the source is synthetic for now. - build/init wiring, an `input` QEMU case, and docs/input.md. Full QEMU suite 48/48, including the new input case and every IPC/endpoint regression (ipc, ipc-call, ipc-cap, vfs, hpet, bus, irqfree).
This commit is contained in:
@@ -0,0 +1,138 @@
|
||||
//! system/services/input — the user-space input service. Shipped in the initial_ramdisk,
|
||||
//! spawned as a ring-3 process, and published under the well-known `input` service id. It
|
||||
//! is the fan-out point between **sources** (keyboard drivers) and **subscribers** (any
|
||||
//! program that wants keyboard events): a source `publish`es a `KeyEvent`, and the service
|
||||
//! pushes it to every subscriber.
|
||||
//!
|
||||
//! The delivery discipline is the whole design (see docs/input.md). The kernel's IPC is a
|
||||
//! synchronous rendezvous: a server holds one pending reply, so it cannot park N
|
||||
//! subscribers blocked in a "wait for next event" call. Broadcasting therefore has to be
|
||||
//! *push* — the service delivering to subscribers. But a synchronous push (`ipc_call`)
|
||||
//! would let one dead or wedged subscriber hang the whole broadcast, since the kernel
|
||||
//! never wakes a sender parked on a dead peer's endpoint. So delivery uses the
|
||||
//! asynchronous `ipc.send`: it posts the event to each subscriber's endpoint queue and
|
||||
//! returns at once, and can never block on a subscriber. That primitive exists for exactly
|
||||
//! this ([ipc.md](../../../docs/ipc.md), "asynchronous / buffered send").
|
||||
//!
|
||||
//! A subscriber registers by handing the service its own endpoint as a capability (M13
|
||||
//! capability passing — this service is its first real consumer). The service keeps that
|
||||
//! handle and `ipc.send`s each event to it.
|
||||
|
||||
const std = @import("std");
|
||||
const runtime = @import("runtime");
|
||||
const protocol = runtime.input_protocol;
|
||||
const ipc = runtime.ipc;
|
||||
const system = runtime.system;
|
||||
|
||||
/// 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,
|
||||
};
|
||||
|
||||
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]system.ProcessDescriptor = undefined;
|
||||
const total = system.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;
|
||||
}
|
||||
}
|
||||
if (!alive) sub.* = .{};
|
||||
}
|
||||
}
|
||||
|
||||
/// Register `endpoint` (owned by task `task_id`) to receive events. Returns false if the
|
||||
/// subscriber table is full.
|
||||
fn addSubscriber(endpoint: ipc.Handle, task_id: u32) bool {
|
||||
for (&subscribers) |*sub| {
|
||||
if (!sub.used) {
|
||||
sub.* = .{ .used = true, .endpoint = endpoint, .task_id = task_id };
|
||||
return true;
|
||||
}
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/// Push `event` to every registered subscriber. `ipc.send` never blocks, so a slow or
|
||||
/// dead subscriber cannot stall delivery to the others.
|
||||
fn broadcast(event: protocol.KeyEvent) void {
|
||||
const bytes = std.mem.asBytes(&event);
|
||||
for (&subscribers) |*sub| {
|
||||
if (sub.used) _ = ipc.send(sub.endpoint, bytes);
|
||||
}
|
||||
}
|
||||
|
||||
/// Handle one request. `got` carries the sender badge (a task id) and, for subscribe, the
|
||||
/// subscriber's endpoint capability in `got.cap`. Writes a `Reply` into `out` and returns
|
||||
/// its length.
|
||||
fn handle(message: []const u8, got: ipc.Received, out: []u8) usize {
|
||||
const reply = struct {
|
||||
fn write(buffer: []u8, status: i32) usize {
|
||||
const header = protocol.Reply{ .status = status };
|
||||
@memcpy(buffer[0..protocol.reply_size], std.mem.asBytes(&header));
|
||||
return protocol.reply_size;
|
||||
}
|
||||
};
|
||||
|
||||
if (message.len < protocol.request_size) return reply.write(out, -1);
|
||||
const request = std.mem.bytesToValue(protocol.Request, message[0..protocol.request_size]);
|
||||
|
||||
switch (@as(protocol.Operation, @enumFromInt(request.operation))) {
|
||||
.subscribe => {
|
||||
const endpoint = got.cap orelse return reply.write(out, -1); // no endpoint passed
|
||||
pruneDeadSubscribers();
|
||||
if (!addSubscriber(endpoint, @intCast(got.badge))) return reply.write(out, -1); // table full
|
||||
return reply.write(out, 0);
|
||||
},
|
||||
.publish => {
|
||||
broadcast(request.event);
|
||||
return reply.write(out, 0);
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
pub fn main() void {
|
||||
const endpoint = ipc.createIpcEndpoint() orelse {
|
||||
_ = system.write("input: no endpoint\n");
|
||||
return;
|
||||
};
|
||||
if (!ipc.register(.input, endpoint)) {
|
||||
_ = system.write("input: register failed\n");
|
||||
return;
|
||||
}
|
||||
_ = system.write("input: ready\n");
|
||||
|
||||
var reply_buffer: [protocol.reply_size]u8 = undefined;
|
||||
var reply_len: usize = 0;
|
||||
var receive: [protocol.request_size]u8 = undefined;
|
||||
while (true) {
|
||||
const got = ipc.replyWait(endpoint, reply_buffer[0..reply_len], &receive, null);
|
||||
// Only synchronous client requests (subscribe/publish) arrive here; nothing sends
|
||||
// this service asynchronous messages, so a notification wake would be spurious.
|
||||
if (got.isNotification()) {
|
||||
reply_len = 0;
|
||||
continue;
|
||||
}
|
||||
reply_len = handle(receive[0..got.len], got, &reply_buffer);
|
||||
}
|
||||
}
|
||||
|
||||
pub const panic = runtime.panic;
|
||||
comptime {
|
||||
_ = &runtime.start._start;
|
||||
}
|
||||
Reference in New Issue
Block a user