authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-03-31 02:10:50-07:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2025-10-02 16:30:59-07:00
log0086d315f53a76c028d3072240ebb3a3b9757a04
tree450612544dc74fd29c0dcffa0790aff271d919ed
parent5508b4c8876074b736ed436f0d20d5cec86f1f85

std.Io: add detached async


2 files changed, 113 insertions(+), 4 deletions(-)

lib/std/Io.zig+43-4
......@@ -580,6 +580,18 @@ pub const VTable = struct {
580580 context_alignment: std.mem.Alignment,
581581 start: *const fn (context: *const anyopaque, result: *anyopaque) void,
582582 ) ?*AnyFuture,
583 /// Executes `start` asynchronously in a manner such that it cleans itself
584 /// up. This mode does not support results, await, or cancel.
585 ///
586 /// Thread-safe.
587 go: *const fn (
588 /// Corresponds to `Io.userdata`.
589 userdata: ?*anyopaque,
590 /// Copied and then passed to `start`.
591 context: []const u8,
592 context_alignment: std.mem.Alignment,
593 start: *const fn (context: *const anyopaque) void,
594 ) void,
583595 /// This function is only called when `async` returns a non-null value.
584596 ///
585597 /// Thread-safe.
......@@ -593,7 +605,6 @@ pub const VTable = struct {
593605 result: []u8,
594606 result_alignment: std.mem.Alignment,
595607 ) void,
596
597608 /// Equivalent to `await` but initiates cancel request.
598609 ///
599610 /// This function is only called when `async` returns a non-null value.
......@@ -671,14 +682,24 @@ pub fn Future(Result: type) type {
671682 /// Idempotent.
672683 pub fn cancel(f: *@This(), io: Io) Result {
673684 const any_future = f.any_future orelse return f.result;
674 io.vtable.cancel(io.userdata, any_future, @ptrCast((&f.result)[0..1]), .of(Result));
685 io.vtable.cancel(
686 io.userdata,
687 any_future,
688 if (@sizeOf(Result) == 0) &.{} else @ptrCast((&f.result)[0..1]), // work around compiler bug
689 .of(Result),
690 );
675691 f.any_future = null;
676692 return f.result;
677693 }
678694
679695 pub fn await(f: *@This(), io: Io) Result {
680696 const any_future = f.any_future orelse return f.result;
681 io.vtable.await(io.userdata, any_future, @ptrCast((&f.result)[0..1]), .of(Result));
697 io.vtable.await(
698 io.userdata,
699 any_future,
700 if (@sizeOf(Result) == 0) &.{} else @ptrCast((&f.result)[0..1]), // work around compiler bug
701 .of(Result),
702 );
682703 f.any_future = null;
683704 return f.result;
684705 }
......@@ -996,7 +1017,7 @@ pub fn async(io: Io, function: anytype, args: anytype) Future(@typeInfo(@TypeOf(
9961017 var future: Future(Result) = undefined;
9971018 future.any_future = io.vtable.async(
9981019 io.userdata,
999 @ptrCast((&future.result)[0..1]),
1020 if (@sizeOf(Result) == 0) &.{} else @ptrCast((&future.result)[0..1]), // work around compiler bug
10001021 .of(Result),
10011022 if (@sizeOf(Args) == 0) &.{} else @ptrCast((&args)[0..1]), // work around compiler bug
10021023 .of(Args),
......@@ -1005,6 +1026,24 @@ pub fn async(io: Io, function: anytype, args: anytype) Future(@typeInfo(@TypeOf(
10051026 return future;
10061027}
10071028
1029/// Calls `function` with `args` asynchronously. The resource cleans itself up
1030/// when the function returns. Does not support await, cancel, or a return value.
1031pub fn go(io: Io, function: anytype, args: anytype) void {
1032 const Args = @TypeOf(args);
1033 const TypeErased = struct {
1034 fn start(context: *const anyopaque) void {
1035 const args_casted: *const Args = @alignCast(@ptrCast(context));
1036 @call(.auto, function, args_casted.*);
1037 }
1038 };
1039 io.vtable.go(
1040 io.userdata,
1041 if (@sizeOf(Args) == 0) &.{} else @ptrCast((&args)[0..1]), // work around compiler bug
1042 .of(Args),
1043 TypeErased.start,
1044 );
1045}
1046
10081047pub fn openFile(io: Io, dir: fs.Dir, sub_path: []const u8, flags: fs.File.OpenFlags) FileOpenError!fs.File {
10091048 return io.vtable.openFile(io.userdata, dir, sub_path, flags);
10101049}
lib/std/Thread/Pool.zig+70
......@@ -332,6 +332,7 @@ pub fn io(pool: *Pool) Io {
332332 .vtable = &.{
333333 .@"async" = @"async",
334334 .@"await" = @"await",
335 .go = go,
335336 .cancel = cancel,
336337 .cancelRequested = cancelRequested,
337338 .mutexLock = mutexLock,
......@@ -472,6 +473,75 @@ fn @"async"(
472473 return @ptrCast(closure);
473474}
474475
476const DetachedClosure = struct {
477 pool: *Pool,
478 func: *const fn (context: *anyopaque) void,
479 run_node: std.Thread.Pool.RunQueue.Node = .{ .data = .{ .runFn = runFn } },
480 context_alignment: std.mem.Alignment,
481 context_len: usize,
482
483 fn runFn(runnable: *std.Thread.Pool.Runnable, _: ?usize) void {
484 const run_node: *std.Thread.Pool.RunQueue.Node = @fieldParentPtr("data", runnable);
485 const closure: *DetachedClosure = @alignCast(@fieldParentPtr("run_node", run_node));
486 closure.func(closure.contextPointer());
487 const gpa = closure.pool.allocator;
488 const base: [*]align(@alignOf(DetachedClosure)) u8 = @ptrCast(closure);
489 gpa.free(base[0..contextEnd(closure.context_alignment, closure.context_len)]);
490 }
491
492 fn contextOffset(context_alignment: std.mem.Alignment) usize {
493 return context_alignment.forward(@sizeOf(DetachedClosure));
494 }
495
496 fn contextEnd(context_alignment: std.mem.Alignment, context_len: usize) usize {
497 return contextOffset(context_alignment) + context_len;
498 }
499
500 fn contextPointer(closure: *DetachedClosure) [*]u8 {
501 const base: [*]u8 = @ptrCast(closure);
502 return base + contextOffset(closure.context_alignment);
503 }
504};
505
506fn go(
507 userdata: ?*anyopaque,
508 context: []const u8,
509 context_alignment: std.mem.Alignment,
510 start: *const fn (context: *const anyopaque) void,
511) void {
512 const pool: *std.Thread.Pool = @alignCast(@ptrCast(userdata));
513 pool.mutex.lock();
514
515 const gpa = pool.allocator;
516 const n = DetachedClosure.contextEnd(context_alignment, context.len);
517 const closure: *DetachedClosure = @alignCast(@ptrCast(gpa.alignedAlloc(u8, @alignOf(DetachedClosure), n) catch {
518 pool.mutex.unlock();
519 start(context.ptr);
520 return;
521 }));
522 closure.* = .{
523 .pool = pool,
524 .func = start,
525 .context_alignment = context_alignment,
526 .context_len = context.len,
527 };
528 @memcpy(closure.contextPointer()[0..context.len], context);
529 pool.run_queue.prepend(&closure.run_node);
530
531 if (pool.threads.items.len < pool.threads.capacity) {
532 pool.threads.addOneAssumeCapacity().* = std.Thread.spawn(.{
533 .stack_size = pool.stack_size,
534 .allocator = gpa,
535 }, worker, .{pool}) catch t: {
536 pool.threads.items.len -= 1;
537 break :t undefined;
538 };
539 }
540
541 pool.mutex.unlock();
542 pool.cond.signal();
543}
544
475545fn @"await"(
476546 userdata: ?*anyopaque,
477547 any_future: *std.Io.AnyFuture,