










在严肃阅读了Zig’s Io.Threaded is Neat这篇文章及其引用的文章后,我反思了一下我之前写的任务队列实现,感觉之前的实现太不 Zig 了,完全没体现出 std.Io 的优势。
之前的队列 Consumer 需要用回调函数来通知任务的发布者一个任务完成了,这其实就相当于在写类似 JavaScript 中 promise.then(cb)这样的代码。而且由于 Zig 完全不支持闭包,所以这块的实现不可避免地需要用到指针:
userdata: *anyopaque,
on_result: OnResultCallback,
const OnResultCallback = *const fn (*anyopaque, Task) void;
// in consumer
self.on_result(self.userdata, t.*);如果说 JavaScript 中回调太多而被称为「回调地狱」的话,那么这种无法使用闭包的回调可谓是「回调十八层地狱」了。
JavaScript 解决「回调地狱」问题靠的是引入 async/await 关键字,与之对应的 Zig 魔法就是 .async/.await 函数了。在目前默认的线程模型下,当我们调用 .await 时,本质上是在让当前线程等待其他线程将其唤醒。放到我们的获取任务执行结果的场景下,我们需要让任务结果读取者调用 .await ,然后让队列的 Consumer 在任务执行结束后唤醒结果读取者。与这个需求对应的 Zig API 就是 std.Io.Cond 。
这个类型的方法主要分为两类:
wait(cond: *Condition, io: Io, mutex: *Mutex):让当前正在执行的函数挂起并进入等待状态waitUncancelable(cond: *Condition, io: Io, mutex: *Mutex):跟 Wait 相同,但会让等待变得不可取消signal(cond: *Condition, io: Io):唤醒一个等待者broadcast(cond: *Condition, io: Io):唤醒全部的等待者这个函数如此简单以至于标准库都没有提供函数说明文档。对我这个之前完全没接触过这方面编程知识的人来说,必须借助 LLM 阅读源码才能理解这个函数应该如何使用。
然而目前的 Cond 实现存在一些问题,不过不太影响我的使用,希望 Zig 能早日修复吧。
掌握了 Cond 的用法,实现一个类似 C# 的 TaskCompletionSource 就很简单了:
const std = @import("std");
const Io = std.Io;
pub fn TaskCompletionSource(comptime T: type, comptime E: type) type {
return struct {
const Self = @This();
mutex: Io.Mutex = .init,
state: union(enum) {
pending: void,
completed: T,
faulted: E,
} = .pending,
cond: Io.Condition = .init,
pub const init = Self{};
pub fn setComplete(self: *Self, io: Io, value: T) !void {
try self.mutex.lock(io);
defer self.mutex.unlock(io);
if (self.state != .pending) return error.AlreadyInFinalState;
self.cond.broadcast(io);
self.state = .{ .completed = value };
}
pub fn setFaulted(self: *Self, io: Io, err: E) !void {
try self.mutex.lock(io);
defer self.mutex.unlock(io);
if (self.state != .pending) return error.AlreadyInFinalState;
self.cond.broadcast(io);
self.state = .{ .faulted = err };
}
pub fn await(self: *Self, io: Io) !void {
try self.mutex.lock(io);
defer self.mutex.unlock(io);
while (self.state == .pending) try self.cond.wait(io, &self.mutex);
}
};
}
const TCS = TaskCompletionSource(i32, []const u8);
const Stub = struct {
fn complete(t: *TCS) void {
t.setComplete(std.testing.io, 42) catch unreachable;
}
fn fault(t: *TCS) void {
t.setFaulted(std.testing.io, "error") catch unreachable;
}
};
test "f: tcs" {
var tcs = TCS{};
_ = std.testing.io.async(Stub.complete, .{&tcs});
try tcs.await(std.testing.io);
switch (tcs.state) {
.completed => |value| try std.testing.expect(value == 42),
else => unreachable,
}
}有了这个新的类型后,在任务队列的 Consumer 中唤醒任务结果读取者就非常简单了:
// task result reader, waits for the task queue picks the task and executes it.
try tcs.await(io);
// task queue consumer, when task completed wake up readers.
try self.tcs.setComplete(self.io, tool_use_result);此内容由惯性聚合(RSS阅读器)自动聚合整理,仅供阅读参考。 原文来自 — 版权归原作者所有。