From 11b52c6e47f18cf0a7f7c3effe92d32a6d432e2d Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:09:49 +0000 Subject: [PATCH 01/12] fix(cubejs): cast to ClickHouse String type in Tesseract SQL Cube 1.7.30 ClickHouseQuery.sqlTemplates() never overrides templates.types.string, so the native planner CASTs multi-column primary keys for count measures without sql to the base 'STRING', which ClickHouse rejects (Unknown data type family: STRING). Extend the version-pinned postinstall patch to set it to 'String'. Co-Authored-By: Claude Opus 5.5 --- .../cubejs/scripts/patchCubeYamlCompiler.mjs | 64 +++++++++++++++---- .../src/__tests__/cube17Regression.test.js | 58 ++++++++++++++++- 2 files changed, 107 insertions(+), 15 deletions(-) diff --git a/services/cubejs/scripts/patchCubeYamlCompiler.mjs b/services/cubejs/scripts/patchCubeYamlCompiler.mjs index 132dde58..6bbc20e2 100644 --- a/services/cubejs/scripts/patchCubeYamlCompiler.mjs +++ b/services/cubejs/scripts/patchCubeYamlCompiler.mjs @@ -22,16 +22,47 @@ const PATCHED = ` else if (typeof obj === 'string') { let code = obj; if (!CubeValidator_1.nonStringFields.has(propertyPath[propertyPath.length - 1])) {`; +const CLICKHOUSE_ORIGINAL = ` templates.types.timestamp = 'DATETIME'; + delete templates.types.time;`; + +const CLICKHOUSE_PATCHED = ` templates.types.timestamp = 'DATETIME'; + // ClickHouse type names are case-sensitive. The base 'STRING' type is + // what Tesseract CASTs multi-column primary keys to for count measures + // without sql, and ClickHouse rejects it ("Unknown data type + // family: STRING"). + templates.types.string = 'String'; + delete templates.types.time;`; + const occurrences = (source, value) => source.split(value).length - 1; -export function patchCompilerSource(source) { - if (source.includes(PATCHED)) return { source, changed: false }; - if (occurrences(source, ORIGINAL) !== 1) { +function applyPatch(source, original, patched, label) { + if (source.includes(patched)) return { source, changed: false }; + if (occurrences(source, original) !== 1) { throw new Error( - "Refusing to patch Cube YAML compiler: expected source anchor was not found exactly once", + `Refusing to patch Cube ${label}: expected source anchor was not found exactly once`, ); } - return { source: source.replace(ORIGINAL, PATCHED), changed: true }; + return { source: source.replace(original, patched), changed: true }; +} + +export function patchCompilerSource(source) { + return applyPatch(source, ORIGINAL, PATCHED, "YAML compiler"); +} + +export function patchClickHouseQuerySource(source) { + return applyPatch( + source, + CLICKHOUSE_ORIGINAL, + CLICKHOUSE_PATCHED, + "ClickHouse query adapter", + ); +} + +async function patchFile(path, patch) { + const current = await readFile(path, "utf8"); + const result = patch(current); + if (result.changed) await writeFile(path, result.source, "utf8"); + return result.changed; } export async function patchInstalledCompiler() { @@ -43,14 +74,19 @@ export async function patchInstalledCompiler() { ); } - const compilerPath = resolve( - dirname(packageJsonPath), - "dist/src/compiler/YamlCompiler.js", + const root = dirname(packageJsonPath); + const compilerPath = resolve(root, "dist/src/compiler/YamlCompiler.js"); + const clickHousePath = resolve(root, "dist/src/adapter/ClickHouseQuery.js"); + const yamlChanged = await patchFile(compilerPath, patchCompilerSource); + const clickHouseChanged = await patchFile( + clickHousePath, + patchClickHouseQuerySource, ); - const current = await readFile(compilerPath, "utf8"); - const result = patchCompilerSource(current); - if (result.changed) await writeFile(compilerPath, result.source, "utf8"); - return { compilerPath, changed: result.changed }; + return { + compilerPath, + clickHousePath, + changed: yamlChanged || clickHouseChanged, + }; } if ( @@ -60,7 +96,7 @@ if ( const result = await patchInstalledCompiler(); console.log( result.changed - ? `Patched Cube ${SUPPORTED_VERSION} YAML metadata handling` - : `Cube ${SUPPORTED_VERSION} YAML metadata patch already applied`, + ? `Patched Cube ${SUPPORTED_VERSION} YAML metadata handling + ClickHouse string type` + : `Cube ${SUPPORTED_VERSION} YAML metadata + ClickHouse string type patches already applied`, ); } diff --git a/services/cubejs/src/__tests__/cube17Regression.test.js b/services/cubejs/src/__tests__/cube17Regression.test.js index 56990da6..ce9bf09a 100644 --- a/services/cubejs/src/__tests__/cube17Regression.test.js +++ b/services/cubejs/src/__tests__/cube17Regression.test.js @@ -6,7 +6,10 @@ import { prepareCompiler } from "@cubejs-backend/schema-compiler"; import { escapeCSVField } from "../utils/csvSerializer.js"; import { validateFormat } from "../utils/formatValidator.js"; -import { patchCompilerSource } from "../../scripts/patchCubeYamlCompiler.mjs"; +import { + patchClickHouseQuerySource, + patchCompilerSource, +} from "../../scripts/patchCubeYamlCompiler.mjs"; const require = createRequire(import.meta.url); const runtimeVersion = require("@cubejs-backend/server-core/package.json").version; @@ -102,6 +105,59 @@ after`; ); }); + it("guards the ClickHouse string-type patch against upstream source drift", () => { + const source = `before + templates.types.timestamp = 'DATETIME'; + delete templates.types.time; +after`; + const first = patchClickHouseQuerySource(source); + assert.equal(first.changed, true); + assert.match(first.source, /templates\.types\.string = 'String';/); + assert.deepEqual(patchClickHouseQuerySource(first.source), { + source: first.source, + changed: false, + }); + assert.throws( + () => patchClickHouseQuerySource("unexpected adapter source"), + /expected source anchor was not found exactly once/, + ); + }); + + it("casts composite-key count measures to ClickHouse String (Tesseract)", async () => { + const { ClickHouseQuery } = require( + "@cubejs-backend/schema-compiler/dist/src/adapter/ClickHouseQuery.js", + ); + const content = `cubes: + - name: points + sql_table: points + dimensions: + - name: series_gid + sql: series_gid + type: string + primary_key: true + - name: ts + sql: ts + type: time + primary_key: true + measures: + - name: count + type: count +`; + const compilers = prepareCompiler( + { dataSchemaFiles: async () => [{ fileName: "points.yml", content }] }, + { adapter: "clickhouse" }, + ); + await compilers.compiler.compile(); + const [sql] = new ClickHouseQuery(compilers, { + measures: ["points.count"], + timezone: "UTC", + useNativeSqlPlanner: true, + }).buildSqlAndParams(); + + assert.match(sql, / AS String\)/); + assert.doesNotMatch(sql, / AS STRING\)/); + }); + it("does not embed a comparison result in the test source", async () => { const source = await readFile(new URL(import.meta.url), "utf8"); assert.doesNotMatch(source, /pass(?:ed)?\s*[:=]\s*(?:true|yes)/i); From 7f7c27de6ca94b6b5d85b6a0ebb09a68f3d6e67b Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:33:48 +0000 Subject: [PATCH 02/12] fix(hasura): record restore source, pin writes to the current version - versions.source_version_id (nullable FK, ON DELETE SET NULL) records which version a rollback restored; insertable/selectable by role user. - dataschemas update/insert (role user) only touch rows of the current version, so history can no longer be rewritten in place. Nested insert_versions_one { dataschemas } saves still pass (verified). - versions insert requires owner/admin of the team (was: any member or the datasource owner, allowing empty versions to become current). - audit_logs.action also allows 'dataschema_update' (PUT dataschema). Co-Authored-By: Claude Opus 5.5 --- services/hasura/metadata/tables.yaml | 75 +++++++++++-------- .../down.sql | 6 ++ .../up.sql | 12 +++ 3 files changed, 61 insertions(+), 32 deletions(-) create mode 100644 services/hasura/migrations/1790522000000_version_history_integrity/down.sql create mode 100644 services/hasura/migrations/1790522000000_version_history_integrity/up.sql diff --git a/services/hasura/metadata/tables.yaml b/services/hasura/metadata/tables.yaml index 92ee7e60..1e04c531 100644 --- a/services/hasura/metadata/tables.yaml +++ b/services/hasura/metadata/tables.yaml @@ -552,17 +552,21 @@ - role: user permission: check: - datasource: - team: - members: - _and: - - member_roles: - team_role: - _in: - - owner - - admin - - user_id: - _eq: X-Hasura-User-Id + _and: + - datasource: + team: + members: + _and: + - member_roles: + team_role: + _in: + - owner + - admin + - user_id: + _eq: X-Hasura-User-Id + - version: + is_current: + _eq: true columns: - code - datasource_id @@ -602,17 +606,21 @@ - code - name filter: - datasource: - team: - members: - _and: - - member_roles: - team_role: - _in: - - owner - - admin - - user_id: - _eq: X-Hasura-User-Id + _and: + - datasource: + team: + members: + _and: + - member_roles: + team_role: + _in: + - owner + - admin + - user_id: + _eq: X-Hasura-User-Id + - version: + is_current: + _eq: true check: datasource: team: @@ -1707,21 +1715,23 @@ - role: user permission: check: - _or: - - branch: - datasource: - user_id: - _eq: X-Hasura-User-Id - - branch: - datasource: - team: - members: - user_id: + branch: + datasource: + team: + members: + _and: + - member_roles: + team_role: + _in: + - owner + - admin + - user_id: _eq: X-Hasura-User-Id columns: - branch_id - checksum - origin + - source_version_id - user_id select_permissions: - role: user @@ -1734,6 +1744,7 @@ - is_current - markdown_doc - origin + - source_version_id - updated_at - user_id filter: diff --git a/services/hasura/migrations/1790522000000_version_history_integrity/down.sql b/services/hasura/migrations/1790522000000_version_history_integrity/down.sql new file mode 100644 index 00000000..71ce04f4 --- /dev/null +++ b/services/hasura/migrations/1790522000000_version_history_integrity/down.sql @@ -0,0 +1,6 @@ +DELETE FROM public.audit_logs WHERE action = 'dataschema_update'; +ALTER TABLE public.audit_logs DROP CONSTRAINT IF EXISTS audit_logs_action_check; +ALTER TABLE public.audit_logs ADD CONSTRAINT audit_logs_action_check + CHECK (action IN ('dataschema_delete', 'version_rollback')); + +ALTER TABLE public.versions DROP COLUMN IF EXISTS source_version_id; diff --git a/services/hasura/migrations/1790522000000_version_history_integrity/up.sql b/services/hasura/migrations/1790522000000_version_history_integrity/up.sql new file mode 100644 index 00000000..487fcd76 --- /dev/null +++ b/services/hasura/migrations/1790522000000_version_history_integrity/up.sql @@ -0,0 +1,12 @@ +-- fix/version-history-integrity +-- 1. Record which version a rollback restored. Nullable: only origin='rollback' +-- rows carry it. SET NULL so pruning an old version never blocks on a restore +-- that referenced it. +ALTER TABLE public.versions + ADD COLUMN IF NOT EXISTS source_version_id uuid + REFERENCES public.versions(id) ON DELETE SET NULL; + +-- 2. PUT /api/v1/dataschema/:id audits under its own action. +ALTER TABLE public.audit_logs DROP CONSTRAINT IF EXISTS audit_logs_action_check; +ALTER TABLE public.audit_logs ADD CONSTRAINT audit_logs_action_check + CHECK (action IN ('dataschema_delete', 'version_rollback', 'dataschema_update')); From 9e89ac5ec87cd1378d70b96d5f38c3c19327df7f Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:34:15 +0000 Subject: [PATCH 03/12] fix(cubejs): rollback records the version it restored Extract commitVersionFiles() from rollbackVersion(): one helper that writes a new version holding exactly the given files with the caller's Hasura token, returns the new rows, and drops the caller's cached user scope so its next request resolves the new version. Rollback now sets source_version_id = toVersionId, and the rollback audit payload carries it. Co-Authored-By: Claude Opus 5.5 --- .../actions/src/rpc/auditVersionRollback.js | 6 +- .../src/utils/__tests__/versionWrites.test.js | 81 +++++++++++++++ .../cubejs/src/utils/dataSourceHelpers.js | 99 ++++++++++++++----- 3 files changed, 162 insertions(+), 24 deletions(-) create mode 100644 services/cubejs/src/utils/__tests__/versionWrites.test.js diff --git a/services/actions/src/rpc/auditVersionRollback.js b/services/actions/src/rpc/auditVersionRollback.js index 52e12c84..ebffd65b 100644 --- a/services/actions/src/rpc/auditVersionRollback.js +++ b/services/actions/src/rpc/auditVersionRollback.js @@ -76,7 +76,11 @@ export default async (session, input) => { target_id: row.id, outcome: "success", error_code: null, - payload: { origin: row.origin, checksum: row.checksum }, + payload: { + origin: row.origin, + checksum: row.checksum, + source_version_id: row.source_version_id ?? null, + }, }); const id = res?.data?.insert_audit_logs_one?.id; return { ok: true, auditLogId: id }; diff --git a/services/cubejs/src/utils/__tests__/versionWrites.test.js b/services/cubejs/src/utils/__tests__/versionWrites.test.js new file mode 100644 index 00000000..f5ff635b --- /dev/null +++ b/services/cubejs/src/utils/__tests__/versionWrites.test.js @@ -0,0 +1,81 @@ +import { describe, it, mock, beforeEach } from "node:test"; +import assert from "node:assert/strict"; + +const fetchGraphQLMock = mock.fn(); +mock.module("../graphql.js", { + namedExports: { fetchGraphQL: fetchGraphQLMock }, +}); + +const { commitVersionFiles, rollbackVersion } = await import( + "../dataSourceHelpers.js" +); + +const DS = "ds-1"; +const BRANCH = "branch-1"; + +describe("commitVersionFiles / rollbackVersion", () => { + beforeEach(() => fetchGraphQLMock.mock.resetCalls()); + + it("inserts a new version with the caller token and returns the new rows", async () => { + fetchGraphQLMock.mock.mockImplementation(async () => ({ + data: { + insert_versions_one: { + id: "v-new", + dataschemas: [{ id: "d-new", name: "a.yml", checksum: "c", version_id: "v-new" }], + }, + }, + })); + const res = await commitVersionFiles({ + branchId: BRANCH, + userId: "u-1", + datasourceId: DS, + files: [{ id: "d-old", name: "a.yml", code: "cubes: []", checksum: "x" }], + authToken: "tok", + }); + assert.deepEqual(res, { + newVersionId: "v-new", + dataschemas: [{ id: "d-new", name: "a.yml", checksum: "c", version_id: "v-new" }], + }); + const [, vars, token, opts] = fetchGraphQLMock.mock.calls[0].arguments; + assert.equal(token, "tok"); + assert.deepEqual(opts, { preserveErrors: true }); + assert.equal(vars.object.origin, "user"); + assert.equal("source_version_id" in vars.object, false); + // Only insertable columns — never the source row's id/checksum. + assert.deepEqual(vars.object.dataschemas.data, [ + { name: "a.yml", code: "cubes: []", user_id: "u-1", datasource_id: DS }, + ]); + }); + + it("passes Hasura errors through for mapping", async () => { + const errors = [{ extensions: { code: "permission-error" } }]; + fetchGraphQLMock.mock.mockImplementation(async () => ({ data: null, errors })); + const res = await commitVersionFiles({ + branchId: BRANCH, + userId: "u-1", + datasourceId: DS, + files: [], + authToken: "tok", + }); + assert.deepEqual(res, { errors }); + }); + + it("rollback records origin=rollback and the restored version as source", async () => { + fetchGraphQLMock.mock.mockImplementation(async (query) => + query.includes("VersionDataschemas") + ? { data: { dataschemas: [{ id: "d-1", name: "a.yml", code: "x" }] } } + : { data: { insert_versions_one: { id: "v-new", dataschemas: [] } } } + ); + const res = await rollbackVersion({ + branchId: BRANCH, + toVersionId: "v-old", + userId: "u-1", + datasourceId: DS, + authToken: "tok", + }); + assert.deepEqual(res, { newVersionId: "v-new", clonedDataschemaCount: 1 }); + const insert = fetchGraphQLMock.mock.calls[1].arguments[1].object; + assert.equal(insert.origin, "rollback"); + assert.equal(insert.source_version_id, "v-old"); + }); +}); diff --git a/services/cubejs/src/utils/dataSourceHelpers.js b/services/cubejs/src/utils/dataSourceHelpers.js index ae887eea..8d3ce04a 100644 --- a/services/cubejs/src/utils/dataSourceHelpers.js +++ b/services/cubejs/src/utils/dataSourceHelpers.js @@ -186,6 +186,22 @@ const upsertVersionMutation = ` } `; +// Same insert as upsertVersionMutation, returning the new rows so a caller +// can report the ids of the files it just wrote. +const insertVersionReturningMutation = ` + mutation ($object: versions_insert_input!) { + insert_versions_one(object: $object) { + id + dataschemas { + id + name + checksum + version_id + } + } + } +`; + const sqlCredentialsQuery = ` query ($username: String!) { sql_credentials(where: {username: {_eq: $username}}) { @@ -363,29 +379,31 @@ export const findVersionBranch = async ({ versionId }) => { }; /** - * Insert a new version on `branchId` whose dataschemas are byte-identical - * clones of `toVersionId`'s dataschemas. FR-013 / FR-013a: - * - only dataschemas are cloned (no explorations/alerts/alerts), - * - the new row uses `origin: 'rollback'`, - * - the caller's minted Hasura token is used so owner/admin permission - * policies are enforced at the database layer. + * Insert a new version on `branchId` holding exactly `files` (`[{name, code}]`). + * Every model change goes through a new version so the previous one stays + * restorable; `versions_flip_is_current_trg` makes the new row current. * - * Returns `{newVersionId, clonedDataschemaCount}` on success, - * `{errors}` on any Hasura permission/constraint failure so the caller can - * map the extensions.code via mapHasuraErrorCode(). + * The caller's minted Hasura token is used so owner/admin permission policies + * are enforced at the database layer too. On success the caller's cached + * user scope is dropped so its next request resolves the new version. + * + * Returns `{newVersionId, dataschemas}` on success, `{errors}` on any Hasura + * permission/constraint failure so the caller can map the extensions.code via + * mapHasuraErrorCode(). */ -export const rollbackVersion = async ({ +export const commitVersionFiles = async ({ branchId, - toVersionId, userId, datasourceId, + files, + origin = "user", + sourceVersionId = null, authToken, }) => { - const originals = await findVersionDataschemas({ versionId: toVersionId }); // The `set_public_dataschemas_checksum` BEFORE-INSERT trigger computes // dataschema-level checksums; we do NOT set one here or Hasura's insert // permission rejects the extra column. - const clonedData = originals.map((row) => ({ + const data = files.map((row) => ({ name: row.name, code: row.code, user_id: userId, @@ -395,23 +413,23 @@ export const rollbackVersion = async ({ // Version-level checksum: md5 over the concatenated dataschema codes in // a stable order. Matches the `version.checksum` NOT NULL constraint. const versionChecksum = md5OfCode( - originals - .map((r) => r.name) - .sort() - .map((n) => `${n}:${originals.find((x) => x.name === n)?.code || ""}`) + [...files] + .sort((a, b) => (a.name < b.name ? -1 : a.name > b.name ? 1 : 0)) + .map((r) => `${r.name}:${r.code || ""}`) .join("\n") ); const object = { branch_id: branchId, user_id: userId, - origin: "rollback", + origin, checksum: versionChecksum, - dataschemas: { data: clonedData }, + dataschemas: { data }, }; + if (sourceVersionId) object.source_version_id = sourceVersionId; const res = await fetchGraphQL( - upsertVersionMutation, + insertVersionReturningMutation, { object }, authToken, { preserveErrors: true } @@ -421,8 +439,8 @@ export const rollbackVersion = async ({ return { errors: res.errors }; } - const newVersionId = res?.data?.insert_versions_one?.id; - if (!newVersionId) { + const row = res?.data?.insert_versions_one; + if (!row?.id) { return { errors: [ { message: "insert_versions_one returned no id", extensions: {} }, @@ -430,7 +448,42 @@ export const rollbackVersion = async ({ }; } - return { newVersionId, clonedDataschemaCount: clonedData.length }; + invalidateUserCache(userId); + return { newVersionId: row.id, dataschemas: row.dataschemas || [] }; +}; + +/** + * Insert a new version on `branchId` whose dataschemas are byte-identical + * clones of `toVersionId`'s dataschemas. FR-013 / FR-013a: + * - only dataschemas are cloned (no explorations/alerts/alerts), + * - the new row uses `origin: 'rollback'` and records + * `source_version_id = toVersionId`. + * + * Returns `{newVersionId, clonedDataschemaCount}` on success, + * `{errors}` on any Hasura permission/constraint failure. + */ +export const rollbackVersion = async ({ + branchId, + toVersionId, + userId, + datasourceId, + authToken, +}) => { + const originals = await findVersionDataschemas({ versionId: toVersionId }); + const res = await commitVersionFiles({ + branchId, + userId, + datasourceId, + files: originals, + origin: "rollback", + sourceVersionId: toVersionId, + authToken, + }); + if (res.errors) return res; + return { + newVersionId: res.newVersionId, + clonedDataschemaCount: originals.length, + }; }; function md5OfCode(code) { From 61ade3c839cce279075695933acb1f3705a15b8f Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:34:15 +0000 Subject: [PATCH 04/12] fix(cubejs): dataschema delete writes a new version DELETE /api/v1/dataschema/:id deleted the row from the current version in place, destroying history. It now writes a new version on the branch with every other file of the current version, so the delete is undoable with a rollback. Guards are unchanged (auth, partition, owner/admin, current + active, reference scan/409) and moved to utils/mutableDataschema.js for reuse. The success audit row is written directly (the delete event trigger no longer fires); the response adds versionId + branchId. Co-Authored-By: Claude Opus 5.5 --- .../cubejs/src/routes/deleteDataschema.js | 338 ++++-------------- .../cubejs/src/utils/mutableDataschema.js | 186 ++++++++++ 2 files changed, 257 insertions(+), 267 deletions(-) create mode 100644 services/cubejs/src/utils/mutableDataschema.js diff --git a/services/cubejs/src/routes/deleteDataschema.js b/services/cubejs/src/routes/deleteDataschema.js index 0e9249d6..8340ef9b 100644 --- a/services/cubejs/src/routes/deleteDataschema.js +++ b/services/cubejs/src/routes/deleteDataschema.js @@ -1,65 +1,20 @@ import YAML from "yaml"; -import { verifyAndProvision } from "../utils/directVerifyAuth.js"; import { emitModelEvent } from "../utils/eventEmitter.js"; -import { findUser } from "../utils/dataSourceHelpers.js"; -import { fetchGraphQL } from "../utils/graphql.js"; -import { mintHasuraToken } from "../utils/mintHasuraToken.js"; -import { mintedTokenCache } from "../utils/mintedTokenCache.js"; -import { requireOwnerOrAdmin } from "../utils/requireOwnerOrAdmin.js"; -import { resolvePartitionTeamIds } from "./discover.js"; +import { + commitVersionFiles, + findVersionDataschemas, +} from "../utils/dataSourceHelpers.js"; +import { + ensureHasuraTokenForUser, + resolveMutableDataschema, + respondError, +} from "../utils/mutableDataschema.js"; import { scanCrossCubeReferences } from "../utils/referenceScanner.js"; -import { writeAuditLog } from "../utils/auditWriter.js"; import { mapHasuraErrorCode } from "../utils/mapHasuraErrorCode.js"; import { ErrorCode } from "../utils/errorCodes.js"; import { parseCubesFromJs } from "../utils/smart-generation/diffModels.js"; -const RESOLVE_TARGET_QUERY = ` - query ResolveTargetDataschema($id: uuid!) { - dataschemas_by_pk(id: $id) { - id - name - code - version_id - version { - id - is_current - branch { - id - status - datasource { - id - team_id - } - } - } - } - } -`; - -const SIBLINGS_QUERY = ` - query Siblings($versionId: uuid!, $excludeId: uuid!) { - dataschemas( - where: { - version_id: {_eq: $versionId} - id: {_neq: $excludeId} - } - ) { - id - name - code - } - } -`; - -const DELETE_MUTATION = ` - mutation DeleteDataschema($id: uuid!) { - delete_dataschemas_by_pk(id: $id) { - id - } - } -`; - function parseCubes(name, code) { if (!code) return []; const isYaml = name?.endsWith(".yml") || name?.endsWith(".yaml"); @@ -75,166 +30,48 @@ function parseCubes(name, code) { } } -async function ensureHasuraTokenForUser(userId) { - let tok = mintedTokenCache.get(userId); - if (tok) return tok; - tok = await mintHasuraToken(userId); - const decoded = JSON.parse( - Buffer.from(tok.split(".")[1], "base64url").toString() - ); - mintedTokenCache.set(userId, tok, decoded.exp); - return tok; -} - -function respondError(res, status, code, message, extra = {}) { - return res.status(status).json({ code, message, ...extra }); -} - /** * DELETE /api/v1/dataschema/:dataschemaId * - * Remove a dataschema row from the currently-active version of its branch. + * Remove a dataschema from the current version of its branch by writing a NEW + * version holding every other file of that version — the version the file + * was deleted from stays intact, so the delete is restorable via rollback. * Enforces, in order: - * - authentication (FR-015 direct-verify) - * - partition gate (FR-015) - * - owner/admin role on the datasource's team (FR-015) - * - version-level immutability via `is_current=true` + `branch.status=active` (FR-007) + * - authentication, partition, owner/admin, current version of the active + * branch (resolveMutableDataschema — FR-015 / FR-007) * - cross-cube reference scan (FR-008, seven kinds) * - * Every rejection path emits a durable audit row with `outcome='failure'` via - * `writeAuditLog` (FR-016). Successful deletes are captured by the - * `delete_dataschema_audit` Hasura event trigger. + * Every outcome writes a durable audit row via `writeAuditLog` (FR-016); the + * `delete_dataschema_audit` event trigger no longer fires because no row is + * deleted. */ export default async function deleteDataschema(req, res) { - const verified = await verifyAndProvision(req); - if (verified.error) { - return respondError( - res, - verified.error.status, - verified.error.code, - verified.error.message - ); - } - const { payload, userId } = verified; - - const dataschemaId = req.params?.dataschemaId; - if (!dataschemaId || typeof dataschemaId !== "string") { - return respondError( - res, - 400, - "delete_invalid_request", - "dataschemaId path parameter is required" - ); - } - - // Resolve the target via admin-secret GraphQL (handler owns enforcement). - let targetRow; - try { - const r = await fetchGraphQL(RESOLVE_TARGET_QUERY, { id: dataschemaId }); - targetRow = r?.data?.dataschemas_by_pk; - } catch (err) { - return respondError( - res, - 503, - "hasura_unavailable", - err?.message || "Hasura unavailable" - ); - } - - if (!targetRow) { - return respondError( - res, - 404, - ErrorCode.VALIDATE_TARGET_NOT_FOUND, - "Dataschema not found" - ); - } - - const version = targetRow.version; - const branch = version?.branch; - const datasource = branch?.datasource; - const teamId = datasource?.team_id; - const datasourceId = datasource?.id; - const branchId = branch?.id; - - const user = await findUser({ userId }); - - // Partition gate. - const partitionTeamIds = resolvePartitionTeamIds( - user.members, - payload.partition - ); - if (partitionTeamIds && !partitionTeamIds.has(teamId)) { - await writeAuditLog({ - action: "dataschema_delete", - userId, - datasourceId, - branchId, - targetId: dataschemaId, - outcome: "failure", - errorCode: ErrorCode.DELETE_BLOCKED_AUTHORIZATION, - payload: { reason: "partition_mismatch" }, - }); - return respondError( - res, - 403, - ErrorCode.DELETE_BLOCKED_AUTHORIZATION, - "Caller's partition does not match the datasource's team" - ); - } - - // Owner/admin gate. - if (!requireOwnerOrAdmin(user, teamId)) { - await writeAuditLog({ - action: "dataschema_delete", - userId, - datasourceId, - branchId, - targetId: dataschemaId, - outcome: "failure", - errorCode: ErrorCode.DELETE_BLOCKED_AUTHORIZATION, - payload: { reason: "insufficient_role" }, - }); - return respondError( - res, - 403, - ErrorCode.DELETE_BLOCKED_AUTHORIZATION, - "Owner or admin role required" - ); - } - - // Version-level immutability (FR-007): only the current version of the - // active branch may be edited. - if (version?.is_current !== true || branch?.status !== "active") { - await writeAuditLog({ - action: "dataschema_delete", - userId, - datasourceId, - branchId, - targetId: dataschemaId, - outcome: "failure", - errorCode: ErrorCode.DELETE_BLOCKED_HISTORICAL_VERSION, - payload: { - is_current: version?.is_current ?? null, - branch_status: branch?.status ?? null, - }, - }); - return respondError( - res, - 409, - ErrorCode.DELETE_BLOCKED_HISTORICAL_VERSION, - "Dataschema is attached to a historical version — only the current version of the active branch is mutable" - ); - } + const ctx = await resolveMutableDataschema(req, res, { + action: "dataschema_delete", + codes: { + invalidRequest: "delete_invalid_request", + authorization: ErrorCode.DELETE_BLOCKED_AUTHORIZATION, + historical: ErrorCode.DELETE_BLOCKED_HISTORICAL_VERSION, + }, + }); + if (!ctx) return; + const { + payload, + userId, + dataschemaId, + target, + versionId, + branchId, + datasourceId, + audit, + } = ctx; // Cross-cube reference scan (FR-008). let siblings; try { - const r = await fetchGraphQL(SIBLINGS_QUERY, { - versionId: version.id, - excludeId: dataschemaId, - }); - siblings = r?.data?.dataschemas || []; + siblings = (await findVersionDataschemas({ versionId })).filter( + (row) => row.id !== dataschemaId + ); } catch (err) { return respondError( res, @@ -244,7 +81,7 @@ export default async function deleteDataschema(req, res) { ); } - const targetCubeNames = parseCubes(targetRow.name, targetRow.code).map( + const targetCubeNames = parseCubes(target.name, target.code).map( (c) => c.name ); const otherCubes = siblings.flatMap((row) => @@ -263,15 +100,8 @@ export default async function deleteDataschema(req, res) { } if (blockingReferences.length > 0) { - await writeAuditLog({ - action: "dataschema_delete", - userId, - datasourceId, - branchId, - targetId: dataschemaId, - outcome: "failure", - errorCode: ErrorCode.DELETE_BLOCKED_BY_REFERENCES, - payload: { blockingReferences }, + await audit("failure", ErrorCode.DELETE_BLOCKED_BY_REFERENCES, { + blockingReferences, }); return respondError( res, @@ -282,8 +112,8 @@ export default async function deleteDataschema(req, res) { ); } - // Fire the actual delete with the caller's minted Hasura token so the - // user-role delete_permissions filter applies at the DB layer too (two-layer + // Write the new version with the caller's minted Hasura token so the + // user-role insert permissions apply at the DB layer too (two-layer // defence per research R4). let hasuraToken; try { @@ -297,14 +127,15 @@ export default async function deleteDataschema(req, res) { ); } - let del; + let result; try { - del = await fetchGraphQL( - DELETE_MUTATION, - { id: dataschemaId }, - hasuraToken, - { preserveErrors: true } - ); + result = await commitVersionFiles({ + branchId, + userId, + datasourceId, + files: siblings, + authToken: hasuraToken, + }); } catch (err) { return respondError( res, @@ -314,18 +145,11 @@ export default async function deleteDataschema(req, res) { ); } - if (del?.errors) { - const mapped = mapHasuraErrorCode(del.errors, { action: "delete" }); + if (result.errors) { + const mapped = mapHasuraErrorCode(result.errors, { action: "delete" }); if (mapped === ErrorCode.DELETE_BLOCKED_AUTHORIZATION) { - await writeAuditLog({ - action: "dataschema_delete", - userId, - datasourceId, - branchId, - targetId: dataschemaId, - outcome: "failure", - errorCode: mapped, - payload: { hasura_code: del.errors?.[0]?.extensions?.code || null }, + await audit("failure", mapped, { + hasura_code: result.errors?.[0]?.extensions?.code || null, }); return respondError( res, @@ -334,16 +158,7 @@ export default async function deleteDataschema(req, res) { "Hasura rejected the delete (permission-error)" ); } - await writeAuditLog({ - action: "dataschema_delete", - userId, - datasourceId, - branchId, - targetId: dataschemaId, - outcome: "failure", - errorCode: "hasura_rejected", - payload: { errors: del.errors }, - }); + await audit("failure", "hasura_rejected", { errors: result.errors }); return respondError( res, 503, @@ -352,28 +167,11 @@ export default async function deleteDataschema(req, res) { ); } - if (!del?.data?.delete_dataschemas_by_pk?.id) { - // Row vanished concurrently — treat as not-found and audit for visibility. - await writeAuditLog({ - action: "dataschema_delete", - userId, - datasourceId, - branchId, - targetId: dataschemaId, - outcome: "failure", - errorCode: ErrorCode.VALIDATE_TARGET_NOT_FOUND, - payload: { reason: "row_not_found_at_delete" }, - }); - return respondError( - res, - 404, - ErrorCode.VALIDATE_TARGET_NOT_FOUND, - "Dataschema not found at delete time" - ); - } - - // Success path: the Hasura delete event trigger `delete_dataschema_audit` - // writes the outcome='success' audit row. Handler does not duplicate. + await audit("success", null, { + name: target.name, + version_id: versionId, + new_version_id: result.newVersionId, + }); // 099 T087 (FR-091): a successful delete is a model lifecycle fact. // Fire-and-forget; never blocks the response (FR-007). @@ -383,14 +181,20 @@ export default async function deleteDataschema(req, res) { partition: payload?.partition ?? null, userId, modelId: dataschemaId, - modelLabel: targetRow.name || null, + modelLabel: target.name || null, status: "ok", properties: { datasource_id: datasourceId, branch_id: branchId, - version_id: version?.id ?? null, + version_id: versionId, + new_version_id: result.newVersionId, }, }); - return res.json({ deleted: true, dataschemaId }); + return res.json({ + deleted: true, + dataschemaId, + versionId: result.newVersionId, + branchId, + }); } diff --git a/services/cubejs/src/utils/mutableDataschema.js b/services/cubejs/src/utils/mutableDataschema.js new file mode 100644 index 00000000..967171c5 --- /dev/null +++ b/services/cubejs/src/utils/mutableDataschema.js @@ -0,0 +1,186 @@ +import { verifyAndProvision } from "./directVerifyAuth.js"; +import { findUser } from "./dataSourceHelpers.js"; +import { fetchGraphQL } from "./graphql.js"; +import { mintHasuraToken } from "./mintHasuraToken.js"; +import { mintedTokenCache } from "./mintedTokenCache.js"; +import { requireOwnerOrAdmin } from "./requireOwnerOrAdmin.js"; +import { resolvePartitionTeamIds } from "../routes/discover.js"; +import { writeAuditLog } from "./auditWriter.js"; +import { ErrorCode } from "./errorCodes.js"; + +const RESOLVE_TARGET_QUERY = ` + query ResolveTargetDataschema($id: uuid!) { + dataschemas_by_pk(id: $id) { + id + name + code + checksum + version_id + version { + id + is_current + branch { + id + status + datasource { + id + team_id + } + } + } + } + } +`; + +export function respondError(res, status, code, message, extra = {}) { + return res.status(status).json({ code, message, ...extra }); +} + +export async function ensureHasuraTokenForUser(userId) { + let tok = mintedTokenCache.get(userId); + if (tok) return tok; + tok = await mintHasuraToken(userId); + const decoded = JSON.parse( + Buffer.from(tok.split(".")[1], "base64url").toString() + ); + mintedTokenCache.set(userId, tok, decoded.exp); + return tok; +} + +/** + * Shared guard for the single-dataschema write routes + * (DELETE / PUT /api/v1/dataschema/:dataschemaId). Enforces, in order: + * - authentication (FR-015 direct-verify) + * - partition gate (FR-015) + * - owner/admin role on the datasource's team (FR-015) + * - the dataschema is on the current version of the active branch (FR-007) + * + * Every rejection after the target resolves writes a durable failure audit + * row (FR-016). Returns `null` when a response has already been sent; + * otherwise the resolved target plus an `audit(outcome, errorCode, payload)` + * helper bound to this request. + * + * @param {{action: 'dataschema_delete'|'dataschema_update', + * codes: {invalidRequest:string, authorization:string, historical:string}}} opts + */ +export async function resolveMutableDataschema(req, res, { action, codes }) { + const verified = await verifyAndProvision(req); + if (verified.error) { + respondError( + res, + verified.error.status, + verified.error.code, + verified.error.message + ); + return null; + } + const { payload, userId } = verified; + + const dataschemaId = req.params?.dataschemaId; + if (!dataschemaId || typeof dataschemaId !== "string") { + respondError( + res, + 400, + codes.invalidRequest, + "dataschemaId path parameter is required" + ); + return null; + } + + // Resolve the target via admin-secret GraphQL (handler owns enforcement). + let target; + try { + const r = await fetchGraphQL(RESOLVE_TARGET_QUERY, { id: dataschemaId }); + target = r?.data?.dataschemas_by_pk; + } catch (err) { + respondError( + res, + 503, + "hasura_unavailable", + err?.message || "Hasura unavailable" + ); + return null; + } + + if (!target) { + respondError( + res, + 404, + ErrorCode.VALIDATE_TARGET_NOT_FOUND, + "Dataschema not found" + ); + return null; + } + + const version = target.version; + const branch = version?.branch; + const teamId = branch?.datasource?.team_id; + const datasourceId = branch?.datasource?.id; + const branchId = branch?.id; + + const audit = (outcome, errorCode, auditPayload) => + writeAuditLog({ + action, + userId, + datasourceId, + branchId, + targetId: dataschemaId, + outcome, + errorCode, + payload: auditPayload, + }); + + const user = await findUser({ userId }); + + const partitionTeamIds = resolvePartitionTeamIds( + user.members, + payload.partition + ); + if (partitionTeamIds && !partitionTeamIds.has(teamId)) { + await audit("failure", codes.authorization, { + reason: "partition_mismatch", + }); + respondError( + res, + 403, + codes.authorization, + "Caller's partition does not match the datasource's team" + ); + return null; + } + + if (!requireOwnerOrAdmin(user, teamId)) { + await audit("failure", codes.authorization, { + reason: "insufficient_role", + }); + respondError(res, 403, codes.authorization, "Owner or admin role required"); + return null; + } + + // Version-level immutability (FR-007): only the current version of the + // active branch may be changed (by writing a new version on top of it). + if (version?.is_current !== true || branch?.status !== "active") { + await audit("failure", codes.historical, { + is_current: version?.is_current ?? null, + branch_status: branch?.status ?? null, + }); + respondError( + res, + 409, + codes.historical, + "Dataschema is attached to a historical version — only the current version of the active branch is mutable" + ); + return null; + } + + return { + payload, + userId, + dataschemaId, + target, + versionId: version.id, + branchId, + datasourceId, + audit, + }; +} From 0c2132cd8dfab19318f3b763baeb6c0c9ba15739 Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:34:32 +0000 Subject: [PATCH 05/12] feat(cubejs): PUT /api/v1/dataschema/:id saves one file as a new version MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Body {code}. Copies every file of the dataschema's version into a new version with this file's code replaced, so the previous version stays restorable — replaces the in-place Hasura update_dataschemas_by_pk the agent refinement script uses. Same guards as DELETE (direct-verify, partition, owner/admin, current version of an active branch → else 409). Identical code returns 200 {unchanged: true} and writes nothing. Response: {versionId, branchId, dataschema: {id (new row), name, checksum, version_id}}. Audits under action dataschema_update and emits Model Saved like createDataSchema. New error codes update_blocked_authorization / update_blocked_historical_version. Co-Authored-By: Claude Opus 5.5 --- CLAUDE.md | 2 +- services/cubejs/src/routes/index.js | 8 + .../cubejs/src/routes/updateDataschema.js | 169 ++++++++++++++++++ .../__tests__/mapHasuraErrorCode.test.js | 8 + services/cubejs/src/utils/auditWriter.js | 2 +- services/cubejs/src/utils/errorCodes.js | 2 + .../cubejs/src/utils/mapHasuraErrorCode.js | 3 +- .../contracts/delete-dataschema.yaml | 83 ++++++++- .../contracts/meta-single-cube.yaml | 2 + .../contracts/refresh-compiler.yaml | 2 + .../contracts/validate-in-branch.yaml | 2 + .../contracts/version-diff.yaml | 2 + .../contracts/version-rollback.yaml | 2 + 13 files changed, 282 insertions(+), 5 deletions(-) create mode 100644 services/cubejs/src/routes/updateDataschema.js diff --git a/CLAUDE.md b/CLAUDE.md index 90546f39..94bc7010 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -113,7 +113,7 @@ Docker Compose files per environment: `docker-compose.dev.yml`, `docker-compose. - `services/cubejs/src/routes/` — 13 REST API endpoints, now including: - `validateInBranch.js` (POST /api/v1/validate-in-branch, US1) - `refreshCompiler.js` (POST /api/v1/internal/refresh-compiler, US2) - - `deleteDataschema.js` (DELETE /api/v1/dataschema/:id, US3) + - `deleteDataschema.js` (DELETE /api/v1/dataschema/:id, US3) + `updateDataschema.js` (PUT, same path) — both write a NEW version via `utils/mutableDataschema.js` + `commitVersionFiles` - `metaSingleCube.js` (GET /api/v1/meta/cube/:cubeName, US4) - `versionDiff.js` + `versionRollback.js` (POST /api/v1/version/{diff,rollback}, US5) - `services/cubejs/src/routes/reconcileTeam.js` — POST /api/v1/internal/reconcile-team: per-team default-models worker (013 — probe→generate→merge→validate→publish, system-user only) diff --git a/services/cubejs/src/routes/index.js b/services/cubejs/src/routes/index.js index 6edeff77..57fe8b2f 100644 --- a/services/cubejs/src/routes/index.js +++ b/services/cubejs/src/routes/index.js @@ -5,6 +5,7 @@ import express from "express"; // POST /api/v1/validate-in-branch (direct-verify, US1) // POST /api/v1/internal/refresh-compiler (direct-verify, US2) // DELETE /api/v1/dataschema/:dataschemaId (direct-verify, US3) +// PUT /api/v1/dataschema/:dataschemaId (direct-verify, save one file) // GET /api/v1/meta/cube/:cubeName (checkAuthMiddleware, US4) // POST /api/v1/version/diff (direct-verify, US5) // POST /api/v1/version/rollback (direct-verify, US5) @@ -42,6 +43,7 @@ import discover from "./discover.js"; import metaAll from "./metaAll.js"; import testConnection from "./testConnection.js"; import deleteDataschema from "./deleteDataschema.js"; +import updateDataschema from "./updateDataschema.js"; import metaSingleCube from "./metaSingleCube.js"; import refreshCompiler from "./refreshCompiler.js"; import reconcileTeam from "./reconcileTeam.js"; @@ -333,6 +335,12 @@ export default ({ basePath, cubejs }) => { async (req, res) => deleteDataschema(req, res) ); + // Save one dataschema's code as a new version. Owner/admin only. + router.put( + `${basePath}/v1/dataschema/:dataschemaId`, + async (req, res) => updateDataschema(req, res) + ); + // Model Management API: single-cube metadata (US4). // Datasource-scoped: runs behind checkAuthMiddleware (x-hasura-datasource-id // is mandatory by contract). The /cube/ path segment prevents collision diff --git a/services/cubejs/src/routes/updateDataschema.js b/services/cubejs/src/routes/updateDataschema.js new file mode 100644 index 00000000..bc5fff79 --- /dev/null +++ b/services/cubejs/src/routes/updateDataschema.js @@ -0,0 +1,169 @@ +import { emitModelEvent } from "../utils/eventEmitter.js"; +import { + commitVersionFiles, + findVersionDataschemas, +} from "../utils/dataSourceHelpers.js"; +import { + ensureHasuraTokenForUser, + resolveMutableDataschema, + respondError, +} from "../utils/mutableDataschema.js"; +import { mapHasuraErrorCode } from "../utils/mapHasuraErrorCode.js"; +import { ErrorCode } from "../utils/errorCodes.js"; + +const dataschemaShape = (row) => ({ + id: row.id, + name: row.name, + checksum: row.checksum, + version_id: row.version_id, +}); + +/** + * PUT /api/v1/dataschema/:dataschemaId body: `{code: string}` + * + * Save one file's new code as a NEW version on the dataschema's branch: every + * file of the dataschema's (current) version is copied, this file's code is + * replaced. The previous version stays intact and restorable. Same guards as + * DELETE (auth, partition, owner/admin, current version of an active branch → + * else 409). Identical code is a no-op: 200 with `unchanged: true`. + * + * Returns `{versionId, branchId, dataschema: {id, name, checksum, version_id}}` + * where `dataschema.id` is the NEW row's id. + */ +export default async function updateDataschema(req, res) { + const ctx = await resolveMutableDataschema(req, res, { + action: "dataschema_update", + codes: { + invalidRequest: "update_invalid_request", + authorization: ErrorCode.UPDATE_BLOCKED_AUTHORIZATION, + historical: ErrorCode.UPDATE_BLOCKED_HISTORICAL_VERSION, + }, + }); + if (!ctx) return; + const { + payload, + userId, + dataschemaId, + target, + versionId, + branchId, + datasourceId, + audit, + } = ctx; + + const code = req.body?.code; + if (typeof code !== "string") { + return respondError( + res, + 400, + "update_invalid_request", + "Body must be {code: string}" + ); + } + + if (code === target.code) { + return res.json({ + versionId, + branchId, + dataschema: dataschemaShape(target), + unchanged: true, + }); + } + + let files; + try { + files = (await findVersionDataschemas({ versionId })).map((row) => + row.id === dataschemaId ? { ...row, code } : row + ); + } catch (err) { + return respondError( + res, + 503, + "hasura_unavailable", + err?.message || "Hasura unavailable" + ); + } + + let hasuraToken; + try { + hasuraToken = await ensureHasuraTokenForUser(userId); + } catch { + return respondError( + res, + 503, + "auth_unavailable", + "Unable to mint Hasura token" + ); + } + + let result; + try { + result = await commitVersionFiles({ + branchId, + userId, + datasourceId, + files, + authToken: hasuraToken, + }); + } catch (err) { + return respondError( + res, + 503, + "hasura_unavailable", + err?.message || "Hasura unavailable" + ); + } + + if (result.errors) { + const mapped = mapHasuraErrorCode(result.errors, { action: "update" }); + if (mapped === ErrorCode.UPDATE_BLOCKED_AUTHORIZATION) { + await audit("failure", mapped, { + hasura_code: result.errors?.[0]?.extensions?.code || null, + }); + return respondError( + res, + 403, + mapped, + "Hasura rejected the update (permission-error)" + ); + } + await audit("failure", "hasura_rejected", { errors: result.errors }); + return respondError( + res, + 503, + "hasura_unavailable", + "Hasura rejected the update" + ); + } + + const saved = result.dataschemas.find((row) => row.name === target.name); + + await audit("success", null, { + name: target.name, + version_id: versionId, + new_version_id: result.newVersionId, + new_dataschema_id: saved?.id ?? null, + checksum: saved?.checksum ?? null, + }); + + // 099 T087 (FR-091): same `Model Saved` persistence fact createDataSchema + // emits for every other server-side version write. Fire-and-forget. + emitModelEvent({ + event: "Model Saved", + accountId: payload?.accountId ?? null, + partition: payload?.partition ?? null, + userId, + modelId: result.newVersionId, + modelLabel: target.name || null, + status: "ok", + properties: { branch_id: branchId, origin: "user" }, + }); + + return res.json({ + versionId: result.newVersionId, + branchId, + dataschema: saved + ? dataschemaShape(saved) + : { id: null, name: target.name, checksum: null, version_id: result.newVersionId }, + }); +} diff --git a/services/cubejs/src/utils/__tests__/mapHasuraErrorCode.test.js b/services/cubejs/src/utils/__tests__/mapHasuraErrorCode.test.js index 29175642..ccf97c6e 100644 --- a/services/cubejs/src/utils/__tests__/mapHasuraErrorCode.test.js +++ b/services/cubejs/src/utils/__tests__/mapHasuraErrorCode.test.js @@ -23,6 +23,14 @@ describe("mapHasuraErrorCode", () => { assert.equal(code, ErrorCode.DELETE_BLOCKED_AUTHORIZATION); }); + it("permission-error maps to update_blocked_authorization for update action", () => { + const code = mapHasuraErrorCode( + [{ extensions: { code: "permission-error" } }], + { action: "update" } + ); + assert.equal(code, ErrorCode.UPDATE_BLOCKED_AUTHORIZATION); + }); + it("permission-error maps to rollback_blocked_authorization for rollback action", () => { const code = mapHasuraErrorCode( [{ extensions: { code: "permission-error" } }], diff --git a/services/cubejs/src/utils/auditWriter.js b/services/cubejs/src/utils/auditWriter.js index 6e2178d3..6081785c 100644 --- a/services/cubejs/src/utils/auditWriter.js +++ b/services/cubejs/src/utils/auditWriter.js @@ -47,7 +47,7 @@ function sleep(ms) { * permission policies on `audit_logs` (admin-only) do not reject the write. * * @param {object} args - * @param {'dataschema_delete'|'version_rollback'} args.action + * @param {'dataschema_delete'|'dataschema_update'|'version_rollback'} args.action * @param {string} args.userId * @param {string} args.targetId * @param {string} [args.datasourceId] diff --git a/services/cubejs/src/utils/errorCodes.js b/services/cubejs/src/utils/errorCodes.js index 6bd08422..890928fb 100644 --- a/services/cubejs/src/utils/errorCodes.js +++ b/services/cubejs/src/utils/errorCodes.js @@ -22,6 +22,8 @@ export const ErrorCode = Object.freeze({ ROLLBACK_BLOCKED_AUTHORIZATION: "rollback_blocked_authorization", ROLLBACK_INVALID_REQUEST: "rollback_invalid_request", ROLLBACK_SOURCE_COLUMNS_MISSING: "rollback_source_columns_missing", + UPDATE_BLOCKED_HISTORICAL_VERSION: "update_blocked_historical_version", + UPDATE_BLOCKED_AUTHORIZATION: "update_blocked_authorization", }); export const ErrorCodeSet = Object.freeze(new Set(Object.values(ErrorCode))); diff --git a/services/cubejs/src/utils/mapHasuraErrorCode.js b/services/cubejs/src/utils/mapHasuraErrorCode.js index 03b55d5c..dd219e97 100644 --- a/services/cubejs/src/utils/mapHasuraErrorCode.js +++ b/services/cubejs/src/utils/mapHasuraErrorCode.js @@ -11,7 +11,7 @@ import { ErrorCode } from "./errorCodes.js"; * "propagate as 503 hasura_unavailable" per R11. * * @param {Array<{extensions?:{code?:string}, message?:string}>|null|undefined} errors - * @param {{action?: 'delete'|'rollback'|'validate'|'meta'|'diff'|'refresh'}} [ctx] + * @param {{action?: 'delete'|'update'|'rollback'|'validate'|'meta'|'diff'|'refresh'}} [ctx] * @returns {string|null} */ export function mapHasuraErrorCode(errors, ctx = {}) { @@ -24,6 +24,7 @@ export function mapHasuraErrorCode(errors, ctx = {}) { if (code === "permission-error" || code === "access-denied") { if (action === "delete") return ErrorCode.DELETE_BLOCKED_AUTHORIZATION; + if (action === "update") return ErrorCode.UPDATE_BLOCKED_AUTHORIZATION; if (action === "rollback") return ErrorCode.ROLLBACK_BLOCKED_AUTHORIZATION; return null; } diff --git a/specs/011-model-mgmt-api/contracts/delete-dataschema.yaml b/specs/011-model-mgmt-api/contracts/delete-dataschema.yaml index dbc58199..a618845b 100644 --- a/specs/011-model-mgmt-api/contracts/delete-dataschema.yaml +++ b/specs/011-model-mgmt-api/contracts/delete-dataschema.yaml @@ -8,7 +8,9 @@ paths: operationId: deleteDataschema summary: Delete a dataschema from the latest version of the active branch. description: | - Removes the dataschema row via Hasura (`delete_dataschemas_by_pk`). Blocked if: + Writes a NEW version on the branch holding every other file of the current + version; the version the file was removed from stays intact, so the delete + can be undone with a rollback. Blocked if: - The caller lacks owner/admin role on the datasource team (FR-006) → 403. - The dataschema belongs to a version that is NOT the latest version of the active branch (FR-007) → 409. @@ -30,7 +32,7 @@ paths: format: uuid responses: "200": - description: Dataschema deleted. + description: Dataschema removed; `versionId` is the new current version. content: application/json: schema: { $ref: '#/components/schemas/DeleteDataschemaResponse' } @@ -55,6 +57,59 @@ paths: content: application/json: schema: { $ref: '#/components/schemas/DeleteBlockedResponse' } + put: + operationId: updateDataschema + summary: Save one dataschema's code as a new version of its branch. + description: | + Copies every file of the dataschema's version into a NEW version and + replaces this file's code; the previous version stays restorable. + Identical code is a no-op (200, `unchanged: true`, no new version). + Same guards as DELETE: owner/admin → else 403 `update_blocked_authorization`; + the dataschema must be on the current version of the active branch → + else 409 `update_blocked_historical_version`. + + AUTH: direct-verify, as DELETE. + security: + - BearerAuth: [] + parameters: + - in: path + name: dataschemaId + required: true + schema: + type: string + format: uuid + requestBody: + required: true + content: + application/json: + schema: + type: object + required: [code] + properties: + code: { type: string } + responses: + "200": + description: Saved (or unchanged). + content: + application/json: + schema: { $ref: '#/components/schemas/UpdateDataschemaResponse' } + "400": + description: "Body is not `{code: string}` (`update_invalid_request`)." + "403": + description: "`update_blocked_authorization` (partition or role)." + content: + application/json: + schema: { $ref: '#/components/schemas/ErrorResponse' } + "404": + description: Dataschema not found. + content: + application/json: + schema: { $ref: '#/components/schemas/ErrorResponse' } + "409": + description: "`update_blocked_historical_version`." + content: + application/json: + schema: { $ref: '#/components/schemas/ErrorResponse' } components: securitySchemes: @@ -84,6 +139,8 @@ components: - rollback_blocked_authorization - rollback_invalid_request - rollback_source_columns_missing + - update_blocked_historical_version + - update_blocked_authorization ErrorResponse: type: object required: [code, message] @@ -96,6 +153,28 @@ components: properties: deleted: { type: boolean } dataschemaId: { type: string, format: uuid } + versionId: + type: string + format: uuid + description: The new current version (without the deleted file). + branchId: { type: string, format: uuid } + UpdateDataschemaResponse: + type: object + required: [versionId, branchId, dataschema] + properties: + versionId: { type: string, format: uuid } + branchId: { type: string, format: uuid } + unchanged: + type: boolean + description: Present (true) when the code was identical and nothing was written. + dataschema: + type: object + required: [id, name, checksum, version_id] + properties: + id: { type: string, format: uuid, description: The NEW row's id. } + name: { type: string } + checksum: { type: string } + version_id: { type: string, format: uuid } DeleteBlockedResponse: type: object required: [code, message] diff --git a/specs/011-model-mgmt-api/contracts/meta-single-cube.yaml b/specs/011-model-mgmt-api/contracts/meta-single-cube.yaml index 496436f4..71550863 100644 --- a/specs/011-model-mgmt-api/contracts/meta-single-cube.yaml +++ b/specs/011-model-mgmt-api/contracts/meta-single-cube.yaml @@ -82,6 +82,8 @@ components: - rollback_blocked_authorization - rollback_invalid_request - rollback_source_columns_missing + - update_blocked_historical_version + - update_blocked_authorization ErrorResponse: type: object required: [code, message] diff --git a/specs/011-model-mgmt-api/contracts/refresh-compiler.yaml b/specs/011-model-mgmt-api/contracts/refresh-compiler.yaml index 93377fc0..65a2d204 100644 --- a/specs/011-model-mgmt-api/contracts/refresh-compiler.yaml +++ b/specs/011-model-mgmt-api/contracts/refresh-compiler.yaml @@ -79,6 +79,8 @@ components: - rollback_blocked_authorization - rollback_invalid_request - rollback_source_columns_missing + - update_blocked_historical_version + - update_blocked_authorization ErrorResponse: type: object required: [code, message] diff --git a/specs/011-model-mgmt-api/contracts/validate-in-branch.yaml b/specs/011-model-mgmt-api/contracts/validate-in-branch.yaml index a32e4f2f..9985c704 100644 --- a/specs/011-model-mgmt-api/contracts/validate-in-branch.yaml +++ b/specs/011-model-mgmt-api/contracts/validate-in-branch.yaml @@ -78,6 +78,8 @@ components: - rollback_blocked_authorization - rollback_invalid_request - rollback_source_columns_missing + - update_blocked_historical_version + - update_blocked_authorization ErrorResponse: type: object required: [code, message] diff --git a/specs/011-model-mgmt-api/contracts/version-diff.yaml b/specs/011-model-mgmt-api/contracts/version-diff.yaml index 9b293d76..8dec0c41 100644 --- a/specs/011-model-mgmt-api/contracts/version-diff.yaml +++ b/specs/011-model-mgmt-api/contracts/version-diff.yaml @@ -73,6 +73,8 @@ components: - rollback_blocked_authorization - rollback_invalid_request - rollback_source_columns_missing + - update_blocked_historical_version + - update_blocked_authorization ErrorResponse: type: object required: [code, message] diff --git a/specs/011-model-mgmt-api/contracts/version-rollback.yaml b/specs/011-model-mgmt-api/contracts/version-rollback.yaml index b0538594..b86b474e 100644 --- a/specs/011-model-mgmt-api/contracts/version-rollback.yaml +++ b/specs/011-model-mgmt-api/contracts/version-rollback.yaml @@ -78,6 +78,8 @@ components: - rollback_blocked_authorization - rollback_invalid_request - rollback_source_columns_missing + - update_blocked_historical_version + - update_blocked_authorization ErrorResponse: type: object required: [code, message] From e6df2438fd6746b2be1ab5c5fef88a45beb395b8 Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:34:44 +0000 Subject: [PATCH 06/12] fix(cubejs): require the admin secret on internal invalidate-cache POST /api/v1/internal/invalidate-cache had no auth and is reachable via the /api ingress. Accept it only when x-hasura-admin-secret equals HASURA_GRAPHQL_ADMIN_SECRET (constant-time compare; 401 otherwise, and always 401 when the env var is unset). Its only caller, the actions service (utils/cubeCache.js), now sends that header. Co-Authored-By: Claude Opus 5.5 --- services/actions/src/utils/cubeCache.js | 12 +++++++++--- services/cubejs/src/routes/index.js | 18 +++++++++++++++++- 2 files changed, 26 insertions(+), 4 deletions(-) diff --git a/services/actions/src/utils/cubeCache.js b/services/actions/src/utils/cubeCache.js index e088694c..159db9ea 100644 --- a/services/actions/src/utils/cubeCache.js +++ b/services/actions/src/utils/cubeCache.js @@ -1,5 +1,11 @@ const CUBEJS_URL = process.env.CUBEJS_URL || "http://cubejs:4000"; +// cubejs rejects /internal/invalidate-cache without the shared admin secret. +const headers = { + "Content-Type": "application/json", + "x-hasura-admin-secret": process.env.HASURA_GRAPHQL_ADMIN_SECRET || "", +}; + /** * Invalidate CubeJS caches after admin mutations. * Fire-and-forget — never blocks the caller. @@ -7,7 +13,7 @@ const CUBEJS_URL = process.env.CUBEJS_URL || "http://cubejs:4000"; export function invalidateUserCache(userId) { fetch(`${CUBEJS_URL}/api/v1/internal/invalidate-cache`, { method: "POST", - headers: { "Content-Type": "application/json" }, + headers, body: JSON.stringify({ type: "user", userId }), }).catch(() => {}); } @@ -15,7 +21,7 @@ export function invalidateUserCache(userId) { export function invalidateAllUserCaches() { fetch(`${CUBEJS_URL}/api/v1/internal/invalidate-cache`, { method: "POST", - headers: { "Content-Type": "application/json" }, + headers, body: JSON.stringify({ type: "user" }), }).catch(() => {}); } @@ -23,7 +29,7 @@ export function invalidateAllUserCaches() { export function invalidateRulesCache() { fetch(`${CUBEJS_URL}/api/v1/internal/invalidate-cache`, { method: "POST", - headers: { "Content-Type": "application/json" }, + headers, body: JSON.stringify({ type: "rules" }), }).catch(() => {}); } diff --git a/services/cubejs/src/routes/index.js b/services/cubejs/src/routes/index.js index 57fe8b2f..fc7003be 100644 --- a/services/cubejs/src/routes/index.js +++ b/services/cubejs/src/routes/index.js @@ -1,3 +1,4 @@ +import { timingSafeEqual } from "node:crypto"; import express from "express"; // Model Management API (feature 011-model-mgmt-api) adds six routes registered @@ -56,6 +57,17 @@ import version from "./version.js"; const router = express.Router(); +// Service-to-service shared secret for /internal/invalidate-cache. Refuses +// everything when HASURA_GRAPHQL_ADMIN_SECRET is unset. +const hasAdminSecret = (req) => { + const expected = process.env.HASURA_GRAPHQL_ADMIN_SECRET; + const given = req.headers["x-hasura-admin-secret"]; + if (!expected || typeof given !== "string") return false; + const a = Buffer.from(given); + const b = Buffer.from(expected); + return a.length === b.length && timingSafeEqual(a, b); +}; + export default ({ basePath, cubejs }) => { router.get(`${basePath}/v1/load`, (req, res, next) => maybeHandleLoadExport(req, res, next, cubejs) @@ -203,8 +215,12 @@ export default ({ basePath, cubejs }) => { next(); }); - // Internal cache invalidation endpoint (called by Actions service, no auth) + // Internal cache invalidation endpoint (called by the Actions service with + // the Hasura admin secret — the route is reachable through the ingress). router.post(`${basePath}/v1/internal/invalidate-cache`, (req, res) => { + if (!hasAdminSecret(req)) { + return res.status(401).json({ code: "unauthorized", message: "Unauthorized" }); + } const { type, userId } = req.body || {}; if (type === "user") { From f8283be21743a7925fae8bb8e31f37c294fdad1a Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:34:44 +0000 Subject: [PATCH 07/12] fix(cubejs): version diff reports every changed attribute, incl. views The diff only flagged a member as modified when its type changed; sql, title, description, format, meta, primary_key, public, segments, joins, pre_aggregations and cube-level sql/sql_table/refresh_key/extends changes were invisible, and views were ignored. Cubes and views are now compared attribute by attribute (named lists member by member). Additive response fields: modifiedCubes[].changedAttributes (e.g. sql_table, measures.revenue.sql, joins.customers), kind (cube|view) on every entry, and modifiedFiles[] listing every file whose code differs. A cube added to an existing file is reported in addedCubes. diffModels' smart-generate preview semantics are unchanged; parseCubesFromJs can optionally collect view() definitions. Co-Authored-By: Claude Opus 5.5 --- .../src/utils/__tests__/versionDiff.test.js | 102 ++++++++ .../src/utils/smart-generation/diffModels.js | 7 +- services/cubejs/src/utils/versionDiff.js | 219 +++++++++++------- .../contracts/version-diff.yaml | 16 ++ 4 files changed, 258 insertions(+), 86 deletions(-) diff --git a/services/cubejs/src/utils/__tests__/versionDiff.test.js b/services/cubejs/src/utils/__tests__/versionDiff.test.js index 762d2a50..365292c4 100644 --- a/services/cubejs/src/utils/__tests__/versionDiff.test.js +++ b/services/cubejs/src/utils/__tests__/versionDiff.test.js @@ -113,4 +113,106 @@ describe("diffVersions", () => { assert.ok(measures, "expected a measures-level change entry"); assert.ok(measures.added.includes("revenue")); }); + + it("reports attribute-level changes the old type-only diff ignored", () => { + const edited = ordersV1 + .replace("sql_table: public.orders", "sql_table: public.orders_v2") + .replace( + " sql: id\n primary_key: true", + " sql: id\n title: Order id\n primary_key: true" + ); + const out = diffVersions({ + fromDataschemas: [mkRow("orders.yml", ordersV1)], + toDataschemas: [mkRow("orders.yml", edited)], + }); + assert.equal(out.modifiedCubes.length, 1); + const mod = out.modifiedCubes[0]; + assert.equal(mod.kind, "cube"); + assert.deepEqual(mod.changedAttributes.sort(), [ + "dimensions.id.title", + "sql_table", + ]); + const dims = mod.changes.find((c) => c.field === "dimensions"); + assert.deepEqual(dims.modified, ["id"]); + assert.deepEqual(out.modifiedFiles, [ + { file: "orders.yml", cubeNames: ["orders"] }, + ]); + }); + + it("diffs joins, pre-aggregations and views member by member", () => { + const from = `cubes: + - name: orders + sql_table: public.orders + joins: + - name: customers + sql: "{CUBE}.customer_id = {customers}.id" + relationship: many_to_one +views: + - name: sales + cubes: + - join_path: orders + includes: "*" +`; + const to = from + .replace("relationship: many_to_one", "relationship: one_to_one") + .replace('includes: "*"', "includes: [count]"); + const out = diffVersions({ + fromDataschemas: [mkRow("m.yml", from)], + toDataschemas: [mkRow("m.yml", to)], + }); + const byName = Object.fromEntries( + out.modifiedCubes.map((c) => [c.cubeName, c]) + ); + assert.deepEqual(byName.orders.changedAttributes, [ + "joins.customers.relationship", + ]); + assert.equal(byName.sales.kind, "view"); + assert.deepEqual(byName.sales.changedAttributes, ["cubes"]); + }); + + it("reports cubes added inside an existing file as added", () => { + const out = diffVersions({ + fromDataschemas: [mkRow("all.yml", ordersV1)], + toDataschemas: [mkRow("all.yml", ordersV1 + customers.slice("cubes:\n".length))], + }); + assert.deepEqual(out.addedCubes, [ + { cubeName: "customers", file: "all.yml", kind: "cube" }, + ]); + assert.deepEqual(out.modifiedCubes, []); + }); + + it("always lists a file whose code differs, even without a semantic change", () => { + const out = diffVersions({ + fromDataschemas: [mkRow("orders.yml", ordersV1)], + toDataschemas: [mkRow("orders.yml", `# note\n${ordersV1}`)], + }); + assert.deepEqual(out.modifiedCubes, []); + assert.deepEqual(out.modifiedFiles, [{ file: "orders.yml", cubeNames: [] }]); + }); + + it("parses JS files that also declare views", () => { + const js = (title) => `cube(\`orders\`, { + sql_table: \`public.orders\`, + measures: { count: { type: \`count\`, title: \`${title}\` } }, +}); +view(\`sales\`, { cubes: [{ join_path: \`orders\`, includes: \`*\` }] }); +`; + const out = diffVersions({ + fromDataschemas: [mkRow("orders.js", js("Orders"))], + toDataschemas: [mkRow("orders.js", js("Order count"))], + }); + assert.equal(out.modifiedCubes.length, 1); + assert.deepEqual(out.modifiedCubes[0].changedAttributes, [ + "measures.count.title", + ]); + }); + + it("does not report cubes of an unparseable file as removed", () => { + const out = diffVersions({ + fromDataschemas: [mkRow("orders.yml", ordersV1)], + toDataschemas: [mkRow("orders.yml", "cubes: [\n broken")], + }); + assert.deepEqual(out.removedCubes, []); + assert.deepEqual(out.modifiedFiles, [{ file: "orders.yml", cubeNames: [] }]); + }); }); diff --git a/services/cubejs/src/utils/smart-generation/diffModels.js b/services/cubejs/src/utils/smart-generation/diffModels.js index 140c408c..b28eca7d 100644 --- a/services/cubejs/src/utils/smart-generation/diffModels.js +++ b/services/cubejs/src/utils/smart-generation/diffModels.js @@ -311,9 +311,11 @@ function objectFieldsToArray(fields) { /** * Parse cube definitions from a JS cube file string. - * Evaluates in a VM sandbox with mock cube() function. + * Evaluates in a VM sandbox with mock cube() function. Pass a `views` array + * to also collect view() definitions into it (without one, a file calling + * view() stays unparseable, as before). */ -export function parseCubesFromJs(jsContent) { +export function parseCubesFromJs(jsContent, views = null) { const cubes = []; const mockCube = (name, def) => { @@ -331,6 +333,7 @@ export function parseCubesFromJs(jsContent) { try { const context = createContext({ cube: mockCube, + ...(views && { view: (name, def) => views.push({ name, ...def }) }), CUBE: '{CUBE}', FILTER_PARAMS: createDeepProxy(), SQL_UTILS: createDeepProxy(), diff --git a/services/cubejs/src/utils/versionDiff.js b/services/cubejs/src/utils/versionDiff.js index 888ac271..aa88dcb5 100644 --- a/services/cubejs/src/utils/versionDiff.js +++ b/services/cubejs/src/utils/versionDiff.js @@ -1,108 +1,131 @@ import YAML from "yaml"; -import { - parseCubesFromJs, - diffModels, -} from "./smart-generation/diffModels.js"; +import { parseCubesFromJs } from "./smart-generation/diffModels.js"; -function parseCubes(name, code) { +/** + * Parse a dataschema file into its cubes and views: + * `[{kind: 'cube'|'view', name, def}]`, or `null` when unparseable. + */ +function parseModels(name, code) { if (!code) return []; const isYaml = name?.endsWith(".yml") || name?.endsWith(".yaml"); try { + let cubes; + let views; if (isYaml) { const parsed = YAML.parse(code); - return Array.isArray(parsed?.cubes) ? parsed.cubes : []; + cubes = Array.isArray(parsed?.cubes) ? parsed.cubes : []; + views = Array.isArray(parsed?.views) ? parsed.views : []; + } else { + views = []; + cubes = parseCubesFromJs(code, views); + if (!cubes && !views.length) return null; // unparseable (or empty) JS + cubes = cubes || []; } - const cubes = parseCubesFromJs(code); - return Array.isArray(cubes) ? cubes : []; + return [ + ...cubes.map((def) => ({ kind: "cube", name: def?.name, def })), + ...views.map((def) => ({ kind: "view", name: def?.name, def })), + ].filter((m) => typeof m.name === "string"); } catch { - return []; + return null; + } +} + +// Order-insensitive for object keys (YAML key order is not semantic); +// functions (JS models) compare by source. +function canon(v) { + if (typeof v === "function") return v.toString(); + if (Array.isArray(v)) return v.map(canon); + if (v && typeof v === "object") { + return Object.fromEntries( + Object.keys(v) + .sort() + .map((k) => [k, canon(v[k])]) + ); } + return v; } +const same = (a, b) => JSON.stringify(canon(a)) === JSON.stringify(canon(b)); + +const isNamedList = (v) => + Array.isArray(v) && + v.every((x) => x && typeof x === "object" && typeof x.name === "string"); + +const byName = (list) => new Map((list || []).map((x) => [x.name, x])); + +// Member lists reported through the contract's `changes[]` entries. +const CHANGE_FIELDS = new Set(["dimensions", "measures", "segments"]); /** - * Group a flat `diffModels` result into per-cube `CubeChange` records. - * - * `diffModels` returns `{fields_added, fields_updated, fields_removed}` arrays - * where every entry carries the `cube` attribute identifying which cube it - * belongs to. This helper re-indexes those flat arrays into the per-cube - * shape required by contracts/version-diff.yaml (`CubeChange.changes[]`). + * Compare two definitions of the same cube/view. Named lists (dimensions, + * measures, segments, joins, pre_aggregations, hierarchies, …) are diffed + * member by member; every other key (sql, sql_table, refresh_key, extends, + * title, description, meta, public, a view's `cubes`, …) is compared whole. * - * @param {string} fileName - * @param {string} fromCode - * @param {string} toCode - * @returns {Array<{cubeName:string, file:string, changes:Array}>} + * @returns {{changes: Array, changedAttributes: string[]}} + * `changedAttributes` holds cube-level keys (`sql_table`) and member-level + * paths (`dimensions..`; `measures.` when added/removed). */ -function diffFilePair(fileName, fromCode, toCode) { - const flat = diffModels(fromCode || "", toCode || "", "replace"); - const byCube = new Map(); +function diffDefinitions(from, to) { + const changes = []; + const changedAttributes = []; + const keys = [...new Set([...Object.keys(from), ...Object.keys(to)])]; - const ensure = (cubeName) => { - if (!byCube.has(cubeName)) { - byCube.set(cubeName, new Map()); - } - return byCube.get(cubeName); - }; - - const bucket = (cubeName, memberType) => { - const cube = ensure(cubeName); - if (!cube.has(memberType)) { - cube.set(memberType, { added: [], removed: [], modified: [] }); + for (const key of keys) { + if (key === "name") continue; + const a = from[key]; + const b = to[key]; + const memberWise = + (a === undefined || isNamedList(a)) && + (b === undefined || isNamedList(b)); + if (!memberWise) { + if (!same(a, b)) changedAttributes.push(key); + continue; } - return cube.get(memberType); - }; - for (const entry of flat.fields_added || []) { - if (!entry?.cube) continue; - const b = bucket(entry.cube, entry.member_type || "meta"); - b.added.push(entry.name); - } - for (const entry of flat.fields_removed || []) { - if (!entry?.cube) continue; - const b = bucket(entry.cube, entry.member_type || "meta"); - b.removed.push(entry.name); - } - for (const entry of flat.fields_updated || []) { - if (!entry?.cube) continue; - const b = bucket(entry.cube, entry.member_type || "meta"); - b.modified.push(entry.name); - } - - const cubes = []; - for (const [cubeName, members] of byCube) { - const changes = []; - for (const [memberType, diff] of members) { - const hasAny = - diff.added.length || diff.removed.length || diff.modified.length; - if (!hasAny) continue; - changes.push({ - field: memberType === "measure" - ? "measures" - : memberType === "dimension" - ? "dimensions" - : memberType === "segment" - ? "segments" - : "meta", - added: diff.added, - removed: diff.removed, - modified: diff.modified, - }); + const fromMembers = byName(a); + const toMembers = byName(b); + const change = { field: key, added: [], removed: [], modified: [] }; + for (const [member, def] of toMembers) { + const old = fromMembers.get(member); + if (!old) { + change.added.push(member); + changedAttributes.push(`${key}.${member}`); + continue; + } + const attrs = [...new Set([...Object.keys(old), ...Object.keys(def)])]; + const changed = attrs.filter((x) => x !== "name" && !same(old[x], def[x])); + if (changed.length) change.modified.push(member); + for (const attr of changed) { + changedAttributes.push(`${key}.${member}.${attr}`); + } } - if (changes.length > 0) { - cubes.push({ cubeName, file: fileName, changes }); + for (const member of fromMembers.keys()) { + if (!toMembers.has(member)) { + change.removed.push(member); + changedAttributes.push(`${key}.${member}`); + } } + const hasAny = + change.added.length || change.removed.length || change.modified.length; + if (hasAny && CHANGE_FIELDS.has(key)) changes.push(change); } - return cubes; + + return { changes, changedAttributes }; } /** * Diff two versions (identified by their dataschema arrays) into the * `{addedCubes, removedCubes, modifiedCubes}` shape demanded by FR-011 - * and contracts/version-diff.yaml. + * and contracts/version-diff.yaml. Views are reported like cubes, with + * `kind: 'view'`. * - * Matching is by dataschema `name` (the file name) — a cube is "added" when - * its file is absent from `fromDataschemas` and "removed" when its file is - * absent from `toDataschemas`. Byte-identical files are skipped. + * Files match by dataschema `name`; cubes/views match by name within a + * file. A cube is "added"/"removed" when it appears on one side only. + * `modifiedCubes[].changedAttributes` lists every changed cube-level key + * and member attribute. `modifiedFiles` lists EVERY file present in both + * versions whose code differs — `cubeNames` is empty when no semantic + * change was found (formatting/comments) or the file does not parse. * * @param {object} args * @param {Array<{id?:string, name:string, code:string, checksum?:string}>} args.fromDataschemas @@ -121,11 +144,13 @@ export function diffVersions({ fromDataschemas, toDataschemas }) { const addedCubes = []; const removedCubes = []; const modifiedCubes = []; + const modifiedFiles = []; + const entry = (m, file) => ({ cubeName: m.name, file, kind: m.kind }); for (const [file, toRow] of toByFile) { if (!fromByFile.has(file)) { - for (const cube of parseCubes(file, toRow.code)) { - addedCubes.push({ cubeName: cube.name, file }); + for (const m of parseModels(file, toRow.code) || []) { + addedCubes.push(entry(m, file)); } continue; } @@ -139,17 +164,43 @@ export function diffVersions({ fromDataschemas, toDataschemas }) { } if (fromRow.code === toRow.code) continue; - const perCube = diffFilePair(file, fromRow.code, toRow.code); - for (const cube of perCube) modifiedCubes.push(cube); + const cubeNames = []; + const fromModels = parseModels(file, fromRow.code); + const toModels = parseModels(file, toRow.code); + if (fromModels && toModels) { + const key = (m) => `${m.kind}:${m.name}`; + const fromMap = new Map(fromModels.map((m) => [key(m), m])); + const toMap = new Map(toModels.map((m) => [key(m), m])); + for (const [k, m] of toMap) { + const old = fromMap.get(k); + if (!old) { + addedCubes.push(entry(m, file)); + cubeNames.push(m.name); + continue; + } + const diff = diffDefinitions(old.def || {}, m.def || {}); + if (diff.changedAttributes.length) { + modifiedCubes.push({ ...entry(m, file), ...diff }); + cubeNames.push(m.name); + } + } + for (const [k, m] of fromMap) { + if (!toMap.has(k)) { + removedCubes.push(entry(m, file)); + cubeNames.push(m.name); + } + } + } + modifiedFiles.push({ file, cubeNames }); } for (const [file, fromRow] of fromByFile) { if (!toByFile.has(file)) { - for (const cube of parseCubes(file, fromRow.code)) { - removedCubes.push({ cubeName: cube.name, file }); + for (const m of parseModels(file, fromRow.code) || []) { + removedCubes.push(entry(m, file)); } } } - return { addedCubes, removedCubes, modifiedCubes }; + return { addedCubes, removedCubes, modifiedCubes, modifiedFiles }; } diff --git a/specs/011-model-mgmt-api/contracts/version-diff.yaml b/specs/011-model-mgmt-api/contracts/version-diff.yaml index 8dec0c41..434a9b7a 100644 --- a/specs/011-model-mgmt-api/contracts/version-diff.yaml +++ b/specs/011-model-mgmt-api/contracts/version-diff.yaml @@ -102,6 +102,7 @@ components: properties: cubeName: { type: string } file: { type: string } + kind: { type: string, enum: [cube, view] } removedCubes: type: array items: @@ -110,6 +111,16 @@ components: properties: cubeName: { type: string } file: { type: string } + kind: { type: string, enum: [cube, view] } + modifiedFiles: + type: array + description: Every file in both versions whose code differs; `cubeNames` is empty when no semantic change was found (formatting/comments) or the file does not parse. + items: + type: object + required: [file, cubeNames] + properties: + file: { type: string } + cubeNames: { type: array, items: { type: string } } modifiedCubes: type: array items: @@ -118,6 +129,11 @@ components: properties: cubeName: { type: string } file: { type: string } + kind: { type: string, enum: [cube, view] } + changedAttributes: + type: array + description: Changed cube-level keys (`sql_table`, `refresh_key`, …) and member paths (`dimensions..`; `joins.` when added/removed). + items: { type: string } changes: type: array items: From a168d0e43d5491f42de9d1e098af0d0851e10311 Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:34:44 +0000 Subject: [PATCH 08/12] fix(cubejs): owner/admin required to save generated models MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit generate-models and non-dry-run smart-generate write versions with the admin secret and trusted body.branchId, so any team member could write models, and a caller could write onto another team's branch. Both now require the branch to belong to the request's datasource and — except smart-generate dry runs — the caller to be owner/admin of its team. The row-type pipeline identity is provisioned as admin, so it keeps working. Co-Authored-By: Claude Opus 5.5 --- .../cubejs/src/routes/generateDataSchema.js | 21 ++++-- services/cubejs/src/routes/smartGenerate.js | 14 ++++ .../src/utils/__tests__/versionWrites.test.js | 75 ++++++++++++++++++- .../cubejs/src/utils/dataSourceHelpers.js | 60 +++++++++++++++ 4 files changed, 161 insertions(+), 9 deletions(-) diff --git a/services/cubejs/src/routes/generateDataSchema.js b/services/cubejs/src/routes/generateDataSchema.js index 60680d8d..92cb1c9e 100644 --- a/services/cubejs/src/routes/generateDataSchema.js +++ b/services/cubejs/src/routes/generateDataSchema.js @@ -1,6 +1,7 @@ import { ScaffoldingTemplate } from "@cubejs-backend/schema-compiler"; import yaml from "js-yaml"; import { + authorizeModelWrite, createDataSchema, findDataSchemas, } from "../utils/dataSourceHelpers.js"; @@ -106,17 +107,25 @@ export default async (req, res, cubejs) => { userId, }; + const { + tables = [], + overwrite = false, + branchId, + format = "yaml", + } = req.body || {}; + + const denied = await authorizeModelWrite({ userId, dataSourceId, branchId }); + if (denied) { + return res + .status(denied.status) + .json({ code: denied.code, message: denied.message }); + } + let driver; try { driver = await tenantDriverFactory(cubejs)({ securityContext }); let schema = removeLegacyEnrichmentSchema(await driver.tablesSchema()); - const { - tables = [], - overwrite = false, - branchId, - format = "yaml", - } = req.body || {}; const { tables: normalizedTables, schema: normalizedSchema } = normalizeTables(schema, tables); diff --git a/services/cubejs/src/routes/smartGenerate.js b/services/cubejs/src/routes/smartGenerate.js index ed084fea..224482ba 100644 --- a/services/cubejs/src/routes/smartGenerate.js +++ b/services/cubejs/src/routes/smartGenerate.js @@ -1,4 +1,5 @@ import { + authorizeModelWrite, createDataSchema, findDataSchemas, } from "../utils/dataSourceHelpers.js"; @@ -215,6 +216,19 @@ export default async (req, res, cubejs) => { }); } + // Dry runs (previews) stay open to members; saves need owner/admin. + const denied = await authorizeModelWrite({ + userId: securityContext.userId, + dataSourceId: securityContext.userScope?.dataSource?.dataSourceId, + branchId, + dryRun, + }); + if (denied) { + return res + .status(denied.status) + .json({ code: denied.code, message: denied.message }); + } + let driver; try { diff --git a/services/cubejs/src/utils/__tests__/versionWrites.test.js b/services/cubejs/src/utils/__tests__/versionWrites.test.js index f5ff635b..ed0d6ab0 100644 --- a/services/cubejs/src/utils/__tests__/versionWrites.test.js +++ b/services/cubejs/src/utils/__tests__/versionWrites.test.js @@ -6,13 +6,38 @@ mock.module("../graphql.js", { namedExports: { fetchGraphQL: fetchGraphQLMock }, }); -const { commitVersionFiles, rollbackVersion } = await import( - "../dataSourceHelpers.js" -); +const { + authorizeModelWrite, + commitVersionFiles, + invalidateUserCache, + rollbackVersion, +} = await import("../dataSourceHelpers.js"); +const TEAM = "team-1"; const DS = "ds-1"; const BRANCH = "branch-1"; +function userQueryResult(teamRole) { + return { + data: { + members: [ + { + id: "m-1", + team_id: TEAM, + team: { + name: "t", + settings: {}, + datasources: [ + { id: DS, team_id: TEAM, branches: [{ id: BRANCH, versions: [] }] }, + ], + }, + member_roles: [{ id: "r-1", team_role: teamRole }], + }, + ], + }, + }; +} + describe("commitVersionFiles / rollbackVersion", () => { beforeEach(() => fetchGraphQLMock.mock.resetCalls()); @@ -79,3 +104,47 @@ describe("commitVersionFiles / rollbackVersion", () => { assert.equal(insert.source_version_id, "v-old"); }); }); + +describe("authorizeModelWrite", () => { + beforeEach(() => { + fetchGraphQLMock.mock.resetCalls(); + invalidateUserCache(null); + }); + + const as = (role) => + fetchGraphQLMock.mock.mockImplementation(async () => userQueryResult(role)); + + it("refuses a branch that is not on the request's datasource", async () => { + as("owner"); + const res = await authorizeModelWrite({ + userId: "u-1", + dataSourceId: DS, + branchId: "someone-elses-branch", + }); + assert.equal(res.status, 404); + }); + + it("refuses a plain member's write but allows the member's dry run", async () => { + as("member"); + const write = await authorizeModelWrite({ userId: "u-1", dataSourceId: DS, branchId: BRANCH }); + assert.equal(write.status, 403); + const dry = await authorizeModelWrite({ + userId: "u-1", + dataSourceId: DS, + branchId: BRANCH, + dryRun: true, + }); + assert.equal(dry, null); + }); + + it("allows owners and admins (the row-type pipeline identity is admin)", async () => { + for (const role of ["owner", "admin"]) { + invalidateUserCache(null); + as(role); + assert.equal( + await authorizeModelWrite({ userId: "u-1", dataSourceId: DS, branchId: BRANCH }), + null + ); + } + }); +}); diff --git a/services/cubejs/src/utils/dataSourceHelpers.js b/services/cubejs/src/utils/dataSourceHelpers.js index 8d3ce04a..e73f5477 100644 --- a/services/cubejs/src/utils/dataSourceHelpers.js +++ b/services/cubejs/src/utils/dataSourceHelpers.js @@ -3,6 +3,7 @@ import { createHash } from "crypto"; import { fetchGraphQL } from "./graphql.js"; import { fetchWorkOSUserProfile } from "./workosAuth.js"; import { emitModelEvent } from "./eventEmitter.js"; +import { requireOwnerOrAdmin } from "./requireOwnerOrAdmin.js"; // --- User scope cache: keyed by userId, 30s TTL --- const userCache = new Map(); @@ -286,6 +287,65 @@ export const findUser = async ({ userId }) => { return result; }; +/** + * Authorise a datasource-scoped model write (generate-models, smart-generate) + * onto the request body's `branchId`. Those routes write with the admin + * secret, so this is the only gate: the branch must belong to the request's + * datasource, and — unless `dryRun` — the caller must be owner/admin of the + * datasource's team (the row-type pipeline identity is provisioned as admin). + * + * Never throws. + * + * @returns {Promise} + */ +export const authorizeModelWrite = async ({ + userId, + dataSourceId, + branchId, + dryRun = false, +}) => { + const lookup = (user) => + user?.dataSources?.find( + (ds) => + ds.id === dataSourceId && + (ds.branches || []).some((b) => b.id === branchId) + ); + + let user; + let dataSource; + try { + user = await findUser({ userId }); + dataSource = lookup(user); + if (!dataSource) { + // A branch created moments ago may not be in the cached scope yet. + invalidateUserCache(userId); + user = await findUser({ userId }); + dataSource = lookup(user); + } + } catch (err) { + return { + status: 503, + code: "hasura_unavailable", + message: err?.message || "Hasura unavailable", + }; + } + if (!dataSource) { + return { + status: 404, + code: "branch_not_found", + message: `Branch "${branchId}" not found on this datasource`, + }; + } + if (!dryRun && !requireOwnerOrAdmin(user, dataSource.team_id)) { + return { + status: 403, + code: "owner_or_admin_required", + message: "Owner or admin role required to write models", + }; + } + return null; +}; + export const findSqlCredentials = async (username) => { const res = await fetchGraphQL(sqlCredentialsQuery, { username }); const sqlCredentials = res?.data?.sql_credentials?.[0]; From c4122888efcd62156fce99b8bbaa7696f84fec3e Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:35:27 +0000 Subject: [PATCH 09/12] test(actions): cubeCache sends the invalidate-cache admin secret Co-Authored-By: Claude Opus 5.5 --- .../src/utils/__tests__/cubeCache.test.js | 28 +++++++++++++++++++ 1 file changed, 28 insertions(+) create mode 100644 services/actions/src/utils/__tests__/cubeCache.test.js diff --git a/services/actions/src/utils/__tests__/cubeCache.test.js b/services/actions/src/utils/__tests__/cubeCache.test.js new file mode 100644 index 00000000..22fb7c11 --- /dev/null +++ b/services/actions/src/utils/__tests__/cubeCache.test.js @@ -0,0 +1,28 @@ +import { describe, it } from "node:test"; +import assert from "node:assert/strict"; + +describe("cubeCache", () => { + it("sends the admin secret cubejs requires on invalidate-cache", async () => { + process.env.HASURA_GRAPHQL_ADMIN_SECRET = "s3cret"; + const calls = []; + const originalFetch = globalThis.fetch; + globalThis.fetch = async (url, init) => { + calls.push({ url, init }); + return { ok: true }; + }; + try { + const { invalidateUserCache, invalidateAllUserCaches, invalidateRulesCache } = + await import("../cubeCache.js"); + invalidateUserCache("u-1"); + invalidateAllUserCaches(); + invalidateRulesCache(); + } finally { + globalThis.fetch = originalFetch; + } + assert.equal(calls.length, 3); + for (const { url, init } of calls) { + assert.match(url, /\/api\/v1\/internal\/invalidate-cache$/); + assert.equal(init.headers["x-hasura-admin-secret"], "s3cret"); + } + }); +}); From 20af604566eec947eed902533d37a9a76a532bc8 Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:46:04 +0000 Subject: [PATCH 10/12] refactor(cubejs): one version-write path, shared write guards - commitVersionFiles() is the only server-side version insert: PUT, DELETE, rollback and createDataSchema (now a thin adapter keeping its throw-on-error / {id} contract and caller checksum) all go through it, and it is the single place that emits Model Saved. The returning-less upsertVersionMutation twin is gone. - utils/modelWriteGuards.js (was mutableDataschema.js): authorizeTeamWrite() = partition + owner/admin gate with failure audit, shared by resolveMutableDataschema (PUT/DELETE) and version rollback. - hasuraTokenForUser() in mintHasuraToken.js replaces four copies of the cached get-or-mint (delete/put, rollback, validate-in-branch, hasura proxy); respondError() lives in errorCodes.js instead of per-route copies. Co-Authored-By: Claude Opus 5.5 --- CLAUDE.md | 2 +- .../cubejs/src/routes/deleteDataschema.js | 11 +- services/cubejs/src/routes/hasuraProxy.js | 15 +- services/cubejs/src/routes/refreshCompiler.js | 6 +- .../cubejs/src/routes/updateDataschema.js | 31 ++-- .../cubejs/src/routes/validateInBranch.js | 17 +-- services/cubejs/src/routes/versionDiff.js | 6 +- services/cubejs/src/routes/versionRollback.js | 83 +++-------- .../src/utils/__tests__/versionWrites.test.js | 34 ++++- .../cubejs/src/utils/dataSourceHelpers.js | 133 ++++++++++-------- services/cubejs/src/utils/errorCodes.js | 5 + services/cubejs/src/utils/mintHasuraToken.js | 20 +++ ...tableDataschema.js => modelWriteGuards.js} | 73 +++++----- 13 files changed, 207 insertions(+), 229 deletions(-) rename services/cubejs/src/utils/{mutableDataschema.js => modelWriteGuards.js} (76%) diff --git a/CLAUDE.md b/CLAUDE.md index 94bc7010..60048354 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -113,7 +113,7 @@ Docker Compose files per environment: `docker-compose.dev.yml`, `docker-compose. - `services/cubejs/src/routes/` — 13 REST API endpoints, now including: - `validateInBranch.js` (POST /api/v1/validate-in-branch, US1) - `refreshCompiler.js` (POST /api/v1/internal/refresh-compiler, US2) - - `deleteDataschema.js` (DELETE /api/v1/dataschema/:id, US3) + `updateDataschema.js` (PUT, same path) — both write a NEW version via `utils/mutableDataschema.js` + `commitVersionFiles` + - `deleteDataschema.js` (DELETE /api/v1/dataschema/:id, US3) + `updateDataschema.js` (PUT, same path) — share guards in `utils/modelWriteGuards.js`; every server-side version write (these, rollback, createDataSchema) goes through `commitVersionFiles` in dataSourceHelpers - `metaSingleCube.js` (GET /api/v1/meta/cube/:cubeName, US4) - `versionDiff.js` + `versionRollback.js` (POST /api/v1/version/{diff,rollback}, US5) - `services/cubejs/src/routes/reconcileTeam.js` — POST /api/v1/internal/reconcile-team: per-team default-models worker (013 — probe→generate→merge→validate→publish, system-user only) diff --git a/services/cubejs/src/routes/deleteDataschema.js b/services/cubejs/src/routes/deleteDataschema.js index 8340ef9b..689ba5d3 100644 --- a/services/cubejs/src/routes/deleteDataschema.js +++ b/services/cubejs/src/routes/deleteDataschema.js @@ -5,14 +5,11 @@ import { commitVersionFiles, findVersionDataschemas, } from "../utils/dataSourceHelpers.js"; -import { - ensureHasuraTokenForUser, - resolveMutableDataschema, - respondError, -} from "../utils/mutableDataschema.js"; +import { resolveMutableDataschema } from "../utils/modelWriteGuards.js"; +import { hasuraTokenForUser } from "../utils/mintHasuraToken.js"; import { scanCrossCubeReferences } from "../utils/referenceScanner.js"; import { mapHasuraErrorCode } from "../utils/mapHasuraErrorCode.js"; -import { ErrorCode } from "../utils/errorCodes.js"; +import { ErrorCode, respondError } from "../utils/errorCodes.js"; import { parseCubesFromJs } from "../utils/smart-generation/diffModels.js"; function parseCubes(name, code) { @@ -117,7 +114,7 @@ export default async function deleteDataschema(req, res) { // defence per research R4). let hasuraToken; try { - hasuraToken = await ensureHasuraTokenForUser(userId); + hasuraToken = await hasuraTokenForUser(userId); } catch { return respondError( res, diff --git a/services/cubejs/src/routes/hasuraProxy.js b/services/cubejs/src/routes/hasuraProxy.js index be4dca5e..8d7d0c21 100644 --- a/services/cubejs/src/routes/hasuraProxy.js +++ b/services/cubejs/src/routes/hasuraProxy.js @@ -3,8 +3,7 @@ import { createProxyMiddleware } from "http-proxy-middleware"; import { detectTokenType, verifyWorkOSToken, verifyFraiOSToken } from "../utils/workosAuth.js"; import { provisionUserFromWorkOS, provisionUserFromFraiOS } from "../utils/dataSourceHelpers.js"; -import { mintHasuraToken } from "../utils/mintHasuraToken.js"; -import { mintedTokenCache } from "../utils/mintedTokenCache.js"; +import { hasuraTokenForUser } from "../utils/mintHasuraToken.js"; // Shared auth path with checkAuth.js — no code duplication. // provisionUserFromWorkOS() manages workosSubCache + inflightProvisions internally. @@ -85,17 +84,7 @@ export default function createHasuraProxy(config = {}) { userId = await provisionUserFromFraiOS(payload); } - // Check minted token cache - let hasuraToken = mintedTokenCache.get(userId); - if (!hasuraToken) { - hasuraToken = await mintHasuraToken(userId); - // Decode exp for cache storage - const parts = hasuraToken.split("."); - const decoded = JSON.parse( - Buffer.from(parts[1], "base64url").toString() - ); - mintedTokenCache.set(userId, hasuraToken, decoded.exp); - } + const hasuraToken = await hasuraTokenForUser(userId); // Swap the Authorization header for Hasura req.headers.authorization = `Bearer ${hasuraToken}`; diff --git a/services/cubejs/src/routes/refreshCompiler.js b/services/cubejs/src/routes/refreshCompiler.js index 47f886f9..068efc3e 100644 --- a/services/cubejs/src/routes/refreshCompiler.js +++ b/services/cubejs/src/routes/refreshCompiler.js @@ -5,11 +5,7 @@ import { resolvePartitionTeamIds } from "./discover.js"; import { requireOwnerOrAdmin } from "../utils/requireOwnerOrAdmin.js"; import defineUserScope from "../utils/defineUserScope.js"; import { invalidateCompilerForBranch } from "../utils/compilerCacheInvalidator.js"; -import { ErrorCode } from "../utils/errorCodes.js"; - -function respondError(res, status, code, message) { - return res.status(status).json({ code, message }); -} +import { ErrorCode, respondError } from "../utils/errorCodes.js"; /** * POST /api/v1/internal/refresh-compiler diff --git a/services/cubejs/src/routes/updateDataschema.js b/services/cubejs/src/routes/updateDataschema.js index bc5fff79..7d32db94 100644 --- a/services/cubejs/src/routes/updateDataschema.js +++ b/services/cubejs/src/routes/updateDataschema.js @@ -1,15 +1,11 @@ -import { emitModelEvent } from "../utils/eventEmitter.js"; import { commitVersionFiles, findVersionDataschemas, } from "../utils/dataSourceHelpers.js"; -import { - ensureHasuraTokenForUser, - resolveMutableDataschema, - respondError, -} from "../utils/mutableDataschema.js"; +import { resolveMutableDataschema } from "../utils/modelWriteGuards.js"; +import { hasuraTokenForUser } from "../utils/mintHasuraToken.js"; import { mapHasuraErrorCode } from "../utils/mapHasuraErrorCode.js"; -import { ErrorCode } from "../utils/errorCodes.js"; +import { ErrorCode, respondError } from "../utils/errorCodes.js"; const dataschemaShape = (row) => ({ id: row.id, @@ -86,7 +82,7 @@ export default async function updateDataschema(req, res) { let hasuraToken; try { - hasuraToken = await ensureHasuraTokenForUser(userId); + hasuraToken = await hasuraTokenForUser(userId); } catch { return respondError( res, @@ -104,6 +100,12 @@ export default async function updateDataschema(req, res) { datasourceId, files, authToken: hasuraToken, + // `Model Saved`, emitted by commitVersionFiles like every other save. + emit: { + accountId: payload?.accountId ?? null, + partition: payload?.partition ?? null, + userId, + }, }); } catch (err) { return respondError( @@ -146,19 +148,6 @@ export default async function updateDataschema(req, res) { checksum: saved?.checksum ?? null, }); - // 099 T087 (FR-091): same `Model Saved` persistence fact createDataSchema - // emits for every other server-side version write. Fire-and-forget. - emitModelEvent({ - event: "Model Saved", - accountId: payload?.accountId ?? null, - partition: payload?.partition ?? null, - userId, - modelId: result.newVersionId, - modelLabel: target.name || null, - status: "ok", - properties: { branch_id: branchId, origin: "user" }, - }); - return res.json({ versionId: result.newVersionId, branchId, diff --git a/services/cubejs/src/routes/validateInBranch.js b/services/cubejs/src/routes/validateInBranch.js index 197c7bce..0593d8ea 100644 --- a/services/cubejs/src/routes/validateInBranch.js +++ b/services/cubejs/src/routes/validateInBranch.js @@ -7,8 +7,7 @@ import { findDataSchemas, } from "../utils/dataSourceHelpers.js"; import { fetchGraphQL } from "../utils/graphql.js"; -import { mintHasuraToken } from "../utils/mintHasuraToken.js"; -import { mintedTokenCache } from "../utils/mintedTokenCache.js"; +import { hasuraTokenForUser } from "../utils/mintHasuraToken.js"; import { requireOwnerOrAdmin } from "../utils/requireOwnerOrAdmin.js"; import { scanCrossCubeReferences } from "../utils/referenceScanner.js"; import { resolvePartitionTeamIds } from "./discover.js"; @@ -77,18 +76,6 @@ function respondJson(res, status, body) { return res.status(status).json(body); } -async function ensureHasuraTokenForUser(userId) { - let hasuraToken = mintedTokenCache.get(userId); - if (hasuraToken) return hasuraToken; - hasuraToken = await mintHasuraToken(userId); - const parts = hasuraToken.split("."); - const payload = JSON.parse( - Buffer.from(parts[1], "base64url").toString() - ); - mintedTokenCache.set(userId, hasuraToken, payload.exp); - return hasuraToken; -} - /** * POST /api/v1/validate-in-branch * @@ -231,7 +218,7 @@ export default async function validateInBranch(req, res) { // the user-role select permission). let existing; try { - const hasuraToken = await ensureHasuraTokenForUser(userId); + const hasuraToken = await hasuraTokenForUser(userId); existing = await findDataSchemas({ branchId, authToken: hasuraToken }); } catch (err) { return respondJson(res, 503, { diff --git a/services/cubejs/src/routes/versionDiff.js b/services/cubejs/src/routes/versionDiff.js index 5cbb8982..3b2a8aa5 100644 --- a/services/cubejs/src/routes/versionDiff.js +++ b/services/cubejs/src/routes/versionDiff.js @@ -6,11 +6,7 @@ import { } from "../utils/dataSourceHelpers.js"; import { resolvePartitionTeamIds } from "./discover.js"; import { diffVersions } from "../utils/versionDiff.js"; -import { ErrorCode } from "../utils/errorCodes.js"; - -function respondError(res, status, code, message) { - return res.status(status).json({ code, message }); -} +import { ErrorCode, respondError } from "../utils/errorCodes.js"; /** * POST /api/v1/version/diff diff --git a/services/cubejs/src/routes/versionRollback.js b/services/cubejs/src/routes/versionRollback.js index 05399247..02454e95 100644 --- a/services/cubejs/src/routes/versionRollback.js +++ b/services/cubejs/src/routes/versionRollback.js @@ -5,28 +5,11 @@ import { findVersionBranch, rollbackVersion as rollbackHelper, } from "../utils/dataSourceHelpers.js"; -import { requireOwnerOrAdmin } from "../utils/requireOwnerOrAdmin.js"; -import { resolvePartitionTeamIds } from "./discover.js"; +import { authorizeTeamWrite } from "../utils/modelWriteGuards.js"; import { writeAuditLog } from "../utils/auditWriter.js"; import { mapHasuraErrorCode } from "../utils/mapHasuraErrorCode.js"; -import { mintHasuraToken } from "../utils/mintHasuraToken.js"; -import { mintedTokenCache } from "../utils/mintedTokenCache.js"; -import { ErrorCode } from "../utils/errorCodes.js"; - -async function ensureHasuraTokenForUser(userId) { - let tok = mintedTokenCache.get(userId); - if (tok) return tok; - tok = await mintHasuraToken(userId); - const decoded = JSON.parse( - Buffer.from(tok.split(".")[1], "base64url").toString() - ); - mintedTokenCache.set(userId, tok, decoded.exp); - return tok; -} - -function respondError(res, status, code, message, extra = {}) { - return res.status(status).json({ code, message, ...extra }); -} +import { hasuraTokenForUser } from "../utils/mintHasuraToken.js"; +import { ErrorCode, respondError } from "../utils/errorCodes.js"; /** * POST /api/v1/version/rollback @@ -96,50 +79,28 @@ export default async function versionRollback(req, res, cubejs) { } const user = await findUser({ userId }); - const partitionTeamIds = resolvePartitionTeamIds( - user.members, - payload.partition - ); - if (partitionTeamIds && !partitionTeamIds.has(meta.teamId)) { - await writeAuditLog({ - action: "version_rollback", - userId, - datasourceId: meta.datasourceId, - branchId, - targetId: toVersionId, - outcome: "failure", - errorCode: ErrorCode.ROLLBACK_BLOCKED_AUTHORIZATION, - payload: { reason: "partition_mismatch" }, - }); - return respondError( - res, - 403, - ErrorCode.ROLLBACK_BLOCKED_AUTHORIZATION, - "Caller's partition does not match the branch's team" - ); - } - if (!requireOwnerOrAdmin(user, meta.teamId)) { - await writeAuditLog({ - action: "version_rollback", - userId, - datasourceId: meta.datasourceId, - branchId, - targetId: toVersionId, - outcome: "failure", - errorCode: ErrorCode.ROLLBACK_BLOCKED_AUTHORIZATION, - payload: { reason: "insufficient_role" }, - }); - return respondError( - res, - 403, - ErrorCode.ROLLBACK_BLOCKED_AUTHORIZATION, - "Owner or admin role required" - ); - } + const allowed = await authorizeTeamWrite(res, { + user, + partition: payload.partition, + teamId: meta.teamId, + code: ErrorCode.ROLLBACK_BLOCKED_AUTHORIZATION, + audit: (outcome, errorCode, auditPayload) => + writeAuditLog({ + action: "version_rollback", + userId, + datasourceId: meta.datasourceId, + branchId, + targetId: toVersionId, + outcome, + errorCode, + payload: auditPayload, + }), + }); + if (!allowed) return; let hasuraToken; try { - hasuraToken = await ensureHasuraTokenForUser(userId); + hasuraToken = await hasuraTokenForUser(userId); } catch { return respondError( res, diff --git a/services/cubejs/src/utils/__tests__/versionWrites.test.js b/services/cubejs/src/utils/__tests__/versionWrites.test.js index ed0d6ab0..dc95854a 100644 --- a/services/cubejs/src/utils/__tests__/versionWrites.test.js +++ b/services/cubejs/src/utils/__tests__/versionWrites.test.js @@ -9,6 +9,7 @@ mock.module("../graphql.js", { const { authorizeModelWrite, commitVersionFiles, + createDataSchema, invalidateUserCache, rollbackVersion, } = await import("../dataSourceHelpers.js"); @@ -64,7 +65,8 @@ describe("commitVersionFiles / rollbackVersion", () => { const [, vars, token, opts] = fetchGraphQLMock.mock.calls[0].arguments; assert.equal(token, "tok"); assert.deepEqual(opts, { preserveErrors: true }); - assert.equal(vars.object.origin, "user"); + // origin left to the column default ('user') + assert.equal("origin" in vars.object, false); assert.equal("source_version_id" in vars.object, false); // Only insertable columns — never the source row's id/checksum. assert.deepEqual(vars.object.dataschemas.data, [ @@ -85,6 +87,36 @@ describe("commitVersionFiles / rollbackVersion", () => { assert.deepEqual(res, { errors }); }); + it("createDataSchema keeps its contract on top of the shared insert", async () => { + fetchGraphQLMock.mock.mockImplementation(async () => ({ + data: { insert_versions_one: { id: "v-new", dataschemas: [] } }, + })); + const out = await createDataSchema({ + branch_id: BRANCH, + user_id: "u-1", + checksum: "caller-checksum", + dataschemas: { + data: [{ name: "a.yml", code: "x", user_id: "u-1", datasource_id: DS }], + }, + }); + assert.deepEqual(out, { id: "v-new" }); + const [, vars, token] = fetchGraphQLMock.mock.calls[0].arguments; + assert.equal(token, undefined); // admin secret, as before + assert.equal(vars.object.checksum, "caller-checksum"); + assert.deepEqual(vars.object.dataschemas.data, [ + { name: "a.yml", code: "x", user_id: "u-1", datasource_id: DS }, + ]); + + fetchGraphQLMock.mock.mockImplementation(async () => ({ + data: null, + errors: [{ message: "boom" }], + })); + await assert.rejects( + createDataSchema({ branch_id: BRANCH, user_id: "u-1", checksum: "c", dataschemas: { data: [] } }), + (err) => err.status === 503 && /boom/.test(err.message) + ); + }); + it("rollback records origin=rollback and the restored version as source", async () => { fetchGraphQLMock.mock.mockImplementation(async (query) => query.includes("VersionDataschemas") diff --git a/services/cubejs/src/utils/dataSourceHelpers.js b/services/cubejs/src/utils/dataSourceHelpers.js index e73f5477..6a674eff 100644 --- a/services/cubejs/src/utils/dataSourceHelpers.js +++ b/services/cubejs/src/utils/dataSourceHelpers.js @@ -177,19 +177,9 @@ const branchSchemasQuery = ` } `; -const upsertVersionMutation = ` - mutation ($object: versions_insert_input!) { - insert_versions_one( - object: $object - ) { - id - } - } -`; - -// Same insert as upsertVersionMutation, returning the new rows so a caller -// can report the ids of the files it just wrote. -const insertVersionReturningMutation = ` +// The one version insert: returns the new rows so callers can report the +// ids of the files they just wrote. +const insertVersionMutation = ` mutation ($object: versions_insert_input!) { insert_versions_one(object: $object) { id @@ -360,38 +350,37 @@ export const getDataSources = async () => { return res; }; -export const createDataSchema = async (object) => { - // `emit` (optional) carries tenant attribution for the 099 T087 `Model Saved` - // lifecycle event — destructured OUT here alongside `authToken` so it never - // reaches the `versions_insert_input` mutation variable (unknown column). - const { authToken, emit, ...version } = object; - - let res = await fetchGraphQL( - upsertVersionMutation, - { object: version }, - authToken - ); - res = res?.data?.insert_versions_one; - - // 099 T087 (FR-091): this is THE server-side persistence chokepoint for a - // model version. Emit `Model Saved` (persistence fact) when a caller supplied - // tenant attribution and a version was actually created. Fire-and-forget; - // never throws / never blocks (FR-007). The pure editor save (raw GraphQL - // straight to Hasura) does not pass through here — it is covered by a Hasura - // event trigger (T087, tables.yaml), a separate mechanism. - if (emit && res?.id) { - emitModelEvent({ - event: "Model Saved", - accountId: emit.accountId ?? null, - partition: emit.partition ?? null, - userId: emit.userId ?? null, - modelId: res.id, - status: "ok", - properties: { branch_id: version.branch_id ?? null, origin: version.origin ?? "save" }, - }); +/** + * Legacy `versions_insert_input`-shaped entry point (generate-models, + * smart-generate, reconcile-team) over {@link commitVersionFiles}. Keeps its + * contract: throws (status 503) on Hasura errors, returns `{id}`. + */ +export const createDataSchema = async ({ + authToken, + emit, + branch_id, + user_id, + checksum, + origin, + dataschemas, +}) => { + const files = dataschemas?.data || []; + const res = await commitVersionFiles({ + branchId: branch_id, + userId: user_id, + datasourceId: files[0]?.datasource_id, + files, + origin, + checksum, + authToken, + emit, + }); + if (res.errors) { + const error = new Error(JSON.stringify(res.errors)); + error.status = 503; + throw error; } - - return res; + return { id: res.newVersionId }; }; export const findDataSchemas = async ({ branchId, authToken }) => { @@ -439,13 +428,18 @@ export const findVersionBranch = async ({ versionId }) => { }; /** - * Insert a new version on `branchId` holding exactly `files` (`[{name, code}]`). - * Every model change goes through a new version so the previous one stays - * restorable; `versions_flip_is_current_trg` makes the new row current. + * THE server-side write of a model version: insert a new version on + * `branchId` holding exactly `files` (`[{name, code}]`). Every model change + * goes through a new version so the previous one stays restorable; + * `versions_flip_is_current_trg` makes the new row current. * - * The caller's minted Hasura token is used so owner/admin permission policies - * are enforced at the database layer too. On success the caller's cached - * user scope is dropped so its next request resolves the new version. + * `authToken` (a caller's minted Hasura token) makes the user-role + * permission policies apply at the database layer too; without it the admin + * secret is used. `origin` defaults to the column default ('user'); + * `checksum` defaults to an md5 over the sorted files. On success the + * caller's cached user scope is dropped so its next request resolves the new + * version, and — when `emit` carries tenant attribution — `Model Saved` is + * emitted. * * Returns `{newVersionId, dataschemas}` on success, `{errors}` on any Hasura * permission/constraint failure so the caller can map the extensions.code via @@ -456,9 +450,11 @@ export const commitVersionFiles = async ({ userId, datasourceId, files, - origin = "user", + origin, + checksum, sourceVersionId = null, authToken, + emit, }) => { // The `set_public_dataschemas_checksum` BEFORE-INSERT trigger computes // dataschema-level checksums; we do NOT set one here or Hasura's insert @@ -472,24 +468,26 @@ export const commitVersionFiles = async ({ // Version-level checksum: md5 over the concatenated dataschema codes in // a stable order. Matches the `version.checksum` NOT NULL constraint. - const versionChecksum = md5OfCode( - [...files] - .sort((a, b) => (a.name < b.name ? -1 : a.name > b.name ? 1 : 0)) - .map((r) => `${r.name}:${r.code || ""}`) - .join("\n") - ); + const versionChecksum = + checksum || + md5OfCode( + [...files] + .sort((a, b) => (a.name < b.name ? -1 : a.name > b.name ? 1 : 0)) + .map((r) => `${r.name}:${r.code || ""}`) + .join("\n") + ); const object = { branch_id: branchId, user_id: userId, - origin, checksum: versionChecksum, dataschemas: { data }, }; + if (origin) object.origin = origin; if (sourceVersionId) object.source_version_id = sourceVersionId; const res = await fetchGraphQL( - insertVersionReturningMutation, + insertVersionMutation, { object }, authToken, { preserveErrors: true } @@ -509,6 +507,23 @@ export const commitVersionFiles = async ({ } invalidateUserCache(userId); + + // 099 T087 (FR-091): `Model Saved` persistence fact. Fire-and-forget; + // never throws / never blocks (FR-007). The pure editor save (raw GraphQL + // straight to Hasura) does not pass through here — it is covered by a Hasura + // event trigger (T087, tables.yaml), a separate mechanism. + if (emit) { + emitModelEvent({ + event: "Model Saved", + accountId: emit.accountId ?? null, + partition: emit.partition ?? null, + userId: emit.userId ?? null, + modelId: row.id, + status: "ok", + properties: { branch_id: branchId ?? null, origin: origin ?? "save" }, + }); + } + return { newVersionId: row.id, dataschemas: row.dataschemas || [] }; }; diff --git a/services/cubejs/src/utils/errorCodes.js b/services/cubejs/src/utils/errorCodes.js index 890928fb..988a35da 100644 --- a/services/cubejs/src/utils/errorCodes.js +++ b/services/cubejs/src/utils/errorCodes.js @@ -28,6 +28,11 @@ export const ErrorCode = Object.freeze({ export const ErrorCodeSet = Object.freeze(new Set(Object.values(ErrorCode))); +/** Send a Model-Management error body `{code, message, ...extra}`. */ +export function respondError(res, status, code, message, extra = {}) { + return res.status(status).json({ code, message, ...extra }); +} + export function isKnownErrorCode(code) { return ErrorCodeSet.has(code); } diff --git a/services/cubejs/src/utils/mintHasuraToken.js b/services/cubejs/src/utils/mintHasuraToken.js index ab0336af..b8686502 100644 --- a/services/cubejs/src/utils/mintHasuraToken.js +++ b/services/cubejs/src/utils/mintHasuraToken.js @@ -1,5 +1,7 @@ import { SignJWT } from "jose"; +import { mintedTokenCache } from "./mintedTokenCache.js"; + const { JWT_EXPIRES_IN, JWT_ALGORITHM, JWT_CLAIMS_NAMESPACE, JWT_KEY } = process.env; @@ -28,3 +30,21 @@ export async function mintHasuraToken(userId) { .setSubject(userId) .sign(secret); } + +/** + * Cached Hasura token for `userId`: reuse a cached, unexpired mint or mint + * and cache a new one. + * + * @param {string} userId + * @returns {Promise} + */ +export async function hasuraTokenForUser(userId) { + const cached = mintedTokenCache.get(userId); + if (cached) return cached; + const token = await mintHasuraToken(userId); + const { exp } = JSON.parse( + Buffer.from(token.split(".")[1], "base64url").toString() + ); + mintedTokenCache.set(userId, token, exp); + return token; +} diff --git a/services/cubejs/src/utils/mutableDataschema.js b/services/cubejs/src/utils/modelWriteGuards.js similarity index 76% rename from services/cubejs/src/utils/mutableDataschema.js rename to services/cubejs/src/utils/modelWriteGuards.js index 967171c5..7790a290 100644 --- a/services/cubejs/src/utils/mutableDataschema.js +++ b/services/cubejs/src/utils/modelWriteGuards.js @@ -1,12 +1,10 @@ import { verifyAndProvision } from "./directVerifyAuth.js"; import { findUser } from "./dataSourceHelpers.js"; import { fetchGraphQL } from "./graphql.js"; -import { mintHasuraToken } from "./mintHasuraToken.js"; -import { mintedTokenCache } from "./mintedTokenCache.js"; import { requireOwnerOrAdmin } from "./requireOwnerOrAdmin.js"; import { resolvePartitionTeamIds } from "../routes/discover.js"; import { writeAuditLog } from "./auditWriter.js"; -import { ErrorCode } from "./errorCodes.js"; +import { ErrorCode, respondError } from "./errorCodes.js"; const RESOLVE_TARGET_QUERY = ` query ResolveTargetDataschema($id: uuid!) { @@ -32,19 +30,29 @@ const RESOLVE_TARGET_QUERY = ` } `; -export function respondError(res, status, code, message, extra = {}) { - return res.status(status).json({ code, message, ...extra }); -} - -export async function ensureHasuraTokenForUser(userId) { - let tok = mintedTokenCache.get(userId); - if (tok) return tok; - tok = await mintHasuraToken(userId); - const decoded = JSON.parse( - Buffer.from(tok.split(".")[1], "base64url").toString() - ); - mintedTokenCache.set(userId, tok, decoded.exp); - return tok; +/** + * Partition gate + owner/admin gate shared by the Model-Management write + * routes (FR-015). On refusal writes a failure audit row via `audit`, + * responds 403 with `code`, and returns false. + * + * @param {import('express').Response} res + * @param {{user: object, partition: string|null, teamId: string, code: string, + * audit: (outcome:string, errorCode:string, payload:object) => Promise}} args + * @returns {Promise} + */ +export async function authorizeTeamWrite(res, { user, partition, teamId, code, audit }) { + const partitionTeamIds = resolvePartitionTeamIds(user.members, partition); + if (partitionTeamIds && !partitionTeamIds.has(teamId)) { + await audit("failure", code, { reason: "partition_mismatch" }); + respondError(res, 403, code, "Caller's partition does not match the datasource's team"); + return false; + } + if (!requireOwnerOrAdmin(user, teamId)) { + await audit("failure", code, { reason: "insufficient_role" }); + respondError(res, 403, code, "Owner or admin role required"); + return false; + } + return true; } /** @@ -131,31 +139,14 @@ export async function resolveMutableDataschema(req, res, { action, codes }) { }); const user = await findUser({ userId }); - - const partitionTeamIds = resolvePartitionTeamIds( - user.members, - payload.partition - ); - if (partitionTeamIds && !partitionTeamIds.has(teamId)) { - await audit("failure", codes.authorization, { - reason: "partition_mismatch", - }); - respondError( - res, - 403, - codes.authorization, - "Caller's partition does not match the datasource's team" - ); - return null; - } - - if (!requireOwnerOrAdmin(user, teamId)) { - await audit("failure", codes.authorization, { - reason: "insufficient_role", - }); - respondError(res, 403, codes.authorization, "Owner or admin role required"); - return null; - } + const allowed = await authorizeTeamWrite(res, { + user, + partition: payload.partition, + teamId, + code: codes.authorization, + audit, + }); + if (!allowed) return null; // Version-level immutability (FR-007): only the current version of the // active branch may be changed (by writing a new version on top of it). From 210bdbb835b18bff67bca5a3eee46d64c08f005f Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 15:51:03 +0000 Subject: [PATCH 11/12] fix(cubejs): version diff requires team membership MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit POST /api/v1/version/diff read both versions with the admin secret and only checked the partition claim, so a caller whose token carries no partition (Hasura HS256, WorkOS without partition) could diff any team's versions. It now requires the branch to be visible to the caller through the same guard the generation routes use — authorizeModelWrite renamed to authorizeBranchAccess({readOnly}) — and answers 404 otherwise. Co-Authored-By: Claude Opus 5.5 --- .../cubejs/src/routes/generateDataSchema.js | 4 ++-- services/cubejs/src/routes/smartGenerate.js | 6 ++--- services/cubejs/src/routes/versionDiff.js | 15 ++++++++++++ .../src/utils/__tests__/versionWrites.test.js | 24 ++++++++++++------- .../cubejs/src/utils/dataSourceHelpers.js | 17 ++++++------- 5 files changed, 45 insertions(+), 21 deletions(-) diff --git a/services/cubejs/src/routes/generateDataSchema.js b/services/cubejs/src/routes/generateDataSchema.js index 92cb1c9e..24804f3e 100644 --- a/services/cubejs/src/routes/generateDataSchema.js +++ b/services/cubejs/src/routes/generateDataSchema.js @@ -1,7 +1,7 @@ import { ScaffoldingTemplate } from "@cubejs-backend/schema-compiler"; import yaml from "js-yaml"; import { - authorizeModelWrite, + authorizeBranchAccess, createDataSchema, findDataSchemas, } from "../utils/dataSourceHelpers.js"; @@ -114,7 +114,7 @@ export default async (req, res, cubejs) => { format = "yaml", } = req.body || {}; - const denied = await authorizeModelWrite({ userId, dataSourceId, branchId }); + const denied = await authorizeBranchAccess({ userId, dataSourceId, branchId }); if (denied) { return res .status(denied.status) diff --git a/services/cubejs/src/routes/smartGenerate.js b/services/cubejs/src/routes/smartGenerate.js index 224482ba..0ae5ef31 100644 --- a/services/cubejs/src/routes/smartGenerate.js +++ b/services/cubejs/src/routes/smartGenerate.js @@ -1,5 +1,5 @@ import { - authorizeModelWrite, + authorizeBranchAccess, createDataSchema, findDataSchemas, } from "../utils/dataSourceHelpers.js"; @@ -217,11 +217,11 @@ export default async (req, res, cubejs) => { } // Dry runs (previews) stay open to members; saves need owner/admin. - const denied = await authorizeModelWrite({ + const denied = await authorizeBranchAccess({ userId: securityContext.userId, dataSourceId: securityContext.userScope?.dataSource?.dataSourceId, branchId, - dryRun, + readOnly: dryRun, }); if (denied) { return res diff --git a/services/cubejs/src/routes/versionDiff.js b/services/cubejs/src/routes/versionDiff.js index 3b2a8aa5..92409491 100644 --- a/services/cubejs/src/routes/versionDiff.js +++ b/services/cubejs/src/routes/versionDiff.js @@ -1,5 +1,6 @@ import { verifyAndProvision } from "../utils/directVerifyAuth.js"; import { + authorizeBranchAccess, findUser, findVersionBranch, findVersionDataschemas, @@ -92,6 +93,20 @@ export default async function versionDiff(req, res) { ); } + // Team membership: versions are read with the admin secret, so without a + // partition claim this is the only thing keeping other teams out. + const denied = await authorizeBranchAccess({ + userId, + dataSourceId: toMeta.datasourceId, + branchId: toMeta.branchId, + readOnly: true, + }); + if (denied) { + return denied.status === 404 + ? respondError(res, 404, ErrorCode.DIFF_INVALID_REQUEST, "One or both versions not found") + : respondError(res, denied.status, denied.code, denied.message); + } + let fromRows; let toRows; try { diff --git a/services/cubejs/src/utils/__tests__/versionWrites.test.js b/services/cubejs/src/utils/__tests__/versionWrites.test.js index dc95854a..52c8879f 100644 --- a/services/cubejs/src/utils/__tests__/versionWrites.test.js +++ b/services/cubejs/src/utils/__tests__/versionWrites.test.js @@ -7,7 +7,7 @@ mock.module("../graphql.js", { }); const { - authorizeModelWrite, + authorizeBranchAccess, commitVersionFiles, createDataSchema, invalidateUserCache, @@ -137,7 +137,7 @@ describe("commitVersionFiles / rollbackVersion", () => { }); }); -describe("authorizeModelWrite", () => { +describe("authorizeBranchAccess", () => { beforeEach(() => { fetchGraphQLMock.mock.resetCalls(); invalidateUserCache(null); @@ -148,23 +148,31 @@ describe("authorizeModelWrite", () => { it("refuses a branch that is not on the request's datasource", async () => { as("owner"); - const res = await authorizeModelWrite({ + const res = await authorizeBranchAccess({ userId: "u-1", dataSourceId: DS, branchId: "someone-elses-branch", }); assert.equal(res.status, 404); + // reads (version diff) are refused the same way + const read = await authorizeBranchAccess({ + userId: "u-1", + dataSourceId: DS, + branchId: "someone-elses-branch", + readOnly: true, + }); + assert.equal(read.status, 404); }); - it("refuses a plain member's write but allows the member's dry run", async () => { + it("refuses a plain member's write but allows the member's read", async () => { as("member"); - const write = await authorizeModelWrite({ userId: "u-1", dataSourceId: DS, branchId: BRANCH }); + const write = await authorizeBranchAccess({ userId: "u-1", dataSourceId: DS, branchId: BRANCH }); assert.equal(write.status, 403); - const dry = await authorizeModelWrite({ + const dry = await authorizeBranchAccess({ userId: "u-1", dataSourceId: DS, branchId: BRANCH, - dryRun: true, + readOnly: true, }); assert.equal(dry, null); }); @@ -174,7 +182,7 @@ describe("authorizeModelWrite", () => { invalidateUserCache(null); as(role); assert.equal( - await authorizeModelWrite({ userId: "u-1", dataSourceId: DS, branchId: BRANCH }), + await authorizeBranchAccess({ userId: "u-1", dataSourceId: DS, branchId: BRANCH }), null ); } diff --git a/services/cubejs/src/utils/dataSourceHelpers.js b/services/cubejs/src/utils/dataSourceHelpers.js index 6a674eff..90071483 100644 --- a/services/cubejs/src/utils/dataSourceHelpers.js +++ b/services/cubejs/src/utils/dataSourceHelpers.js @@ -278,21 +278,22 @@ export const findUser = async ({ userId }) => { }; /** - * Authorise a datasource-scoped model write (generate-models, smart-generate) - * onto the request body's `branchId`. Those routes write with the admin - * secret, so this is the only gate: the branch must belong to the request's - * datasource, and — unless `dryRun` — the caller must be owner/admin of the - * datasource's team (the row-type pipeline identity is provisioned as admin). + * Authorise access to `branchId` on `dataSourceId`: the branch must belong to + * that datasource and the datasource must be visible to the caller (team + * membership, via findUser). Unless `readOnly`, the caller must also be + * owner/admin of the datasource's team (the row-type pipeline identity is + * provisioned as admin). Used by routes that read or write through the admin + * secret, where this is the only gate. * * Never throws. * * @returns {Promise} */ -export const authorizeModelWrite = async ({ +export const authorizeBranchAccess = async ({ userId, dataSourceId, branchId, - dryRun = false, + readOnly = false, }) => { const lookup = (user) => user?.dataSources?.find( @@ -326,7 +327,7 @@ export const authorizeModelWrite = async ({ message: `Branch "${branchId}" not found on this datasource`, }; } - if (!dryRun && !requireOwnerOrAdmin(user, dataSource.team_id)) { + if (!readOnly && !requireOwnerOrAdmin(user, dataSource.team_id)) { return { status: 403, code: "owner_or_admin_required", From 56b3baad5861e4ae7ad2648ad5b512a7c4c28af2 Mon Sep 17 00:00:00 2001 From: stefanbaxter Date: Sun, 27 Sep 2026 16:00:09 +0000 Subject: [PATCH 12/12] fix(cubejs): build on EOL bullseye via the final security snapshot Debian 11 LTS ended 2026-08-31. deb.debian.org purged the bullseye-security pool (between Sep 3 and 5) but still serves its index, so apt-get install libssl1.1 resolves to 1.1.1w-0+deb11u8 and 404s on every arch; archive.debian.org has no bullseye-security yet. Read that suite from snapshot.debian.org's frozen final state (20260903T000000Z, index dated 2026-08-31) with check-valid-until=no, keeping the same final security versions the 2026-09-03 images shipped. Stop-gap: the base gets no further security updates; moving to a bookworm base is the real fix. Verified: buildx linux/amd64 and linux/arm64 builds pass, CI's in-image yarn test is 695/695 on both, amd64 image serves /livez on the vi stack. Co-Authored-By: Claude Opus 5.5 --- services/cubejs/Dockerfile | 10 ++++++++++ 1 file changed, 10 insertions(+) diff --git a/services/cubejs/Dockerfile b/services/cubejs/Dockerfile index 51c6efcf..c58853ad 100644 --- a/services/cubejs/Dockerfile +++ b/services/cubejs/Dockerfile @@ -2,6 +2,16 @@ FROM node:22.14.0-bullseye@sha256:1a956cf7ae435192dcdd18b8aba7349d649f0947795eac ARG DATABRICKS_JDBC_URL=https://databricks-bi-artifacts.s3.us-east-2.amazonaws.com/simbaspark-drivers/jdbc/2.6.32/DatabricksJDBC42-2.6.32.1054.zip +# Debian 11 (bullseye) LTS ended 2026-08-31. deb.debian.org has since purged +# the bullseye-security pool but still serves its index, so any install that +# resolves to a security update (libssl1.1 1.1.1w-0+deb11u8, ...) 404s. Read +# that suite from snapshot.debian.org's frozen copy of its final state (index +# dated 2026-08-31 21:13 UTC); its Valid-Until has passed, hence the flag. +# ponytail: stop-gap on an EOL base with no further security updates; move to +# a bookworm base (OpenSSL 3) to get updates again. +RUN sed -i 's#^deb http://deb.debian.org/debian-security bullseye-security main$#deb [check-valid-until=no] http://snapshot.debian.org/archive/debian-security/20260903T000000Z bullseye-security main#' /etc/apt/sources.list \ + && grep -q '^deb \[check-valid-until=no\] http://snapshot.debian.org/archive/debian-security/' /etc/apt/sources.list + RUN DEBIAN_FRONTEND=noninteractive \ && apt-get update \ && apt-get install -y --no-install-recommends rxvt-unicode libssl1.1 \