Fan-out
Description
Fan-out.
- Main sends 5 Events and 4 Sensors into one mailbox.
- 3 worker threads share the mailbox, compete for items.
- Main closes the mailbox, frees any items left unclaimed.
- Verifies every item was either received or freed.
Diagram
main ──Event×5 + Sensor×4──► mailbox ──► worker A
├──► worker B (compete; each item goes to one)
└──► worker C
mailbox.close ──► remaining list ──► freeItem (main)
Source
pub fn fan_out(allocator: std.mem.Allocator, io: std.Io) !void {
const mbh: MailboxHandle = try mailbox.new(io, allocator);
defer mailbox.destroy(mbh, allocator);
var ctx_a: WorkerCtx = .{ .mbh = mbh, .alloc = allocator };
var ctx_b: WorkerCtx = .{ .mbh = mbh, .alloc = allocator };
var ctx_c: WorkerCtx = .{ .mbh = mbh, .alloc = allocator };
var fa = try io.concurrent(fanOutWorkerFn, .{&ctx_a});
var fb = try io.concurrent(fanOutWorkerFn, .{&ctx_b});
var fc = try io.concurrent(fanOutWorkerFn, .{&ctx_c});
const n_events: usize = 5;
const n_sensors: usize = 4;
var i: usize = 0;
while (i < n_events) : (i += 1) {
var slot: Slot = null;
defer items.Event.EventPolyHelper.destroy(allocator, &slot);
try items.Event.EventPolyHelper.create(allocator, &slot);
items.Event.EventPolyHelper.mustIdentifySlotAs(&slot).code = @intCast(i);
try mailbox.send(mbh, &slot);
}
i = 0;
while (i < n_sensors) : (i += 1) {
var slot: Slot = null;
defer items.Sensor.SensorPolyHelper.destroy(allocator, &slot);
try items.Sensor.SensorPolyHelper.create(allocator, &slot);
items.Sensor.SensorPolyHelper.mustIdentifySlotAs(&slot).value = @as(f64, @floatFromInt(i));
try mailbox.send(mbh, &slot);
}
var rem: std.DoublyLinkedList = mailbox.close(mbh);
var remaining: usize = 0;
while (rem.popFirst()) |node| {
items.freeItem(@fieldParentPtr("node", node), allocator);
remaining += 1;
}
fa.await(io);
fb.await(io);
fc.await(io);
const total: usize = ctx_a.received + ctx_b.received + ctx_c.received;
std.log.info("fan-out: a={d} b={d} c={d} remaining={d}", .{ ctx_a.received, ctx_b.received, ctx_c.received, remaining });
try helpers.expect(error.FanOutFailed, total + remaining == n_events + n_sensors, "wrong total");
}
const WorkerCtx = struct {
mbh: MailboxHandle,
alloc: std.mem.Allocator,
received: usize = 0,
};
fn fanOutWorkerFn(ctx: *WorkerCtx) void {
while (true) {
var slot: Slot = null;
defer items.freeSlot(&slot, ctx.alloc);
mailbox.receive(ctx.mbh, &slot, null) catch return;
ctx.received += 1;
}
}
const items = @import("../items/items.zig");
const helpers = @import("../helpers/helpers.zig");
const matryoshka = @import("matryoshka");
const std = @import("std");
const polynode = matryoshka.polynode;
const mailbox = matryoshka.mailbox;
const Slot = polynode.Slot;
const MailboxHandle = mailbox.MailboxHandle;