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-07-20 10:38:39-07:00
logebf92042e3a081ce84668a7edc1aa74c9ad7e9e5
tree377ed810f6e21ad79e6ea5dd6380ca0f28775417
parenta7790bd32e1e8caf6f2f0bedede8a7cb7b35c443

std.Io: add detached async


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

lib/std/Io.zig+43-4
...@@ -933,6 +933,18 @@ pub const VTable = struct {...@@ -933,6 +933,18 @@ pub const VTable = struct {
933 context_alignment: std.mem.Alignment,933 context_alignment: std.mem.Alignment,
934 start: *const fn (context: *const anyopaque, result: *anyopaque) void,934 start: *const fn (context: *const anyopaque, result: *anyopaque) void,
935 ) ?*AnyFuture,935 ) ?*AnyFuture,
936 /// Executes `start` asynchronously in a manner such that it cleans itself
937 /// up. This mode does not support results, await, or cancel.
938 ///
939 /// Thread-safe.
940 go: *const fn (
941 /// Corresponds to `Io.userdata`.
942 userdata: ?*anyopaque,
943 /// Copied and then passed to `start`.
944 context: []const u8,
945 context_alignment: std.mem.Alignment,
946 start: *const fn (context: *const anyopaque) void,
947 ) void,
936 /// This function is only called when `async` returns a non-null value.948 /// This function is only called when `async` returns a non-null value.
937 ///949 ///
938 /// Thread-safe.950 /// Thread-safe.
...@@ -946,7 +958,6 @@ pub const VTable = struct {...@@ -946,7 +958,6 @@ pub const VTable = struct {
946 result: []u8,958 result: []u8,
947 result_alignment: std.mem.Alignment,959 result_alignment: std.mem.Alignment,
948 ) void,960 ) void,
949
950 /// Equivalent to `await` but initiates cancel request.961 /// Equivalent to `await` but initiates cancel request.
951 ///962 ///
952 /// This function is only called when `async` returns a non-null value.963 /// This function is only called when `async` returns a non-null value.
...@@ -1024,14 +1035,24 @@ pub fn Future(Result: type) type {...@@ -1024,14 +1035,24 @@ pub fn Future(Result: type) type {
1024 /// Idempotent.1035 /// Idempotent.
1025 pub fn cancel(f: *@This(), io: Io) Result {1036 pub fn cancel(f: *@This(), io: Io) Result {
1026 const any_future = f.any_future orelse return f.result;1037 const any_future = f.any_future orelse return f.result;
1027 io.vtable.cancel(io.userdata, any_future, @ptrCast((&f.result)[0..1]), .of(Result));1038 io.vtable.cancel(
1039 io.userdata,
1040 any_future,
1041 if (@sizeOf(Result) == 0) &.{} else @ptrCast((&f.result)[0..1]), // work around compiler bug
1042 .of(Result),
1043 );
1028 f.any_future = null;1044 f.any_future = null;
1029 return f.result;1045 return f.result;
1030 }1046 }
10311047
1032 pub fn await(f: *@This(), io: Io) Result {1048 pub fn await(f: *@This(), io: Io) Result {
1033 const any_future = f.any_future orelse return f.result;1049 const any_future = f.any_future orelse return f.result;
1034 io.vtable.await(io.userdata, any_future, @ptrCast((&f.result)[0..1]), .of(Result));1050 io.vtable.await(
1051 io.userdata,
1052 any_future,
1053 if (@sizeOf(Result) == 0) &.{} else @ptrCast((&f.result)[0..1]), // work around compiler bug
1054 .of(Result),
1055 );
1035 f.any_future = null;1056 f.any_future = null;
1036 return f.result;1057 return f.result;
1037 }1058 }
...@@ -1349,7 +1370,7 @@ pub fn async(io: Io, function: anytype, args: anytype) Future(@typeInfo(@TypeOf(...@@ -1349,7 +1370,7 @@ pub fn async(io: Io, function: anytype, args: anytype) Future(@typeInfo(@TypeOf(
1349 var future: Future(Result) = undefined;1370 var future: Future(Result) = undefined;
1350 future.any_future = io.vtable.async(1371 future.any_future = io.vtable.async(
1351 io.userdata,1372 io.userdata,
1352 @ptrCast((&future.result)[0..1]),1373 if (@sizeOf(Result) == 0) &.{} else @ptrCast((&future.result)[0..1]), // work around compiler bug
1353 .of(Result),1374 .of(Result),
1354 if (@sizeOf(Args) == 0) &.{} else @ptrCast((&args)[0..1]), // work around compiler bug1375 if (@sizeOf(Args) == 0) &.{} else @ptrCast((&args)[0..1]), // work around compiler bug
1355 .of(Args),1376 .of(Args),
...@@ -1358,6 +1379,24 @@ pub fn async(io: Io, function: anytype, args: anytype) Future(@typeInfo(@TypeOf(...@@ -1358,6 +1379,24 @@ pub fn async(io: Io, function: anytype, args: anytype) Future(@typeInfo(@TypeOf(
1358 return future;1379 return future;
1359}1380}
13601381
1382/// Calls `function` with `args` asynchronously. The resource cleans itself up
1383/// when the function returns. Does not support await, cancel, or a return value.
1384pub fn go(io: Io, function: anytype, args: anytype) void {
1385 const Args = @TypeOf(args);
1386 const TypeErased = struct {
1387 fn start(context: *const anyopaque) void {
1388 const args_casted: *const Args = @alignCast(@ptrCast(context));
1389 @call(.auto, function, args_casted.*);
1390 }
1391 };
1392 io.vtable.go(
1393 io.userdata,
1394 if (@sizeOf(Args) == 0) &.{} else @ptrCast((&args)[0..1]), // work around compiler bug
1395 .of(Args),
1396 TypeErased.start,
1397 );
1398}
1399
1361pub fn openFile(io: Io, dir: fs.Dir, sub_path: []const u8, flags: fs.File.OpenFlags) FileOpenError!fs.File {1400pub fn openFile(io: Io, dir: fs.Dir, sub_path: []const u8, flags: fs.File.OpenFlags) FileOpenError!fs.File {
1362 return io.vtable.openFile(io.userdata, dir, sub_path, flags);1401 return io.vtable.openFile(io.userdata, dir, sub_path, flags);
1363}1402}
lib/std/Thread/Pool.zig+70
...@@ -332,6 +332,7 @@ pub fn io(pool: *Pool) Io {...@@ -332,6 +332,7 @@ pub fn io(pool: *Pool) Io {
332 .vtable = &.{332 .vtable = &.{
333 .@"async" = @"async",333 .@"async" = @"async",
334 .@"await" = @"await",334 .@"await" = @"await",
335 .go = go,
335 .cancel = cancel,336 .cancel = cancel,
336 .cancelRequested = cancelRequested,337 .cancelRequested = cancelRequested,
337 .mutexLock = mutexLock,338 .mutexLock = mutexLock,
...@@ -472,6 +473,75 @@ fn @"async"(...@@ -472,6 +473,75 @@ fn @"async"(
472 return @ptrCast(closure);473 return @ptrCast(closure);
473}474}
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
475fn @"await"(545fn @"await"(
476 userdata: ?*anyopaque,546 userdata: ?*anyopaque,
477 any_future: *std.Io.AnyFuture,547 any_future: *std.Io.AnyFuture,