| ... | @@ -101,30 +101,48 @@ const DebugEvent = struct { | ... | @@ -101,30 +101,48 @@ const DebugEvent = struct { |
| 101 | }; | 101 | }; |
| 102 | | 102 | |
| 103 | const PosixEvent = struct { | 103 | const PosixEvent = struct { |
| 104 | sem: c.sem_t, | 104 | sem: c.sem_t = undefined, |
| | 105 | /// Sadly this is needed because pthreads semaphore API does not |
| | 106 | /// support static initialization. |
| | 107 | init_mutex: std.mutex.PthreadMutex = .{}, |
| | 108 | state: enum { uninit, init } = .uninit, |
| 105 | | 109 | |
| 106 | fn init() PosixEvent { | 110 | fn init() PosixEvent { |
| 107 | return PosixEvent{ | 111 | return .{}; |
| 108 | .sem = c.sem_t.init(0, 0), | | |
| 109 | }; | | |
| 110 | } | 112 | } |
| 111 | | 113 | |
| | 114 | /// Not thread-safe. |
| 112 | fn deinit(self: *PosixEvent) void { | 115 | fn deinit(self: *PosixEvent) void { |
| 113 | assert(c.sem_destroy(&self.sem) == 0); | 116 | switch (self.state) { |
| | 117 | .uninit => {}, |
| | 118 | .init => { |
| | 119 | assert(c.sem_destroy(&self.sem) == 0); |
| | 120 | }, |
| | 121 | } |
| | 122 | self.* = undefined; |
| 114 | } | 123 | } |
| 115 | | 124 | |
| 116 | fn reset(self: *PosixEvent) void { | 125 | fn reset(self: *PosixEvent) void { |
| 117 | self.deinit(); | 126 | const sem = self.getInitializedSem(); |
| 118 | assert(c.sem_init(&self.sem, 0, 0) == 0); | 127 | while (true) { |
| | 128 | switch (c.getErrno(c.sem_trywait(sem))) { |
| | 129 | 0 => continue, // Need to make it go to zero. |
| | 130 | c.EINTR => continue, |
| | 131 | c.EINVAL => unreachable, |
| | 132 | c.EAGAIN => return, // The semaphore currently has the value zero. |
| | 133 | else => unreachable, |
| | 134 | } |
| | 135 | } |
| 119 | } | 136 | } |
| 120 | | 137 | |
| 121 | fn set(self: *PosixEvent) void { | 138 | fn set(self: *PosixEvent) void { |
| 122 | assert(c.sem_post(&self.sem) == 0); | 139 | assert(c.sem_post(self.getInitializedSem()) == 0); |
| 123 | } | 140 | } |
| 124 | | 141 | |
| 125 | fn wait(self: *PosixEvent) void { | 142 | fn wait(self: *PosixEvent) void { |
| | 143 | const sem = self.getInitializedSem(); |
| 126 | while (true) { | 144 | while (true) { |
| 127 | switch (c.getErrno(c.sem_wait(&self.sem))) { | 145 | switch (c.getErrno(c.sem_wait(sem))) { |
| 128 | 0 => return, | 146 | 0 => return, |
| 129 | c.EINTR => continue, | 147 | c.EINTR => continue, |
| 130 | c.EINVAL => unreachable, | 148 | c.EINVAL => unreachable, |
| ... | @@ -148,6 +166,7 @@ const PosixEvent = struct { | ... | @@ -148,6 +166,7 @@ const PosixEvent = struct { |
| 148 | } | 166 | } |
| 149 | ts.tv_sec = @intCast(@TypeOf(ts.tv_sec), @divFloor(timeout_abs, time.ns_per_s)); | 167 | ts.tv_sec = @intCast(@TypeOf(ts.tv_sec), @divFloor(timeout_abs, time.ns_per_s)); |
| 150 | ts.tv_nsec = @intCast(@TypeOf(ts.tv_nsec), @mod(timeout_abs, time.ns_per_s)); | 168 | ts.tv_nsec = @intCast(@TypeOf(ts.tv_nsec), @mod(timeout_abs, time.ns_per_s)); |
| | 169 | const sem = self.getInitializedSem(); |
| 151 | while (true) { | 170 | while (true) { |
| 152 | switch (c.getErrno(c.sem_timedwait(&self.sem, &ts))) { | 171 | switch (c.getErrno(c.sem_timedwait(&self.sem, &ts))) { |
| 153 | 0 => return, | 172 | 0 => return, |
| ... | @@ -158,6 +177,20 @@ const PosixEvent = struct { | ... | @@ -158,6 +177,20 @@ const PosixEvent = struct { |
| 158 | } | 177 | } |
| 159 | } | 178 | } |
| 160 | } | 179 | } |
| | 180 | |
| | 181 | fn getInitializedSem(self: *PosixEvent) *c.sem_t { |
| | 182 | const held = self.init_mutex.acquire(); |
| | 183 | defer held.release(); |
| | 184 | |
| | 185 | switch (self.state) { |
| | 186 | .init => return &self.sem, |
| | 187 | .uninit => { |
| | 188 | self.state = .init; |
| | 189 | assert(c.sem_init(&self.sem, 0, 0) == 0); |
| | 190 | return &self.sem; |
| | 191 | }, |
| | 192 | } |
| | 193 | } |
| 161 | }; | 194 | }; |
| 162 | | 195 | |
| 163 | const AtomicEvent = struct { | 196 | const AtomicEvent = struct { |