This is an automated email from the ASF dual-hosted git repository.

Cole-Greer pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/tinkerpop.git


The following commit(s) were added to refs/heads/master by this push:
     new d56f99f3f9 Reduce async overhead in GraphBinary reads (#3630)
d56f99f3f9 is described below

commit d56f99f3f901536f64885ffb3c516e05b5fd3dfe
Author: kirill-stepanishin <[email protected]>
AuthorDate: Tue Sep 15 10:39:13 2026 -0700

    Reduce async overhead in GraphBinary reads (#3630)
    
    Adds a synchronous fast path to gremlin-javascript StreamReader
    
    Every primitive read in StreamReader awaits #ensure(n) first. For a 
buffered submit() response a chunk of
    the payload is already in memory, so that await often does nothing but 
still costs a Promise and a
    microtask per read, which dominates client CPU on large responses. Reads 
now return right away when the bytes are already buffered and only await when 
the buffer needs refilling.
    
    Assisted-by: Claude Code:claude-opus-4-8
---
 .../structure/io/binary/internals/StreamReader.js  | 94 +++++++++++++---------
 1 file changed, 57 insertions(+), 37 deletions(-)

diff --git 
a/gremlin-js/gremlin-javascript/lib/structure/io/binary/internals/StreamReader.js
 
b/gremlin-js/gremlin-javascript/lib/structure/io/binary/internals/StreamReader.js
index b1f1274812..c134afb54d 100644
--- 
a/gremlin-js/gremlin-javascript/lib/structure/io/binary/internals/StreamReader.js
+++ 
b/gremlin-js/gremlin-javascript/lib/structure/io/binary/internals/StreamReader.js
@@ -23,7 +23,8 @@ import { Buffer } from 'buffer';
  * Async byte reader that provides a uniform interface over both a complete 
Buffer
  * (for non-streaming submit()) and a ReadableStream (for streaming HTTP 
responses).
  *
- * Handles chunk boundaries transparently and blocks (awaits) until the 
requested bytes are available.
+ * Chunk boundaries are handled transparently. A read returns immediately when 
the
+ * requested bytes are already buffered, and otherwise awaits more of the 
stream.
  */
 export default class StreamReader {
   /** @type {Buffer} */
@@ -106,6 +107,28 @@ export default class StreamReader {
     this.#offset = 0;
   }
 
+  /**
+   * @param {number} n
+   * @returns {boolean} true when the next `n` bytes are already buffered and 
can be read
+   *   without pulling more of the stream
+   */
+  #available(n) {
+    return this.#offset + n <= this.#buffer.length;
+  }
+
+  /**
+   * Advance past the next `n` bytes and return the offset where they start. 
The caller must
+   * have already made those bytes available, via #available or #ensure.
+   * @param {number} n
+   * @returns {number} offset of the first of the `n` bytes
+   */
+  #advance(n) {
+    const at = this.#offset;
+    this.#offset = at + n;
+    this.#position += n;
+    return at;
+  }
+
   /**
    * Total number of bytes consumed so far (monotonically increasing).
    * Useful for error diagnostics.
@@ -121,83 +144,80 @@ export default class StreamReader {
    * @returns {Promise<Buffer>}
    */
   async readBytes(n) {
-    await this.#ensure(n);
-    const result = this.#buffer.subarray(this.#offset, this.#offset + n);
-    this.#offset += n;
-    this.#position += n;
-    return result;
+    if (!this.#available(n)) {
+      await this.#ensure(n);
+    }
+    const at = this.#advance(n);
+    return this.#buffer.subarray(at, at + n);
   }
 
   /**
    * @returns {Promise<number>} unsigned 8-bit integer
    */
   async readUInt8() {
-    await this.#ensure(1);
-    this.#position++;
-    return this.#buffer[this.#offset++];
+    if (!this.#available(1)) {
+      await this.#ensure(1);
+    }
+    return this.#buffer[this.#advance(1)];
   }
 
   /**
    * @returns {Promise<number>} signed 8-bit integer
    */
   async readByte() {
-    await this.#ensure(1);
-    this.#position++;
-    return this.#buffer.readInt8(this.#offset++);
+    if (!this.#available(1)) {
+      await this.#ensure(1);
+    }
+    return this.#buffer.readInt8(this.#advance(1));
   }
 
   /**
    * @returns {Promise<number>} signed 16-bit big-endian integer
    */
   async readInt16BE() {
-    await this.#ensure(2);
-    const v = this.#buffer.readInt16BE(this.#offset);
-    this.#offset += 2;
-    this.#position += 2;
-    return v;
+    if (!this.#available(2)) {
+      await this.#ensure(2);
+    }
+    return this.#buffer.readInt16BE(this.#advance(2));
   }
 
   /**
    * @returns {Promise<number>} signed 32-bit big-endian integer
    */
   async readInt32BE() {
-    await this.#ensure(4);
-    const v = this.#buffer.readInt32BE(this.#offset);
-    this.#offset += 4;
-    this.#position += 4;
-    return v;
+    if (!this.#available(4)) {
+      await this.#ensure(4);
+    }
+    return this.#buffer.readInt32BE(this.#advance(4));
   }
 
   /**
    * @returns {Promise<bigint>} signed 64-bit big-endian integer
    */
   async readBigInt64BE() {
-    await this.#ensure(8);
-    const v = this.#buffer.readBigInt64BE(this.#offset);
-    this.#offset += 8;
-    this.#position += 8;
-    return v;
+    if (!this.#available(8)) {
+      await this.#ensure(8);
+    }
+    return this.#buffer.readBigInt64BE(this.#advance(8));
   }
 
   /**
    * @returns {Promise<number>} 32-bit big-endian float
    */
   async readFloatBE() {
-    await this.#ensure(4);
-    const v = this.#buffer.readFloatBE(this.#offset);
-    this.#offset += 4;
-    this.#position += 4;
-    return v;
+    if (!this.#available(4)) {
+      await this.#ensure(4);
+    }
+    return this.#buffer.readFloatBE(this.#advance(4));
   }
 
   /**
    * @returns {Promise<number>} 64-bit big-endian double
    */
   async readDoubleBE() {
-    await this.#ensure(8);
-    const v = this.#buffer.readDoubleBE(this.#offset);
-    this.#offset += 8;
-    this.#position += 8;
-    return v;
+    if (!this.#available(8)) {
+      await this.#ensure(8);
+    }
+    return this.#buffer.readDoubleBE(this.#advance(8));
   }
 }

Reply via email to