authorgravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2019-11-08 01:21:22-05:00
committergravatar for andrew@ziglang.orgAndrew Kelley <andrew@ziglang.org> 2019-11-08 01:21:22-05:00
logfbbcf2f30d3fe0a9b0c41de9b737c13737497a3b
treed214c16b689f6002c2e1de3741e5fcec83fd374c
parent9b0536e6f43ba916b2b488377f1e87d0ecf4ccf9
parente2189b6e5d5a5644eb937b682cdfe4e658fe27e8
signaturelock-open Commit is signed but in an unrecognized format.

Merge branch 'adaptive_lock' of https://github.com/kprotty/zig into kprotty-adaptive_lock


10 files changed, 462 insertions(+), 184 deletions(-)

lib/std/c.zig+2
......@@ -203,3 +203,5 @@ pub extern "c" fn dn_expand(
203203 exp_dn: [*]u8,
204204 length: c_int,
205205) c_int;
206
207pub extern "c" fn sched_yield() c_int;
lib/std/mutex.zig+64-74
......@@ -1,19 +1,13 @@
11const std = @import("std.zig");
22const builtin = @import("builtin");
3const AtomicOrder = builtin.AtomicOrder;
4const AtomicRmwOp = builtin.AtomicRmwOp;
53const testing = std.testing;
64const SpinLock = std.SpinLock;
7const linux = std.os.linux;
8const windows = std.os.windows;
5const ThreadParker = std.ThreadParker;
96
107/// Lock may be held only once. If the same thread
118/// tries to acquire the same mutex twice, it deadlocks.
12/// This type must be initialized at runtime, and then deinitialized when no
13/// longer needed, to free resources.
14/// If you need static initialization, use std.StaticallyInitializedMutex.
15/// The Linux implementation is based on mutex3 from
16/// https://www.akkadia.org/drepper/futex.pdf
9/// This type supports static initialization and is based off of Golang 1.13 runtime.lock_futex:
10/// https://github.com/golang/go/blob/master/src/runtime/lock_futex.go
1711/// When an application is built in single threaded release mode, all the functions are
1812/// no-ops. In single threaded debug mode, there is deadlock detection.
1913pub const Mutex = if (builtin.single_threaded)
......@@ -43,83 +37,79 @@ pub const Mutex = if (builtin.single_threaded)
4337 return Held{ .mutex = self };
4438 }
4539 }
46else switch (builtin.os) {
47 builtin.Os.linux => struct {
48 /// 0: unlocked
49 /// 1: locked, no waiters
50 /// 2: locked, one or more waiters
51 lock: i32,
52
53 pub const Held = struct {
54 mutex: *Mutex,
55
56 pub fn release(self: Held) void {
57 const c = @atomicRmw(i32, &self.mutex.lock, AtomicRmwOp.Sub, 1, AtomicOrder.Release);
58 if (c != 1) {
59 _ = @atomicRmw(i32, &self.mutex.lock, AtomicRmwOp.Xchg, 0, AtomicOrder.Release);
60 const rc = linux.futex_wake(&self.mutex.lock, linux.FUTEX_WAKE | linux.FUTEX_PRIVATE_FLAG, 1);
61 switch (linux.getErrno(rc)) {
62 0 => {},
63 linux.EINVAL => unreachable,
64 else => unreachable,
65 }
66 }
67 }
40else struct {
41 state: u32, // TODO: make this an enum
42 parker: ThreadParker,
43
44 const Unlocked = 0;
45 const Sleeping = 1;
46 const Locked = 2;
47
48 /// number of iterations to spin yielding the cpu
49 const SPIN_CPU = 4;
50 /// number of iterations to perform in the cpu yield loop
51 const SPIN_CPU_COUNT = 30;
52 /// number of iterations to spin yielding the thread
53 const SPIN_THREAD = 1;
54
55 pub fn init() Mutex {
56 return Mutex{
57 .state = Unlocked,
58 .parker = ThreadParker.init(),
6859 };
60 }
6961
70 pub fn init() Mutex {
71 return Mutex{ .lock = 0 };
72 }
62 pub fn deinit(self: *Mutex) void {
63 self.parker.deinit();
64 }
7365
74 pub fn deinit(self: *Mutex) void {}
66 pub const Held = struct {
67 mutex: *Mutex,
7568
76 pub fn acquire(self: *Mutex) Held {
77 var c = @cmpxchgWeak(i32, &self.lock, 0, 1, AtomicOrder.Acquire, AtomicOrder.Monotonic) orelse
78 return Held{ .mutex = self };
79 if (c != 2)
80 c = @atomicRmw(i32, &self.lock, AtomicRmwOp.Xchg, 2, AtomicOrder.Acquire);
81 while (c != 0) {
82 const rc = linux.futex_wait(&self.lock, linux.FUTEX_WAIT | linux.FUTEX_PRIVATE_FLAG, 2, null);
83 switch (linux.getErrno(rc)) {
84 0, linux.EINTR, linux.EAGAIN => {},
85 linux.EINVAL => unreachable,
86 else => unreachable,
87 }
88 c = @atomicRmw(i32, &self.lock, AtomicRmwOp.Xchg, 2, AtomicOrder.Acquire);
69 pub fn release(self: Held) void {
70 switch (@atomicRmw(u32, &self.mutex.state, .Xchg, Unlocked, .Release)) {
71 Locked => {},
72 Sleeping => self.mutex.parker.unpark(&self.mutex.state),
73 Unlocked => unreachable, // unlocking an unlocked mutex
74 else => unreachable, // should never be anything else
8975 }
90 return Held{ .mutex = self };
9176 }
92 },
93 // TODO once https://github.com/ziglang/zig/issues/287 (copy elision) is solved, we can make a
94 // better implementation of this. The problem is we need the init() function to have access to
95 // the address of the CRITICAL_SECTION, and then have it not move.
96 builtin.Os.windows => std.StaticallyInitializedMutex,
97 else => struct {
98 /// TODO better implementation than spin lock.
99 /// When changing this, one must also change the corresponding
100 /// std.StaticallyInitializedMutex code, since it aliases this type,
101 /// under the assumption that it works both statically and at runtime.
102 lock: SpinLock,
77 };
10378
104 pub const Held = struct {
105 mutex: *Mutex,
79 pub fn acquire(self: *Mutex) Held {
80 // Try and speculatively grab the lock.
81 // If it fails, the state is either Locked or Sleeping
82 // depending on if theres a thread stuck sleeping below.
83 var state = @atomicRmw(u32, &self.state, .Xchg, Locked, .Acquire);
84 if (state == Unlocked)
85 return Held{ .mutex = self };
10686
107 pub fn release(self: Held) void {
108 SpinLock.Held.release(SpinLock.Held{ .spinlock = &self.mutex.lock });
87 while (true) {
88 // try and acquire the lock using cpu spinning on failure
89 var spin: usize = 0;
90 while (spin < SPIN_CPU) : (spin += 1) {
91 var value = @atomicLoad(u32, &self.state, .Monotonic);
92 while (value == Unlocked)
93 value = @cmpxchgWeak(u32, &self.state, Unlocked, state, .Acquire, .Monotonic) orelse return Held{ .mutex = self };
94 SpinLock.yield(SPIN_CPU_COUNT);
10995 }
110 };
11196
112 pub fn init() Mutex {
113 return Mutex{ .lock = SpinLock.init() };
114 }
115
116 pub fn deinit(self: *Mutex) void {}
97 // try and acquire the lock using thread rescheduling on failure
98 spin = 0;
99 while (spin < SPIN_THREAD) : (spin += 1) {
100 var value = @atomicLoad(u32, &self.state, .Monotonic);
101 while (value == Unlocked)
102 value = @cmpxchgWeak(u32, &self.state, Unlocked, state, .Acquire, .Monotonic) orelse return Held{ .mutex = self };
103 std.os.sched_yield();
104 }
117105
118 pub fn acquire(self: *Mutex) Held {
119 _ = self.lock.acquire();
120 return Held{ .mutex = self };
106 // failed to acquire the lock, go to sleep until woken up by `Held.release()`
107 if (@atomicRmw(u32, &self.state, .Xchg, Sleeping, .Acquire) == Unlocked)
108 return Held{ .mutex = self };
109 state = Sleeping;
110 self.parker.park(&self.state, Sleeping);
121111 }
122 },
112 }
123113};
124114
125115const TestContext = struct {
lib/std/os.zig+7
......@@ -3171,3 +3171,10 @@ pub fn dn_expand(
31713171 }
31723172 return error.InvalidDnsPacket;
31733173}
3174
3175pub fn sched_yield() void {
3176 switch (builtin.os) {
3177 .windows => _ = windows.kernel32.SwitchToThread(),
3178 else => assert(system.sched_yield() == 0),
3179 }
3180}
lib/std/os/linux.zig+4
......@@ -954,6 +954,10 @@ pub fn fremovexattr(fd: usize, name: [*]const u8) usize {
954954 return syscall2(SYS_fremovexattr, fd, @ptrToInt(name));
955955}
956956
957pub fn sched_yield() usize {
958 return syscall0(SYS_sched_yield);
959}
960
957961pub fn sched_getaffinity(pid: i32, size: usize, set: *cpu_set_t) usize {
958962 const rc = syscall3(SYS_sched_getaffinity, @bitCast(usize, isize(pid)), size, @ptrToInt(set));
959963 if (@bitCast(isize, rc) < 0) return rc;
lib/std/os/windows/kernel32.zig+2
......@@ -184,6 +184,8 @@ pub extern "kernel32" stdcallcc fn SetHandleInformation(hObject: HANDLE, dwMask:
184184
185185pub extern "kernel32" stdcallcc fn Sleep(dwMilliseconds: DWORD) void;
186186
187pub extern "kernel32" stdcallcc fn SwitchToThread() BOOL;
188
187189pub extern "kernel32" stdcallcc fn TerminateProcess(hProcess: HANDLE, uExitCode: UINT) BOOL;
188190
189191pub extern "kernel32" stdcallcc fn TlsAlloc() DWORD;
lib/std/os/windows/ntdll.zig+18
......@@ -43,3 +43,21 @@ pub extern "NtDll" stdcallcc fn NtQueryDirectoryFile(
4343 FileName: ?*UNICODE_STRING,
4444 RestartScan: BOOLEAN,
4545) NTSTATUS;
46pub extern "NtDll" stdcallcc fn NtCreateKeyedEvent(
47 KeyedEventHandle: *HANDLE,
48 DesiredAccess: ACCESS_MASK,
49 ObjectAttributes: ?PVOID,
50 Flags: ULONG,
51) NTSTATUS;
52pub extern "NtDll" stdcallcc fn NtReleaseKeyedEvent(
53 EventHandle: HANDLE,
54 Key: *const c_void,
55 Alertable: BOOLEAN,
56 Timeout: ?*LARGE_INTEGER,
57) NTSTATUS;
58pub extern "NtDll" stdcallcc fn NtWaitForKeyedEvent(
59 EventHandle: HANDLE,
60 Key: *const c_void,
61 Alertable: BOOLEAN,
62 Timeout: ?*LARGE_INTEGER,
63) NTSTATUS;
lib/std/parker.zig created+322
......@@ -0,0 +1,322 @@
1const std = @import("std.zig");
2const builtin = @import("builtin");
3const time = std.time;
4const testing = std.testing;
5const assert = std.debug.assert;
6const SpinLock = std.SpinLock;
7const linux = std.os.linux;
8const windows = std.os.windows;
9
10pub const ThreadParker = switch (builtin.os) {
11 .macosx,
12 .tvos,
13 .ios,
14 .watchos,
15 .netbsd,
16 .openbsd,
17 .freebsd,
18 .kfreebsd,
19 .dragonfly,
20 .haiku,
21 .hermit,
22 .solaris,
23 .minix,
24 .fuchsia,
25 .emscripten => if (builtin.link_libc) PosixParker else SpinParker,
26 .linux => if (builtin.link_libc) PosixParker else LinuxParker,
27 .windows => WindowsParker,
28 else => SpinParker,
29};
30
31const SpinParker = struct {
32 pub fn init() SpinParker {
33 return SpinParker{};
34 }
35 pub fn deinit(self: *SpinParker) void {}
36
37 pub fn unpark(self: *SpinParker, ptr: *const u32) void {}
38
39 pub fn park(self: *SpinParker, ptr: *const u32, expected: u32) void {
40 var backoff = SpinLock.Backoff.init();
41 while (@atomicLoad(u32, ptr, .Acquire) == expected)
42 backoff.yield();
43 }
44};
45
46const LinuxParker = struct {
47 pub fn init() LinuxParker {
48 return LinuxParker{};
49 }
50 pub fn deinit(self: *LinuxParker) void {}
51
52 pub fn unpark(self: *LinuxParker, ptr: *const u32) void {
53 const rc = linux.futex_wake(@ptrCast(*const i32, ptr), linux.FUTEX_WAKE | linux.FUTEX_PRIVATE_FLAG, 1);
54 assert(linux.getErrno(rc) == 0);
55 }
56
57 pub fn park(self: *LinuxParker, ptr: *const u32, expected: u32) void {
58 const value = @intCast(i32, expected);
59 while (@atomicLoad(u32, ptr, .Acquire) == expected) {
60 const rc = linux.futex_wait(@ptrCast(*const i32, ptr), linux.FUTEX_WAIT | linux.FUTEX_PRIVATE_FLAG, value, null);
61 switch (linux.getErrno(rc)) {
62 0, linux.EAGAIN => return,
63 linux.EINTR => continue,
64 linux.EINVAL => unreachable,
65 else => unreachable,
66 }
67 }
68 }
69};
70
71const WindowsParker = struct {
72 waiters: u32,
73
74 pub fn init() WindowsParker {
75 return WindowsParker{ .waiters = 0 };
76 }
77 pub fn deinit(self: *WindowsParker) void {}
78
79 pub fn unpark(self: *WindowsParker, ptr: *const u32) void {
80 const key = @ptrCast(*const c_void, ptr);
81 const handle = getEventHandle() orelse return;
82
83 var waiting = @atomicLoad(u32, &self.waiters, .Monotonic);
84 while (waiting != 0) {
85 waiting = @cmpxchgWeak(u32, &self.waiters, waiting, waiting - 1, .Acquire, .Monotonic) orelse {
86 const rc = windows.ntdll.NtReleaseKeyedEvent(handle, key, windows.FALSE, null);
87 assert(rc == 0);
88 return;
89 };
90 }
91 }
92
93 pub fn park(self: *WindowsParker, ptr: *const u32, expected: u32) void {
94 var spin = SpinLock.Backoff.init();
95 const ev_handle = getEventHandle();
96 const key = @ptrCast(*const c_void, ptr);
97
98 while (@atomicLoad(u32, ptr, .Monotonic) == expected) {
99 if (ev_handle) |handle| {
100 _ = @atomicRmw(u32, &self.waiters, .Add, 1, .Release);
101 const rc = windows.ntdll.NtWaitForKeyedEvent(handle, key, windows.FALSE, null);
102 assert(rc == 0);
103 } else {
104 spin.yield();
105 }
106 }
107 }
108
109 var event_handle = std.lazyInit(windows.HANDLE);
110
111 fn getEventHandle() ?windows.HANDLE {
112 if (event_handle.get()) |handle_ptr|
113 return handle_ptr.*;
114 defer event_handle.resolve();
115
116 const access_mask = windows.GENERIC_READ | windows.GENERIC_WRITE;
117 if (windows.ntdll.NtCreateKeyedEvent(&event_handle.data, access_mask, null, 0) != 0)
118 return null;
119 return event_handle.data;
120 }
121};
122
123const PosixParker = struct {
124 cond: pthread_cond_t,
125 mutex: pthread_mutex_t,
126
127 pub fn init() PosixParker {
128 return PosixParker{
129 .cond = PTHREAD_COND_INITIALIZER,
130 .mutex = PTHREAD_MUTEX_INITIALIZER,
131 };
132 }
133
134 pub fn deinit(self: *PosixParker) void {
135 // On dragonfly, the destroy functions return EINVAL if they were initialized statically.
136 const retm = pthread_mutex_destroy(&self.mutex);
137 assert(retm == 0 or retm == (if (builtin.os == .dragonfly) os.EINVAL else 0));
138 const retc = pthread_cond_destroy(&self.cond);
139 assert(retc == 0 or retc == (if (builtin.os == .dragonfly) os.EINVAL else 0));
140 }
141
142 pub fn unpark(self: *PosixParker, ptr: *const u32) void {
143 assert(pthread_mutex_lock(&self.mutex) == 0);
144 defer assert(pthread_mutex_unlock(&self.mutex) == 0);
145 assert(pthread_cond_signal(&self.cond) == 0);
146 }
147
148 pub fn park(self: *PosixParker, ptr: *const u32, expected: u32) void {
149 assert(pthread_mutex_lock(&self.mutex) == 0);
150 defer assert(pthread_mutex_unlock(&self.mutex) == 0);
151 while (@atomicLoad(u32, ptr, .Acquire) == expected)
152 assert(pthread_cond_wait(&self.cond, &self.mutex) == 0);
153 }
154
155 const PTHREAD_MUTEX_INITIALIZER = pthread_mutex_t{};
156 extern "c" fn pthread_mutex_lock(mutex: *pthread_mutex_t) c_int;
157 extern "c" fn pthread_mutex_unlock(mutex: *pthread_mutex_t) c_int;
158 extern "c" fn pthread_mutex_destroy(mutex: *pthread_mutex_t) c_int;
159
160 const PTHREAD_COND_INITIALIZER = pthread_cond_t{};
161 extern "c" fn pthread_cond_wait(noalias cond: *pthread_cond_t, noalias mutex: *pthread_mutex_t) c_int;
162 extern "c" fn pthread_cond_signal(cond: *pthread_cond_t) c_int;
163 extern "c" fn pthread_cond_destroy(cond: *pthread_cond_t) c_int;
164
165 // https://github.com/rust-lang/libc
166 usingnamespace switch (builtin.os) {
167 .macosx, .tvos, .ios, .watchos => struct {
168 pub const pthread_mutex_t = extern struct {
169 __sig: c_long = 0x32AAABA7,
170 __opaque: [__PTHREAD_MUTEX_SIZE__]u8 = [_]u8{0} ** __PTHREAD_MUTEX_SIZE__,
171 };
172 pub const pthread_cond_t = extern struct {
173 __sig: c_long = 0x3CB0B1BB,
174 __opaque: [__PTHREAD_COND_SIZE__]u8 = [_]u8{0} ** __PTHREAD_COND_SIZE__,
175 };
176 const __PTHREAD_MUTEX_SIZE__ = if (@sizeOf(usize) == 8) 56 else 40;
177 const __PTHREAD_COND_SIZE__ = if (@sizeOf(usize) == 8) 40 else 24;
178 },
179 .netbsd => struct {
180 pub const pthread_mutex_t = extern struct {
181 ptm_magic: c_uint = 0x33330003,
182 ptm_errorcheck: padded_spin_t = 0,
183 ptm_unused: padded_spin_t = 0,
184 ptm_owner: usize = 0,
185 ptm_waiters: ?*u8 = null,
186 ptm_recursed: c_uint = 0,
187 ptm_spare2: ?*c_void = null,
188 };
189 pub const pthread_cond_t = extern struct {
190 ptc_magic: c_uint = 0x55550005,
191 ptc_lock: pthread_spin_t = 0,
192 ptc_waiters_first: ?*u8 = null,
193 ptc_waiters_last: ?*u8 = null,
194 ptc_mutex: ?*pthread_mutex_t = null,
195 ptc_private: ?*c_void = null,
196 };
197 const pthread_spin_t = if (builtin.arch == .arm or .arch == .powerpc) c_int else u8;
198 const padded_spin_t = switch (builtin.arch) {
199 .sparc, .sparcel, .sparcv9, .i386, .x86_64, .le64 => u32,
200 else => spin_t,
201 };
202 },
203 .openbsd, .freebsd, .kfreebsd, .dragonfly => struct {
204 pub const pthread_mutex_t = extern struct {
205 inner: ?*c_void = null,
206 };
207 pub const pthread_cond_t = extern struct {
208 inner: ?*c_void = null,
209 };
210 },
211 .haiku => struct {
212 pub const pthread_mutex_t = extern struct {
213 flags: u32 = 0,
214 lock: i32 = 0,
215 unused: i32 = -42,
216 owner: i32 = -1,
217 owner_count: i32 = 0,
218 };
219 pub const pthread_cond_t = extern struct {
220 flags: u32 = 0,
221 unused: i32 = -42,
222 mutex: ?*c_void = null,
223 waiter_count: i32 = 0,
224 lock: i32 = 0,
225 };
226 },
227 .hermit => struct {
228 pub const pthread_mutex_t = extern struct {
229 inner: usize = ~usize(0),
230 };
231 pub const pthread_cond_t = extern struct {
232 inner: usize = ~usize(0),
233 };
234 },
235 .solaris => struct {
236 pub const pthread_mutex_t = extern struct {
237 __pthread_mutex_flag1: u16 = 0,
238 __pthread_mutex_flag2: u8 = 0,
239 __pthread_mutex_ceiling: u8 = 0,
240 __pthread_mutex_type: u16 = 0,
241 __pthread_mutex_magic: u16 = 0x4d58,
242 __pthread_mutex_lock: u64 = 0,
243 __pthread_mutex_data: u64 = 0,
244 };
245 pub const pthread_cond_t = extern struct {
246 __pthread_cond_flag: u32 = 0,
247 __pthread_cond_type: u16 = 0,
248 __pthread_cond_magic: u16 = 0x4356,
249 __pthread_cond_data: u64 = 0,
250 };
251 },
252 .fuchsia, .minix, .linux => struct {
253 pub const pthread_mutex_t = extern struct {
254 size: [__SIZEOF_PTHREAD_MUTEX_T]u8 align(@alignOf(usize)) = [_]u8{0} ** __SIZEOF_PTHREAD_MUTEX_T,
255 };
256 pub const pthread_cond_t = extern struct {
257 size: [__SIZEOF_PTHREAD_COND_T]u8 align(@alignOf(usize)) = [_]u8{0} ** __SIZEOF_PTHREAD_COND_T,
258 };
259 const __SIZEOF_PTHREAD_COND_T = 48;
260 const __SIZEOF_PTHREAD_MUTEX_T = if (builtin.os == .fuchsia) 40 else switch (builtin.abi) {
261 .musl, .musleabi, .musleabihf => if (@sizeOf(usize) == 8) 40 else 24,
262 .gnu, .gnuabin32, .gnuabi64, .gnueabi, .gnueabihf, .gnux32 => switch (builtin.arch) {
263 .aarch64 => 48,
264 .x86_64 => if (builtin.abi == .gnux32) 40 else 32,
265 .mips64, .powerpc64, .powerpc64le, .sparcv9 => 40,
266 else => if (@sizeOf(usize) == 8) 40 else 24,
267 },
268 else => unreachable,
269 };
270 },
271 .emscripten => struct {
272 pub const pthread_mutex_t = extern struct {
273 size: [__SIZEOF_PTHREAD_MUTEX_T]u8 align(4) = [_]u8{0} ** __SIZEOF_PTHREAD_MUTEX_T,
274 };
275 pub const pthread_cond_t = extern struct {
276 size: [__SIZEOF_PTHREAD_COND_T]u8 align(@alignOf(usize)) = [_]u8{0} ** __SIZEOF_PTHREAD_COND_T,
277 };
278 const __SIZEOF_PTHREAD_COND_T = 48;
279 const __SIZEOF_PTHREAD_MUTEX_T = 28;
280 },
281 else => unreachable,
282 };
283};
284
285test "std.ThreadParker" {
286 if (builtin.single_threaded)
287 return error.SkipZigTest;
288
289 const Context = struct {
290 parker: ThreadParker,
291 data: u32,
292
293 fn receiver(self: *@This()) void {
294 self.parker.park(&self.data, 0); // receives 1
295 assert(@atomicRmw(u32, &self.data, .Xchg, 2, .SeqCst) == 1); // sends 2
296 self.parker.unpark(&self.data); // wakes up waiters on 2
297 self.parker.park(&self.data, 2); // receives 3
298 assert(@atomicRmw(u32, &self.data, .Xchg, 4, .SeqCst) == 3); // sends 4
299 self.parker.unpark(&self.data); // wakes up waiters on 4
300 }
301
302 fn sender(self: *@This()) void {
303 assert(@atomicRmw(u32, &self.data, .Xchg, 1, .SeqCst) == 0); // sends 1
304 self.parker.unpark(&self.data); // wakes up waiters on 1
305 self.parker.park(&self.data, 1); // receives 2
306 assert(@atomicRmw(u32, &self.data, .Xchg, 3, .SeqCst) == 2); // sends 3
307 self.parker.unpark(&self.data); // wakes up waiters on 3
308 self.parker.park(&self.data, 3); // receives 4
309 }
310 };
311
312 var context = Context{
313 .parker = ThreadParker.init(),
314 .data = 0,
315 };
316 defer context.parker.deinit();
317
318 var receiver = try std.Thread.spawn(&context, Context.receiver);
319 defer receiver.wait();
320
321 context.sender();
322}
\ No newline at end of file
lib/std/spinlock.zig+42-4
......@@ -1,8 +1,8 @@
11const std = @import("std.zig");
22const builtin = @import("builtin");
3const AtomicOrder = builtin.AtomicOrder;
4const AtomicRmwOp = builtin.AtomicRmwOp;
53const assert = std.debug.assert;
4const time = std.time;
5const os = std.os;
66
77pub const SpinLock = struct {
88 lock: u8, // TODO use a bool or enum
......@@ -11,7 +11,8 @@ pub const SpinLock = struct {
1111 spinlock: *SpinLock,
1212
1313 pub fn release(self: Held) void {
14 assert(@atomicRmw(u8, &self.spinlock.lock, builtin.AtomicRmwOp.Xchg, 0, AtomicOrder.SeqCst) == 1);
14 // TODO: @atomicStore() https://github.com/ziglang/zig/issues/2995
15 assert(@atomicRmw(u8, &self.spinlock.lock, .Xchg, 0, .Release) == 1);
1516 }
1617 };
1718
......@@ -20,9 +21,46 @@ pub const SpinLock = struct {
2021 }
2122
2223 pub fn acquire(self: *SpinLock) Held {
23 while (@atomicRmw(u8, &self.lock, builtin.AtomicRmwOp.Xchg, 1, AtomicOrder.SeqCst) != 0) {}
24 var backoff = Backoff.init();
25 while (@atomicRmw(u8, &self.lock, .Xchg, 1, .Acquire) != 0)
26 backoff.yield();
2427 return Held{ .spinlock = self };
2528 }
29
30 pub fn yield(iterations: usize) void {
31 var i = iterations;
32 while (i != 0) : (i -= 1) {
33 switch (builtin.arch) {
34 .i386, .x86_64 => asm volatile("pause"),
35 .arm, .aarch64 => asm volatile("yield"),
36 else => time.sleep(0),
37 }
38 }
39 }
40
41 /// Provides a method to incrementally yield longer each time its called.
42 pub const Backoff = struct {
43 iteration: usize,
44
45 pub fn init() @This() {
46 return @This(){ .iteration = 0 };
47 }
48
49 /// Modified hybrid yielding from
50 /// http://www.1024cores.net/home/lock-free-algorithms/tricks/spinning
51 pub fn yield(self: *@This()) void {
52 defer self.iteration +%= 1;
53 if (self.iteration < 20) {
54 SpinLock.yield(self.iteration);
55 } else if (self.iteration < 24) {
56 os.sched_yield();
57 } else if (self.iteration < 26) {
58 time.sleep(1 * time.millisecond);
59 } else {
60 time.sleep(10 * time.millisecond);
61 }
62 }
63 };
2664};
2765
2866test "spinlock" {
lib/std/statically_initialized_mutex.zig deleted-105
......@@ -1,105 +0,0 @@
1const std = @import("std.zig");
2const builtin = @import("builtin");
3const AtomicOrder = builtin.AtomicOrder;
4const AtomicRmwOp = builtin.AtomicRmwOp;
5const assert = std.debug.assert;
6const expect = std.testing.expect;
7const windows = std.os.windows;
8
9/// Lock may be held only once. If the same thread
10/// tries to acquire the same mutex twice, it deadlocks.
11/// This type is intended to be initialized statically. If you don't
12/// require static initialization, use std.Mutex.
13/// On Windows, this mutex allocates resources when it is
14/// first used, and the resources cannot be freed.
15/// On Linux, this is an alias of std.Mutex.
16pub const StaticallyInitializedMutex = switch (builtin.os) {
17 builtin.Os.linux => std.Mutex,
18 builtin.Os.windows => struct {
19 lock: windows.CRITICAL_SECTION,
20 init_once: windows.RTL_RUN_ONCE,
21
22 pub const Held = struct {
23 mutex: *StaticallyInitializedMutex,
24
25 pub fn release(self: Held) void {
26 windows.kernel32.LeaveCriticalSection(&self.mutex.lock);
27 }
28 };
29
30 pub fn init() StaticallyInitializedMutex {
31 return StaticallyInitializedMutex{
32 .lock = undefined,
33 .init_once = windows.INIT_ONCE_STATIC_INIT,
34 };
35 }
36
37 extern fn initCriticalSection(
38 InitOnce: *windows.RTL_RUN_ONCE,
39 Parameter: ?*c_void,
40 Context: ?*c_void,
41 ) windows.BOOL {
42 const lock = @ptrCast(*windows.CRITICAL_SECTION, @alignCast(@alignOf(windows.CRITICAL_SECTION), Parameter));
43 windows.kernel32.InitializeCriticalSection(lock);
44 return windows.TRUE;
45 }
46
47 /// TODO: once https://github.com/ziglang/zig/issues/287 is solved and std.Mutex has a better
48 /// implementation of a runtime initialized mutex, remove this function.
49 pub fn deinit(self: *StaticallyInitializedMutex) void {
50 windows.InitOnceExecuteOnce(&self.init_once, initCriticalSection, &self.lock, null);
51 windows.kernel32.DeleteCriticalSection(&self.lock);
52 }
53
54 pub fn acquire(self: *StaticallyInitializedMutex) Held {
55 windows.InitOnceExecuteOnce(&self.init_once, initCriticalSection, &self.lock, null);
56 windows.kernel32.EnterCriticalSection(&self.lock);
57 return Held{ .mutex = self };
58 }
59 },
60 else => std.Mutex,
61};
62
63test "std.StaticallyInitializedMutex" {
64 const TestContext = struct {
65 data: i128,
66
67 const TestContext = @This();
68 const incr_count = 10000;
69
70 var mutex = StaticallyInitializedMutex.init();
71
72 fn worker(ctx: *TestContext) void {
73 var i: usize = 0;
74 while (i != TestContext.incr_count) : (i += 1) {
75 const held = mutex.acquire();
76 defer held.release();
77
78 ctx.data += 1;
79 }
80 }
81 };
82
83 var plenty_of_memory = try std.heap.direct_allocator.alloc(u8, 300 * 1024);
84 defer std.heap.direct_allocator.free(plenty_of_memory);
85
86 var fixed_buffer_allocator = std.heap.ThreadSafeFixedBufferAllocator.init(plenty_of_memory);
87 var a = &fixed_buffer_allocator.allocator;
88
89 var context = TestContext{ .data = 0 };
90
91 if (builtin.single_threaded) {
92 TestContext.worker(&context);
93 expect(context.data == TestContext.incr_count);
94 } else {
95 const thread_count = 10;
96 var threads: [thread_count]*std.Thread = undefined;
97 for (threads) |*t| {
98 t.* = try std.Thread.spawn(&context, TestContext.worker);
99 }
100 for (threads) |t|
101 t.wait();
102
103 expect(context.data == thread_count * TestContext.incr_count);
104 }
105}
lib/std/std.zig+1-1
......@@ -19,11 +19,11 @@ pub const Progress = @import("progress.zig").Progress;
1919pub const SegmentedList = @import("segmented_list.zig").SegmentedList;
2020pub const SinglyLinkedList = @import("linked_list.zig").SinglyLinkedList;
2121pub const SpinLock = @import("spinlock.zig").SpinLock;
22pub const StaticallyInitializedMutex = @import("statically_initialized_mutex.zig").StaticallyInitializedMutex;
2322pub const StringHashMap = @import("hash_map.zig").StringHashMap;
2423pub const TailQueue = @import("linked_list.zig").TailQueue;
2524pub const Target = @import("target.zig").Target;
2625pub const Thread = @import("thread.zig").Thread;
26pub const ThreadParker = @import("parker.zig").ThreadParker;
2727
2828pub const atomic = @import("atomic.zig");
2929pub const base64 = @import("base64.zig");