pinojs / pinojs/thread-stream

Simpler algorithm

Open
#59 9 comments 0 reactions 0 assignees View on GitHub

Nobody has claimed this yet.

Dominant language
JavaScript
Stars
261
Forks
31
Avg merge
2d 9h
Merged PRs (30d)
4

Description

Couldn't quite get thread-stream running well in our services so we re-implemented something similar. We used a simpler concurrent queue implementation that might be useful for thread-stream as well.

Basically we do it like this.

const MAX_LEN = 8 * 1024
const BUF_END = 256 * MAX_LEN
const BUF_LEN = BUF_END + MAX_LEN

const WRITE_INDEX = 4
const READ_INDEX = 8

const sharedBuffer = new SharedArrayBuffer(BUF_LEN)
const sharedState = new SharedArrayBuffer(128)

const state = new Int32Array(sharedState)
const buffer = Buffer.from(sharedBuffer)


buffer[0] = 31 // Item header
Atomics.store(state, WRITE_INDEX, 1)
Atomics.notify(state, WRITE_INDEX)

// producer

let draining = false
async function drain () {
  draining = true

  let write = Atomics.load(state, WRITE_INDEX)
  let read = Atomics.load(state, READ_INDEX)

  while (write <= read && write + MAX_LEN > read) {
    const { async, value } = Atomics.waitAsync(state, READ_INDEX, read)
    if (async) {
      await value
    }
    read = Atomics.load(state, READ_INDEX)
  }

  draining = false
  emit('drain')
}

send(data) {
  const len = MAX_LEN // or Buffer.byteLength(name)

  let read = Atomics.load(this._state, READ_INDEX)

  while (write < read && write + len > read) {
    Atomics.wait(this._state, READ_INDEX, read)
    read = Atomics.load(state, READ_INDEX)
  }

  write += this._buffer.write(name, this._write)
  buffer[write++] = 31

  if (write > BUF_END) {
   write = 0
  }

  Atomics.store(state, WRITE_INDEX, write)
  Atomics.notify(state, WRITE_INDEX)

  const needDrain = write + MAX_LEN >= BUF_END

  if (needDrain && !draining) {
    draining = true
    drain()
  }

  return !draining
}

// consumer
async function * receive () {
  while (true) {
    let write = Atomics.load(state, WRITE_INDEX)

    if (read > BUF_END) {
      read = 0
    }

    while (read === write) {
      const { async, value } = Atomics.waitAsync(state, WRITE_INDEX, write)
      if (async) {
        await value
      }
      write = Atomics.load(state, WRITE_INDEX)
    }

    if (write < read) {
      write = BUF_END
    }

    const arr = []
    while (read < write) {
      const idx = buffer.indexOf(31, read)
      arr.push(buffer.toString('utf-8', read, idx))
      read = idx + 1
    }

    Atomics.store(state, READ_INDEX, read)
    Atomics.notify(state, READ_INDEX)

    yield* arr
  }
}

Contributor guide

No contributing guide indexed for this repository

First steps

  1. Read the whole issue, then the project's contributing guide.
  2. Comment on the issue to say you are picking it up — it saves two people doing the same work.
  3. Fork the repository and make your change on a branch.
  4. Open a pull request that references the issue number.

Research direction

The issue names no file, test, or entry point; start by locating the existing thread-stream queue implementation and comparing it with the producer and consumer sketch in the report. Clarify the intended algorithm and acceptance criteria before changing code, then validate the agreed behavior with the project's existing tests.

Written by the indexing model from the issue text.

Assessment

Tech stack
javascript
Domain
backend
Issue type
Feature
Difficulty
5/5
Estimated time
Over a week
Activity status
Stale
Clarity
Needs clarification
Newbie friendliness
25/100

Get new issues in your inbox

A short digest of beginner-friendly GitHub issues.