Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
902 changes: 902 additions & 0 deletions .agents/skills/review-pr/SKILL.md

Large diffs are not rendered by default.

4 changes: 4 additions & 0 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,3 +102,7 @@ Integration tests use TestContainers to spin up QuestDB instances for realistic
- Each worker thread needs its own Sender instance (buffers cannot be shared)
- Protocol version 2 or higher is recommended for new implementations; v2 adds array columns and v3 adds DECIMAL
- Run `pnpm test:dist` after package-boundary changes; it checks both npm tarballs, ESM/CJS loading, browser bundling, and the absence of Node modules from the browser artifact.

## Workflow

- Do not watch or poll CI after pushing unless explicitly asked.
22 changes: 21 additions & 1 deletion benchmarks/sender.bench.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,10 +7,17 @@ import {
type QwpIngressResponse,
type QwpTableBuffer,
} from "../packages/client-core/src/_qwp/_core";
import { QwpSender, type QwpSenderSession } from "../packages/client-core/src/_qwp/sender";
import {
QwpSender,
type QwpSenderSession,
} from "../packages/client-core/src/_qwp/sender";
import { BENCHMARK_WORKLOADS, type BenchmarkRow } from "./workloads";

const ROWS = 10_000;
const NULLISH_COLUMN_NAMES = Array.from(
{ length: 20 },
(_, index) => `optional_${index}_${"x".repeat(82)}`,
);
let sink = 0;

class EncodingSession implements QwpSenderSession {
Expand Down Expand Up @@ -159,4 +166,17 @@ describe("high-level symbol dictionary modes", () => {
});
});

describe("high-level QwpSender with never-populated optional columns", () => {
bench("validate nullish names and publish", async () => {
const sender = senderFor(new EncodingSession(), "full");
await sender.table("bench_nullish").longColumn("present", 1n).atNow();
for (let row = 0; row < ROWS; row++) {
sender.table("bench_nullish");
for (const name of NULLISH_COLUMN_NAMES) sender.longColumn(name, null);
await sender.atNow();
}
await sender.flush();
});
});

export const senderBenchmarkSink = (): number => sink;
65 changes: 56 additions & 9 deletions packages/client-core/src/_qwp/_core/ingress.ts
Original file line number Diff line number Diff line change
Expand Up @@ -62,8 +62,9 @@ interface ColumnEncodeOptions {
deltaSymbols: boolean;
dictionary?: QwpSymbolDictionary;
/**
* Scoped to a single encodeQwpIngressFrame() call, so a column mutated
* between calls can never be sized from a stale plan.
* Scoped to a single frame plan -- one encodeQwpIngressFrame() call, or one
* measureQwpIngressFrame() and its encode() -- so a column mutated between
* frames can never be sized from a stale plan.
*/
plans: Map<QwpColumnBuffer, ColumnPlan>;
}
Expand Down Expand Up @@ -922,27 +923,73 @@ function planQwpIngressFrameEncoding(
};
}

/** @internal Measures without allocating the full frame output buffer. */
/** @internal A sized ingress frame that can be encoded without replanning. */
export interface QwpMeasuredIngressFrame {
/** Exact byte length encode() will produce. */
readonly byteLength: number;
/**
* Encodes the frame from the plan that measured it. Call it at most once,
* before the tables or the dictionary change, in the same synchronous
* section as the measurement.
*/
encode(): Uint8Array;
}

/**
* @internal Sizes a frame without allocating its output buffer.
*
* Measuring and then encoding used to plan the frame twice: the plan resolves
* every symbol against the dictionary and sizes every column, which made it
* the costliest step on the flush path. The measurement now holds on to its
* plan so encode() writes straight from it.
*
* On failure, either step truncates the dictionary to its size before
* measuring. A caller that discards a measured frame without encoding it must
* truncate the dictionary itself.
*/
export function measureQwpIngressFrame(
tables: readonly QwpTableBuffer[],
options: QwpIngressEncodeOptions = {},
): number {
): QwpMeasuredIngressFrame {
const dictionarySize = options.dictionary?.size;
try {
const plan = planQwpIngressFrameEncoding(tables, options);
return QWP_HEADER_SIZE + plan.payloadLength;
} catch (error) {
const restoreDictionary = (): void => {
if (dictionarySize !== undefined)
options.dictionary!.truncate(dictionarySize);
};
let plan: QwpIngressFrameEncodingPlan;
try {
plan = planQwpIngressFrameEncoding(tables, options);
} catch (error) {
restoreDictionary();
throw error;
}
return {
byteLength: QWP_HEADER_SIZE + plan.payloadLength,
encode: () => {
try {
return writeQwpIngressFrame(tables, plan);
} catch (error) {
restoreDictionary();
throw error;
}
},
};
}

function encodeQwpIngressFrameInternal(
tables: readonly QwpTableBuffer[],
options: QwpIngressEncodeOptions,
): Uint8Array {
const plan = planQwpIngressFrameEncoding(tables, options);
return writeQwpIngressFrame(
tables,
planQwpIngressFrameEncoding(tables, options),
);
}

function writeQwpIngressFrame(
tables: readonly QwpTableBuffer[],
plan: QwpIngressFrameEncodingPlan,
): Uint8Array {
const writer = new QwpByteWriter(QWP_HEADER_SIZE + plan.payloadLength);
writeQwpFrameHeader(writer, {
flags: plan.flags,
Expand Down
7 changes: 5 additions & 2 deletions packages/client-core/src/_qwp/ingress-session.ts
Original file line number Diff line number Diff line change
Expand Up @@ -142,12 +142,15 @@ function planIngressFrames(
// oversized candidate before every bisection briefly allocated many
// multiples of the negotiated cap and could exhaust the process before
// splitting had a chance to help.
frameByteLength = measureQwpIngressFrame(candidate, candidateOptions);
// A candidate that fits is encoded from the plan that measured it,
// in this same synchronous section, rather than planned a second time.
const measured = measureQwpIngressFrame(candidate, candidateOptions);
frameByteLength = measured.byteLength;
if (frameByteLength > maxBatchSizeBytes) {
if (dictionarySize !== undefined)
dictionary!.truncate(dictionarySize);
} else {
const frame = encodeQwpIngressFrame(candidate, candidateOptions);
const frame = measured.encode();
frames.push(frame);
if (dictionary) confirmedMaxSymbolId = dictionary.size - 1;
return;
Expand Down
Loading
Loading