From 8b1433ed8defcbe04eac330a0694e8b3d89d79c3 Mon Sep 17 00:00:00 2001 From: Greg Baker Date: Fri, 4 Sep 2026 17:59:56 -0700 Subject: [PATCH 1/5] work --- apps/zero-throughput/api/index.ts | 133 ++++++++++++++++ apps/zero-throughput/package.json | 6 +- apps/zero-throughput/src/analyze.ts | 21 ++- apps/zero-throughput/src/client.ts | 2 +- apps/zero-throughput/src/config.ts | 19 +++ apps/zero-throughput/src/main.ts | 15 +- apps/zero-throughput/src/permissions.ts | 62 -------- apps/zero-throughput/src/processes.ts | 90 ++++++++--- apps/zero-throughput/src/profile-queries.ts | 151 ++++++------------- apps/zero-throughput/src/queries.ts | 158 ++++++++++++++++++++ apps/zero-throughput/src/schema.ts | 4 +- apps/zero-throughput/src/server.ts | 31 ++++ apps/zero-throughput/tsconfig.json | 2 +- pnpm-lock.yaml | 6 + 14 files changed, 498 insertions(+), 202 deletions(-) create mode 100644 apps/zero-throughput/api/index.ts delete mode 100644 apps/zero-throughput/src/permissions.ts create mode 100644 apps/zero-throughput/src/queries.ts create mode 100644 apps/zero-throughput/src/server.ts diff --git a/apps/zero-throughput/api/index.ts b/apps/zero-throughput/api/index.ts new file mode 100644 index 0000000000..cca8b10f34 --- /dev/null +++ b/apps/zero-throughput/api/index.ts @@ -0,0 +1,133 @@ +import {mustGetQuery, type ReadonlyJSONValue} from '@rocicorp/zero'; +import { + handleMutateRequest, + handleQueryRequest, + type Database, + type QueryRequestHandler, + type TransactionProviderHooks, +} from '@rocicorp/zero/server'; +import Fastify, {type FastifyReply, type FastifyRequest} from 'fastify'; +import {queries} from '../src/queries.ts'; +import {schema} from '../src/schema.ts'; + +export const fastify = Fastify({ + logger: process.env.NODE_ENV !== 'test', +}); + +const dummyDbProvider: Database = { + transaction: ( + callback: ( + tx: unknown, + transactionHooks: TransactionProviderHooks, + ) => Promise | R, + ): Promise => + Promise.resolve( + callback( + {}, + { + updateClientMutationID: () => Promise.resolve({lastMutationID: 0}), + writeMutationResult: () => Promise.resolve(), + deleteMutationResults: () => Promise.resolve(), + }, + ), + ), +}; + +fastify.get('/health', (_req, reply) => { + reply.send({status: 'ok'}); +}); + +fastify.get('/', (_req, reply) => { + reply.send({status: 'ok', service: 'zero-throughput-api'}); +}); + +fastify.post<{ + Querystring: Record; + Body: ReadonlyJSONValue; +}>('/api/push', mutateHandler); + +fastify.post<{ + Querystring: Record; + Body: ReadonlyJSONValue; +}>('/api/mutate', mutateHandler); + +function extractUserID( + headers: Record, + query: Record, +): string | undefined { + const authHeader = headers['authorization']; + if (typeof authHeader === 'string') { + if (authHeader.startsWith('Bearer ')) { + return authHeader.slice('Bearer '.length).trim(); + } + return authHeader.trim(); + } + return ( + (headers['x-user-id'] as string | undefined) ?? query.userID ?? undefined + ); +} + +async function mutateHandler( + request: FastifyRequest<{ + Querystring: Record; + Body: ReadonlyJSONValue; + }>, + reply: FastifyReply, +) { + const authUserID = extractUserID(request.headers, request.query); + + const response = await handleMutateRequest>({ + dbProvider: dummyDbProvider, + handler: transact => transact(() => Promise.resolve()), + query: request.query, + body: request.body, + userID: authUserID, + logLevel: 'info', + }); + reply.send(response); +} + +fastify.post<{ + Querystring: Record; + Body: ReadonlyJSONValue; +}>('/api/get-queries', queryHandler); + +fastify.post<{ + Querystring: Record; + Body: ReadonlyJSONValue; +}>('/api/query', queryHandler); + +type AnyQuery = ReturnType; + +const queryTransformHandler: QueryRequestHandler = (name, args) => { + const query = mustGetQuery(queries, name); + return query.fn({args, ctx: undefined}) as unknown as AnyQuery; +}; + +async function queryHandler( + request: FastifyRequest<{ + Querystring: Record; + Body: ReadonlyJSONValue; + }>, + reply: FastifyReply, +) { + const authUserID = extractUserID(request.headers, request.query); + + const response = await handleQueryRequest({ + handler: queryTransformHandler, + schema, + query: request.query, + body: request.body, + userID: authUserID, + logLevel: 'info', + }); + reply.send(response); +} + +export default async function handler( + req: FastifyRequest, + reply: FastifyReply, +) { + await fastify.ready(); + fastify.server.emit('request', req, reply); +} diff --git a/apps/zero-throughput/package.json b/apps/zero-throughput/package.json index 731316b71b..df92af40af 100644 --- a/apps/zero-throughput/package.json +++ b/apps/zero-throughput/package.json @@ -12,6 +12,8 @@ "sweep:num-view-syncers": "node src/linear-sweep.ts --topology distributed --num-view-syncers 1,2,3 --sync-workers 2 --write-rates 300 --users 12 --profiles forum --duration-ms 15000", "sweep:users": "node src/linear-sweep.ts --users 5,10,20,50 --sync-workers 2 --write-rates 200 --profiles forum --duration-ms 15000", "go": "node src/main.ts", + "server": "node src/server.ts", + "start:server": "node src/server.ts", "db-up": "cd docker && docker compose up", "db-down": "cd docker && docker compose down", "check-types": "tsc", @@ -24,8 +26,10 @@ "dependencies": { "@dotenvx/dotenvx": "^1.39.0", "@rocicorp/zero": "workspace:*", + "fastify": "^5.0.0", "postgres": "3.4.7", - "ws": "^8.18.1" + "ws": "^8.18.1", + "zod": "^4.1.11" }, "devDependencies": { "@types/node": "^22.10.5", diff --git a/apps/zero-throughput/src/analyze.ts b/apps/zero-throughput/src/analyze.ts index 28b5f15db5..810c709fb7 100644 --- a/apps/zero-throughput/src/analyze.ts +++ b/apps/zero-throughput/src/analyze.ts @@ -1,7 +1,6 @@ import {mapAST, type AST} from '../../../packages/zero-protocol/src/ast.ts'; import {clientToServer} from '../../../packages/zero-schema/src/name-mapper.ts'; import {runAnalyzeCLI} from '../../../packages/zero/src/analyze.ts'; -import {createBuilder} from '../../../packages/zql/src/query/create-builder.ts'; import type {BenchmarkModel, BenchmarkProfile} from './config.ts'; import { buildProfileQuery, @@ -53,7 +52,6 @@ if (config.help) { } else { const {profile, queryIndex} = resolveProfileQuery(config); const {name, query} = buildProfileQuery( - createBuilder(schema), profile, config.model, queryIndex, @@ -176,10 +174,23 @@ function parseArgs(argv: readonly string[]): AnalyzeConfig { } function queryAST(query: unknown): AST { - if (query === null || typeof query !== 'object' || !('ast' in query)) { - throw new Error('Profile query did not expose an AST'); + if (query !== null && typeof query === 'object') { + if ('ast' in query) { + return (query as {readonly ast: AST}).ast; + } + if ( + 'query' in query && + typeof (query as {query: unknown}).query === 'function' && + 'fn' in (query as {query: {fn?: unknown}}).query + ) { + const q = query as { + args: unknown; + query: {fn: (opts: {args: unknown; ctx: unknown}) => {ast: AST}}; + }; + return q.query.fn({args: q.args, ctx: undefined}).ast; + } } - return (query as {readonly ast: AST}).ast; + throw new Error('Profile query did not expose an AST'); } function parseOption(arg: string): { diff --git a/apps/zero-throughput/src/client.ts b/apps/zero-throughput/src/client.ts index f6483d6294..5269622e2c 100644 --- a/apps/zero-throughput/src/client.ts +++ b/apps/zero-throughput/src/client.ts @@ -62,6 +62,7 @@ export class SyntheticClient { schema, cacheURL: targetCacheURL, userID, + auth: userID, storageKey: `${config.runID}-${clientIndex}`, kvStore: 'mem', logLevel: 'error', @@ -81,7 +82,6 @@ export class SyntheticClient { #registerProfileQuery(config: BenchmarkConfig, queryIndex: number): void { const {name, query} = buildProfileQuery( - this.#zero.query, config.profile, config.model, queryIndex, diff --git a/apps/zero-throughput/src/config.ts b/apps/zero-throughput/src/config.ts index d8328fde0f..5439b4b995 100644 --- a/apps/zero-throughput/src/config.ts +++ b/apps/zero-throughput/src/config.ts @@ -39,6 +39,9 @@ const options = { reset: v.boolean().default(true), cacheURL: v.string().optional(), cacheURLs: v.string().optional(), + appServerPort: v.number().default(3_000), + queryURL: v.string().optional(), + mutateURL: v.string().optional(), topology: v.literalUnion('single', 'distributed').default('single'), numViewSyncers: v.number().default(1), @@ -97,6 +100,9 @@ export type BenchmarkConfig = { readonly profileVS: boolean; readonly processLogMode: 'file' | 'inherit' | 'ignore'; readonly reset: boolean; + readonly appServerPort: number; + readonly queryURL: string | undefined; + readonly mutateURL: string | undefined; readonly cacheURL: string; readonly cacheURLs: readonly string[]; readonly pg: { @@ -196,6 +202,19 @@ export function loadConfig(): BenchmarkConfig { profileVS: parsed.profileVS, processLogMode: parsed.processLogMode, reset: parsed.reset, + appServerPort: parsed.appServerPort, + queryURL: + parsed.queryURL ?? + process.env.ZERO_QUERY_URL ?? + (isZeroManaged + ? `http://127.0.0.1:${parsed.appServerPort}/api/query` + : undefined), + mutateURL: + parsed.mutateURL ?? + process.env.ZERO_MUTATE_URL ?? + (isZeroManaged + ? `http://127.0.0.1:${parsed.appServerPort}/api/mutate` + : undefined), cacheURL: cacheURLs[0], cacheURLs, pg: { diff --git a/apps/zero-throughput/src/main.ts b/apps/zero-throughput/src/main.ts index d26aa88c81..3d75c41c88 100644 --- a/apps/zero-throughput/src/main.ts +++ b/apps/zero-throughput/src/main.ts @@ -9,12 +9,13 @@ import { import {OTelMetricsCollector} from './metrics.ts'; import { analyzeProfileQueries, - deployPermissions, queryPlanAnalysisLogPath, removeReplicaFiles, + startAppServer, startPostgres, startZeroTopology, stopPostgres, + waitForAppServer, waitForZeroCache, type ProcessCommand, } from './processes.ts'; @@ -76,8 +77,16 @@ async function main(): Promise { await removeReplicaFiles(config.zero.replicaFile); } - log('Deploying benchmark permissions...'); - processes.push(await deployPermissions(config)); + if (config.zero.start) { + log(`Starting app server on port ${config.appServerPort}...`); + const appServer = startAppServer(config); + processes.push(appServer); + cleanup.push(() => appServer.stop()); + if (appServer.logPath !== undefined) { + log(`app-server logs: ${appServer.logPath}`); + } + await waitForAppServer(config.appServerPort, 30_000, appServer); + } if (config.zero.start) { log( diff --git a/apps/zero-throughput/src/permissions.ts b/apps/zero-throughput/src/permissions.ts deleted file mode 100644 index 548f52d843..0000000000 --- a/apps/zero-throughput/src/permissions.ts +++ /dev/null @@ -1,62 +0,0 @@ -import {ANYONE_CAN, definePermissions} from '@rocicorp/zero'; -import {schema} from './schema.ts'; - -export {schema}; - -export const permissions = await definePermissions(schema, () => ({ - event: { - row: { - select: ANYONE_CAN, - }, - }, - emailThread: { - row: { - select: ANYONE_CAN, - }, - }, - emailMessage: { - row: { - select: ANYONE_CAN, - }, - }, - forumUser: { - row: { - select: ANYONE_CAN, - }, - }, - forumCategory: { - row: { - select: ANYONE_CAN, - }, - }, - forumThread: { - row: { - select: ANYONE_CAN, - }, - }, - forumPost: { - row: { - select: ANYONE_CAN, - }, - }, - relOrg: { - row: { - select: ANYONE_CAN, - }, - }, - relAccount: { - row: { - select: ANYONE_CAN, - }, - }, - relContact: { - row: { - select: ANYONE_CAN, - }, - }, - relActivity: { - row: { - select: ANYONE_CAN, - }, - }, -})); diff --git a/apps/zero-throughput/src/processes.ts b/apps/zero-throughput/src/processes.ts index b477619fbc..b7f6a692c7 100644 --- a/apps/zero-throughput/src/processes.ts +++ b/apps/zero-throughput/src/processes.ts @@ -32,35 +32,75 @@ export type ManagedProcess = ProcessCommand & { stop(): Promise; }; -export async function deployPermissions( - config: BenchmarkConfig, -): Promise { - const deployPermissionsMain = fileURLToPath( - new URL( - '../../../packages/zero-cache/src/scripts/deploy-permissions.ts', - import.meta.url, - ), - ); - const schemaPath = join(appRoot, 'src/permissions.ts'); +export function startAppServer(config: BenchmarkConfig): ManagedProcess { + const serverMain = fileURLToPath(new URL('server.ts', import.meta.url)); const command = [ process.execPath, - deployPermissionsMain, - '--schema-path', - schemaPath, - '--upstream-db', - config.pg.url, - '--app-id', - config.zero.appID, - '--force', + serverMain, + `--port=${config.appServerPort}`, ]; - await runCommand(command[0], command.slice(1), repoRoot); + const logPath = + config.processLogMode === 'file' + ? join(processLogsDir(config), `${config.runID}-app-server.log`) + : undefined; + const logStream = + logPath === undefined ? undefined : createWriteStream(logPath); + const child = spawn(command[0], command.slice(1), { + cwd: appRoot, + env: { + ...process.env, + NODE_ENV: 'development', + }, + stdio: + config.processLogMode === 'inherit' + ? 'inherit' + : [ + 'ignore', + config.processLogMode === 'file' ? 'pipe' : 'ignore', + config.processLogMode === 'file' ? 'pipe' : 'ignore', + ], + }); + pipeProcessLogs(child, logStream); + return { - name: 'zero-deploy-permissions', + name: 'app-server', command, - cwd: repoRoot, + cwd: appRoot, + logPath, + child, + stop: () => stopChild(child, 'SIGTERM'), }; } +export async function waitForAppServer( + port: number, + timeoutMs: number, + processContext?: ProcessCommand, +): Promise { + const url = `http://127.0.0.1:${port}/health`; + const deadline = Date.now() + timeoutMs; + let lastError: unknown; + while (Date.now() < deadline) { + try { + const response = await fetch(url); + if (response.ok) { + return; + } + lastError = new Error(`Health returned status ${response.status}`); + } catch (error) { + lastError = error; + } + await sleep(200); + } + const procDetails = + processContext?.logPath !== undefined + ? ` (check logs: ${processContext.logPath})` + : ''; + throw new Error( + `Timed out waiting for app-server at ${url} after ${timeoutMs}ms: ${String(lastError)}${procDetails}`, + ); +} + export async function startPostgres(): Promise { const cwd = join(appRoot, 'docker'); await runCommand('docker', ['compose', 'up', '-d', 'postgres'], cwd); @@ -322,9 +362,15 @@ function spawnZeroProcess(args: { ZERO_CHANGE_MAX_CONNS: String(config.zero.changeMaxConns), ZERO_LOG_LEVEL: config.zero.logLevel, ZERO_LOG_FORMAT: 'text', - ZERO_ALLOW_LEGACY_QUERIES: 'true', }; + if (config.queryURL) { + env.ZERO_QUERY_URL = config.queryURL; + } + if (config.mutateURL) { + env.ZERO_MUTATE_URL = config.mutateURL; + } + if (args.changeStreamerURI) { env.ZERO_CHANGE_STREAMER_URI = args.changeStreamerURI; } diff --git a/apps/zero-throughput/src/profile-queries.ts b/apps/zero-throughput/src/profile-queries.ts index 78ebd563a0..1f895f1006 100644 --- a/apps/zero-throughput/src/profile-queries.ts +++ b/apps/zero-throughput/src/profile-queries.ts @@ -1,5 +1,6 @@ -import type {Query, SchemaQuery} from '@rocicorp/zero'; +import type {Query} from '@rocicorp/zero'; import type {BenchmarkModel, BenchmarkProfile} from './config.ts'; +import {queries} from './queries.ts'; import type {schema} from './schema.ts'; import { emailOwnerIDForClient, @@ -37,7 +38,6 @@ export const PROFILE_QUERY_NAMES = { } as const satisfies Record; export function buildProfileQuery( - builder: SchemaQuery, profile: BenchmarkProfile, model: BenchmarkModel, queryIndex: number, @@ -48,38 +48,20 @@ export function buildProfileQuery( case 'feed-append': return { name: profileQueryName(profile, queryIndex), - query: builder.event - .where('bucket', feedBucketForClient(model, clientIndex)) - .orderBy('seq', 'desc') - .limit(rowsPerQuery) as ThroughputQuery, + query: queries.feedRecentEvents({ + bucket: feedBucketForClient(model, clientIndex), + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; case 'email': - return buildEmailQuery( - builder, - model, - queryIndex, - rowsPerQuery, - clientIndex, - ); + return buildEmailQuery(model, queryIndex, rowsPerQuery, clientIndex); case 'forum': - return buildForumQuery( - builder, - model, - queryIndex, - rowsPerQuery, - clientIndex, - ); + return buildForumQuery(model, queryIndex, rowsPerQuery, clientIndex); case 'relational': - return buildRelationalQuery( - builder, - model, - queryIndex, - rowsPerQuery, - clientIndex, - ); + return buildRelationalQuery(model, queryIndex, rowsPerQuery, clientIndex); } } @@ -131,7 +113,6 @@ export function findProfileQuery( } function buildEmailQuery( - builder: SchemaQuery, model: BenchmarkModel, queryIndex: number, rowsPerQuery: number, @@ -143,43 +124,34 @@ function buildEmailQuery( case 0: return { name, - query: builder.emailThread - .where('ownerID', ownerID) - .where('mailbox', 'inbox') - .related('messages', q => q.orderBy('seq', 'desc').limit(5)) - .orderBy('seq', 'desc') - .limit(rowsPerQuery) as ThroughputQuery, + query: queries.emailThreadListWithMessages({ + ownerID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; case 1: return { name, - query: builder.emailMessage - .where('ownerID', ownerID) - .where('mailbox', 'inbox') - .related('thread') - .orderBy('seq', 'desc') - .limit(rowsPerQuery) as ThroughputQuery, + query: queries.emailMessageListWithThread({ + ownerID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; case 2: return { name, - query: builder.emailThread - .where('ownerID', ownerID) - .where('mailbox', 'inbox') - .related('messages', q => - q.where('unread', true).orderBy('seq', 'desc').limit(10), - ) - .orderBy('seq', 'desc') - .limit(rowsPerQuery) as ThroughputQuery, + query: queries.emailUnreadThreadList({ + ownerID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; } throw new Error(`Invalid email query index: ${queryIndex}`); } function buildForumQuery( - builder: SchemaQuery, model: BenchmarkModel, queryIndex: number, rowsPerQuery: number, @@ -191,49 +163,34 @@ function buildForumQuery( case 0: return { name, - query: builder.forumCategory - .where('id', categoryID) - .related('threads', q => - q - .orderBy('seq', 'desc') - .limit(rowsPerQuery) - .related('author') - .related('posts', p => - p.orderBy('seq', 'desc').limit(3).related('author'), - ), - ) as ThroughputQuery, + query: queries.forumCategoryThreadTree({ + categoryID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; case 1: return { name, - query: builder.forumThread - .where('categoryID', categoryID) - .related('category') - .related('author') - .related('posts', q => - q.orderBy('seq', 'desc').limit(5).related('author'), - ) - .orderBy('seq', 'desc') - .limit(rowsPerQuery) as ThroughputQuery, + query: queries.forumThreadListWithPosts({ + categoryID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; case 2: return { name, - query: builder.forumPost - .where('categoryID', categoryID) - .related('thread', q => q.related('author').related('category')) - .related('author') - .orderBy('seq', 'desc') - .limit(rowsPerQuery) as ThroughputQuery, + query: queries.forumPostListWithThread({ + categoryID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; } throw new Error(`Invalid forum query index: ${queryIndex}`); } function buildRelationalQuery( - builder: SchemaQuery, model: BenchmarkModel, queryIndex: number, rowsPerQuery: number, @@ -245,46 +202,28 @@ function buildRelationalQuery( case 0: return { name, - query: builder.relOrg - .where('id', orgID) - .related('accounts', q => - q - .orderBy('seq', 'desc') - .limit(rowsPerQuery) - .related('contacts') - .related('activities', a => - a.orderBy('seq', 'desc').limit(5).related('contact'), - ), - ) - .related('activities', q => - q.orderBy('seq', 'desc').limit(rowsPerQuery), - ) as ThroughputQuery, + query: queries.relationalOrgAccountTree({ + orgID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; case 1: return { name, - query: builder.relAccount - .where('orgID', orgID) - .related('org') - .related('contacts') - .related('activities', q => - q.orderBy('seq', 'desc').limit(5).related('contact'), - ) - .orderBy('seq', 'desc') - .limit(rowsPerQuery) as ThroughputQuery, + query: queries.relationalAccountList({ + orgID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; case 2: return { name, - query: builder.relActivity - .where('orgID', orgID) - .related('org') - .related('account', q => q.related('contacts')) - .related('contact') - .orderBy('seq', 'desc') - .limit(rowsPerQuery) as ThroughputQuery, + query: queries.relationalActivityList({ + orgID, + limit: rowsPerQuery, + }) as unknown as ThroughputQuery, }; } throw new Error(`Invalid relational query index: ${queryIndex}`); diff --git a/apps/zero-throughput/src/queries.ts b/apps/zero-throughput/src/queries.ts new file mode 100644 index 0000000000..eed0fc24a5 --- /dev/null +++ b/apps/zero-throughput/src/queries.ts @@ -0,0 +1,158 @@ +import {defineQueries, defineQuery} from '@rocicorp/zero'; +import * as z from 'zod/mini'; +import {builder} from './schema.ts'; + +export const queries = defineQueries({ + feedRecentEvents: defineQuery( + z.object({ + bucket: z.number(), + limit: z.number(), + }), + ({args: {bucket, limit}}) => + builder.event.where('bucket', bucket).orderBy('seq', 'desc').limit(limit), + ), + + emailThreadListWithMessages: defineQuery( + z.object({ + ownerID: z.string(), + limit: z.number(), + }), + ({args: {ownerID, limit}}) => + builder.emailThread + .where('ownerID', ownerID) + .where('mailbox', 'inbox') + .related('messages', q => q.orderBy('seq', 'desc').limit(5)) + .orderBy('seq', 'desc') + .limit(limit), + ), + + emailMessageListWithThread: defineQuery( + z.object({ + ownerID: z.string(), + limit: z.number(), + }), + ({args: {ownerID, limit}}) => + builder.emailMessage + .where('ownerID', ownerID) + .where('mailbox', 'inbox') + .related('thread') + .orderBy('seq', 'desc') + .limit(limit), + ), + + emailUnreadThreadList: defineQuery( + z.object({ + ownerID: z.string(), + limit: z.number(), + }), + ({args: {ownerID, limit}}) => + builder.emailThread + .where('ownerID', ownerID) + .where('mailbox', 'inbox') + .related('messages', q => + q.where('unread', true).orderBy('seq', 'desc').limit(10), + ) + .orderBy('seq', 'desc') + .limit(limit), + ), + + forumCategoryThreadTree: defineQuery( + z.object({ + categoryID: z.string(), + limit: z.number(), + }), + ({args: {categoryID, limit}}) => + builder.forumCategory.where('id', categoryID).related('threads', q => + q + .orderBy('seq', 'desc') + .limit(limit) + .related('author') + .related('posts', p => + p.orderBy('seq', 'desc').limit(3).related('author'), + ), + ), + ), + + forumThreadListWithPosts: defineQuery( + z.object({ + categoryID: z.string(), + limit: z.number(), + }), + ({args: {categoryID, limit}}) => + builder.forumThread + .where('categoryID', categoryID) + .related('category') + .related('author') + .related('posts', q => + q.orderBy('seq', 'desc').limit(5).related('author'), + ) + .orderBy('seq', 'desc') + .limit(limit), + ), + + forumPostListWithThread: defineQuery( + z.object({ + categoryID: z.string(), + limit: z.number(), + }), + ({args: {categoryID, limit}}) => + builder.forumPost + .where('categoryID', categoryID) + .related('thread', q => q.related('author').related('category')) + .related('author') + .orderBy('seq', 'desc') + .limit(limit), + ), + + relationalOrgAccountTree: defineQuery( + z.object({ + orgID: z.string(), + limit: z.number(), + }), + ({args: {orgID, limit}}) => + builder.relOrg + .where('id', orgID) + .related('accounts', q => + q + .orderBy('seq', 'desc') + .limit(limit) + .related('contacts') + .related('activities', a => + a.orderBy('seq', 'desc').limit(5).related('contact'), + ), + ) + .related('activities', q => q.orderBy('seq', 'desc').limit(limit)), + ), + + relationalAccountList: defineQuery( + z.object({ + orgID: z.string(), + limit: z.number(), + }), + ({args: {orgID, limit}}) => + builder.relAccount + .where('orgID', orgID) + .related('org') + .related('contacts') + .related('activities', q => + q.orderBy('seq', 'desc').limit(5).related('contact'), + ) + .orderBy('seq', 'desc') + .limit(limit), + ), + + relationalActivityList: defineQuery( + z.object({ + orgID: z.string(), + limit: z.number(), + }), + ({args: {orgID, limit}}) => + builder.relActivity + .where('orgID', orgID) + .related('org') + .related('account', q => q.related('contacts')) + .related('contact') + .orderBy('seq', 'desc') + .limit(limit), + ), +}); diff --git a/apps/zero-throughput/src/schema.ts b/apps/zero-throughput/src/schema.ts index e8a661c202..36c2196781 100644 --- a/apps/zero-throughput/src/schema.ts +++ b/apps/zero-throughput/src/schema.ts @@ -1,5 +1,6 @@ import { boolean, + createBuilder, createSchema, json, number, @@ -318,5 +319,6 @@ export const schema = createSchema({ relContactRelationships, relActivityRelationships, ], - enableLegacyQueries: true, }); + +export const builder = createBuilder(schema); diff --git a/apps/zero-throughput/src/server.ts b/apps/zero-throughput/src/server.ts new file mode 100644 index 0000000000..b2f2c6e697 --- /dev/null +++ b/apps/zero-throughput/src/server.ts @@ -0,0 +1,31 @@ +import {fastify} from '../api/index.ts'; + +const port = Number( + process.env.PORT ?? parsePortArg(process.argv.slice(2)) ?? 3000, +); +const host = process.env.HOST ?? '0.0.0.0'; + +function parsePortArg(args: readonly string[]): number | undefined { + for (let i = 0; i < args.length; i++) { + const arg = args[i]; + if (arg === '--port' && i + 1 < args.length) { + return Number(args[i + 1]); + } + if (arg.startsWith('--port=')) { + return Number(arg.slice('--port='.length)); + } + } + return undefined; +} + +async function main() { + try { + const address = await fastify.listen({port, host}); + fastify.log.info(`zero-throughput API server listening on ${address}`); + } catch (err) { + fastify.log.error(err); + process.exit(1); + } +} + +void main(); diff --git a/apps/zero-throughput/tsconfig.json b/apps/zero-throughput/tsconfig.json index d024983ea2..f0bdbc08f9 100644 --- a/apps/zero-throughput/tsconfig.json +++ b/apps/zero-throughput/tsconfig.json @@ -8,5 +8,5 @@ }, "types": ["node"] }, - "include": ["src"] + "include": ["src", "api"] } diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 258bd0cf31..58376d4460 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -292,12 +292,18 @@ importers: '@rocicorp/zero': specifier: workspace:* version: link:../../packages/zero + fastify: + specifier: ^5.0.0 + version: 5.8.5 postgres: specifier: 3.4.7 version: 3.4.7(patch_hash=2bd119d361250d5ee010b23af4c98504e9c6861b692d2e2ae23b9a9e5de4dbce) ws: specifier: ^8.18.1 version: 8.20.1 + zod: + specifier: ^4.1.11 + version: 4.4.3 devDependencies: '@types/node': specifier: ^22.10.5 From 7fea3892626a560574a073a1a03d0e314df86a32 Mon Sep 17 00:00:00 2001 From: Greg Baker Date: Fri, 4 Sep 2026 18:10:35 -0700 Subject: [PATCH 2/5] work --- apps/zero-throughput/src/config.ts | 16 ---------------- apps/zero-throughput/src/processes.ts | 9 ++------- 2 files changed, 2 insertions(+), 23 deletions(-) diff --git a/apps/zero-throughput/src/config.ts b/apps/zero-throughput/src/config.ts index 5439b4b995..2267710bc5 100644 --- a/apps/zero-throughput/src/config.ts +++ b/apps/zero-throughput/src/config.ts @@ -40,8 +40,6 @@ const options = { cacheURL: v.string().optional(), cacheURLs: v.string().optional(), appServerPort: v.number().default(3_000), - queryURL: v.string().optional(), - mutateURL: v.string().optional(), topology: v.literalUnion('single', 'distributed').default('single'), numViewSyncers: v.number().default(1), @@ -101,8 +99,6 @@ export type BenchmarkConfig = { readonly processLogMode: 'file' | 'inherit' | 'ignore'; readonly reset: boolean; readonly appServerPort: number; - readonly queryURL: string | undefined; - readonly mutateURL: string | undefined; readonly cacheURL: string; readonly cacheURLs: readonly string[]; readonly pg: { @@ -203,18 +199,6 @@ export function loadConfig(): BenchmarkConfig { processLogMode: parsed.processLogMode, reset: parsed.reset, appServerPort: parsed.appServerPort, - queryURL: - parsed.queryURL ?? - process.env.ZERO_QUERY_URL ?? - (isZeroManaged - ? `http://127.0.0.1:${parsed.appServerPort}/api/query` - : undefined), - mutateURL: - parsed.mutateURL ?? - process.env.ZERO_MUTATE_URL ?? - (isZeroManaged - ? `http://127.0.0.1:${parsed.appServerPort}/api/mutate` - : undefined), cacheURL: cacheURLs[0], cacheURLs, pg: { diff --git a/apps/zero-throughput/src/processes.ts b/apps/zero-throughput/src/processes.ts index b7f6a692c7..7ec3fab33f 100644 --- a/apps/zero-throughput/src/processes.ts +++ b/apps/zero-throughput/src/processes.ts @@ -362,15 +362,10 @@ function spawnZeroProcess(args: { ZERO_CHANGE_MAX_CONNS: String(config.zero.changeMaxConns), ZERO_LOG_LEVEL: config.zero.logLevel, ZERO_LOG_FORMAT: 'text', + ZERO_QUERY_URL: `http://127.0.0.1:${config.appServerPort}/api/query`, + ZERO_MUTATE_URL: `http://127.0.0.1:${config.appServerPort}/api/mutate`, }; - if (config.queryURL) { - env.ZERO_QUERY_URL = config.queryURL; - } - if (config.mutateURL) { - env.ZERO_MUTATE_URL = config.mutateURL; - } - if (args.changeStreamerURI) { env.ZERO_CHANGE_STREAMER_URI = args.changeStreamerURI; } From 20bd0143f88e4e71d3f5aac3b599e89a0ce2b082 Mon Sep 17 00:00:00 2001 From: Greg Baker Date: Fri, 4 Sep 2026 18:12:46 -0700 Subject: [PATCH 3/5] docs(zero-throughput): update README for query-transform app server --- apps/zero-throughput/README.md | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) diff --git a/apps/zero-throughput/README.md b/apps/zero-throughput/README.md index 821823685b..616acd6a78 100644 --- a/apps/zero-throughput/README.md +++ b/apps/zero-throughput/README.md @@ -6,7 +6,7 @@ The default run: 1. Starts a dedicated PostgreSQL 16 Docker container on port `6436`. 2. Resets the benchmark table and Zero metadata for app id `zero_throughput`. -3. Deploys allow-read permissions for the benchmark table. +3. Starts the Fastify app server on port `3000` for query transformations. 4. Starts `zero-cache` on port `4848`. 5. Runs analyze-query for each distinct live query shape in the selected profile. 6. Starts synthetic Zero clients with live queries for the selected profile. @@ -294,11 +294,12 @@ pnpm --filter zero-throughput run sweep:write-rates -- \ ### Options Reference -| CLI Option | Environment Variable | Default | Description | -| :-------------------- | :--------------------------- | :---------------------- | :----------------------------------------------------------------- | -| `--cache-url ` | `ZERO_THROUGHPUT_CACHE_URL` | `http://127.0.0.1:4848` | Primary Zero cache endpoint or load balancer (disables local Zero) | -| `--cache-urls ` | `ZERO_THROUGHPUT_CACHE_URLS` | `undefined` | Comma-separated View-Syncer URLs for client partitioning | -| `--pg-url ` | `ZERO_THROUGHPUT_PG_URL` | `postgresql://...:6436` | Upstream database connection string (disables local Postgres) | -| `--reset ` | `ZERO_THROUGHPUT_RESET` | `true` | When `false`, skips dropping/resetting the benchmark table | +| CLI Option | Environment Variable | Default | Description | +| :------------------------ | :-------------------------------- | :---------------------- | :----------------------------------------------------------------- | +| `--cache-url ` | `ZERO_THROUGHPUT_CACHE_URL` | `http://127.0.0.1:4848` | Primary Zero cache endpoint or load balancer (disables local Zero) | +| `--cache-urls ` | `ZERO_THROUGHPUT_CACHE_URLS` | `undefined` | Comma-separated View-Syncer URLs for client partitioning | +| `--pg-url ` | `ZERO_THROUGHPUT_PG_URL` | `postgresql://...:6436` | Upstream database connection string (disables local Postgres) | +| `--app-server-port ` | `ZERO_THROUGHPUT_APP_SERVER_PORT` | `3000` | Local query-transform app server port | +| `--reset ` | `ZERO_THROUGHPUT_RESET` | `true` | When `false`, skips dropping/resetting the benchmark table | Run `pnpm --filter zero-throughput start -- --help` for all options. From 1fac7c266a05e1c84b2506d799780d51d864250f Mon Sep 17 00:00:00 2001 From: Greg Baker Date: Fri, 4 Sep 2026 18:17:12 -0700 Subject: [PATCH 4/5] refactor(zero-throughput): simplify mutate endpoint to 501 response --- apps/zero-throughput/api/index.ts | 77 +++++++------------------------ 1 file changed, 16 insertions(+), 61 deletions(-) diff --git a/apps/zero-throughput/api/index.ts b/apps/zero-throughput/api/index.ts index cca8b10f34..21c9409f84 100644 --- a/apps/zero-throughput/api/index.ts +++ b/apps/zero-throughput/api/index.ts @@ -1,10 +1,7 @@ import {mustGetQuery, type ReadonlyJSONValue} from '@rocicorp/zero'; import { - handleMutateRequest, handleQueryRequest, - type Database, type QueryRequestHandler, - type TransactionProviderHooks, } from '@rocicorp/zero/server'; import Fastify, {type FastifyReply, type FastifyRequest} from 'fastify'; import {queries} from '../src/queries.ts'; @@ -14,25 +11,6 @@ export const fastify = Fastify({ logger: process.env.NODE_ENV !== 'test', }); -const dummyDbProvider: Database = { - transaction: ( - callback: ( - tx: unknown, - transactionHooks: TransactionProviderHooks, - ) => Promise | R, - ): Promise => - Promise.resolve( - callback( - {}, - { - updateClientMutationID: () => Promise.resolve({lastMutationID: 0}), - writeMutationResult: () => Promise.resolve(), - deleteMutationResults: () => Promise.resolve(), - }, - ), - ), -}; - fastify.get('/health', (_req, reply) => { reply.send({status: 'ok'}); }); @@ -41,15 +19,29 @@ fastify.get('/', (_req, reply) => { reply.send({status: 'ok', service: 'zero-throughput-api'}); }); +fastify.post('/api/push', mutateHandler); +fastify.post('/api/mutate', mutateHandler); + +function mutateHandler(_request: FastifyRequest, reply: FastifyReply) { + reply.status(501).send({error: 'Mutations not supported'}); +} + fastify.post<{ Querystring: Record; Body: ReadonlyJSONValue; -}>('/api/push', mutateHandler); +}>('/api/get-queries', queryHandler); fastify.post<{ Querystring: Record; Body: ReadonlyJSONValue; -}>('/api/mutate', mutateHandler); +}>('/api/query', queryHandler); + +type AnyQuery = ReturnType; + +const queryTransformHandler: QueryRequestHandler = (name, args) => { + const query = mustGetQuery(queries, name); + return query.fn({args, ctx: undefined}) as unknown as AnyQuery; +}; function extractUserID( headers: Record, @@ -67,43 +59,6 @@ function extractUserID( ); } -async function mutateHandler( - request: FastifyRequest<{ - Querystring: Record; - Body: ReadonlyJSONValue; - }>, - reply: FastifyReply, -) { - const authUserID = extractUserID(request.headers, request.query); - - const response = await handleMutateRequest>({ - dbProvider: dummyDbProvider, - handler: transact => transact(() => Promise.resolve()), - query: request.query, - body: request.body, - userID: authUserID, - logLevel: 'info', - }); - reply.send(response); -} - -fastify.post<{ - Querystring: Record; - Body: ReadonlyJSONValue; -}>('/api/get-queries', queryHandler); - -fastify.post<{ - Querystring: Record; - Body: ReadonlyJSONValue; -}>('/api/query', queryHandler); - -type AnyQuery = ReturnType; - -const queryTransformHandler: QueryRequestHandler = (name, args) => { - const query = mustGetQuery(queries, name); - return query.fn({args, ctx: undefined}) as unknown as AnyQuery; -}; - async function queryHandler( request: FastifyRequest<{ Querystring: Record; From 41ba5a0ee321d8860c23a987b07b5421d4c5f4a9 Mon Sep 17 00:00:00 2001 From: Greg Baker Date: Fri, 4 Sep 2026 18:33:01 -0700 Subject: [PATCH 5/5] fix --- apps/zero-throughput/src/profile-queries.ts | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/apps/zero-throughput/src/profile-queries.ts b/apps/zero-throughput/src/profile-queries.ts index 1f895f1006..e052797801 100644 --- a/apps/zero-throughput/src/profile-queries.ts +++ b/apps/zero-throughput/src/profile-queries.ts @@ -1,4 +1,4 @@ -import type {Query} from '@rocicorp/zero'; +import type {Query} from '../../../packages/zql/src/query/query.ts'; import type {BenchmarkModel, BenchmarkProfile} from './config.ts'; import {queries} from './queries.ts'; import type {schema} from './schema.ts';