| 248 | } |
| 249 | |
| 250 | function advanceReadIndexCAS( |
| 251 | controlView: Uint32Array, |
| 252 | readIdx: number, |
| 253 | nextIdx: number, |
| 254 | ): void { |
| 255 | const MAX_ADVANCE_RETRIES = 16; |
| 256 | let expectedIdx = readIdx; |
| 257 | |
| 258 | for (let attempt = 0; attempt < MAX_ADVANCE_RETRIES; attempt++) { |
| 259 | const exchanged = Atomics.compareExchange( |
| 260 | controlView, |
| 261 | CONTROL_READ_INDEX, |
| 262 | expectedIdx, |
| 263 | nextIdx, |
| 264 | ); |
| 265 | if (exchanged === expectedIdx) { |
| 266 | return; |
| 267 | } |
| 268 | const currentIdx = Atomics.load(controlView, CONTROL_READ_INDEX); |
| 269 | const hasProgressed = |
| 270 | currentIdx !== readIdx && |
| 271 | ((nextIdx > readIdx && (currentIdx >= nextIdx || currentIdx < readIdx)) || |
| 272 | (nextIdx < readIdx && currentIdx >= nextIdx && currentIdx < readIdx)); |
| 273 | if (hasProgressed) { |
| 274 | return; |
| 275 | } |
| 276 | expectedIdx = exchanged; |
| 277 | } |
| 278 | } |
| 279 | |
| 280 | export function createConsumer(buffer: SharedArrayBuffer): Consumer { |
| 281 | const controlView = new Uint32Array(buffer, 0, 8); |