diff --git a/CLAUDE.md b/CLAUDE.md index 90546f39..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) + - `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/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/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"); + } + }); +}); 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/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 \ 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); diff --git a/services/cubejs/src/routes/deleteDataschema.js b/services/cubejs/src/routes/deleteDataschema.js index 0e9249d6..689ba5d3 100644 --- a/services/cubejs/src/routes/deleteDataschema.js +++ b/services/cubejs/src/routes/deleteDataschema.js @@ -1,65 +1,17 @@ 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 { resolveMutableDataschema } from "../utils/modelWriteGuards.js"; +import { hasuraTokenForUser } from "../utils/mintHasuraToken.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 { ErrorCode, respondError } 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 +27,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 +78,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 +97,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,12 +109,12 @@ 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 { - hasuraToken = await ensureHasuraTokenForUser(userId); + hasuraToken = await hasuraTokenForUser(userId); } catch { return respondError( res, @@ -297,14 +124,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 +142,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 +155,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 +164,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 +178,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/routes/generateDataSchema.js b/services/cubejs/src/routes/generateDataSchema.js index 60680d8d..24804f3e 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 { + authorizeBranchAccess, 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 authorizeBranchAccess({ 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/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/index.js b/services/cubejs/src/routes/index.js index 6edeff77..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 @@ -5,6 +6,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 +44,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"; @@ -54,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) @@ -201,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") { @@ -333,6 +351,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/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/smartGenerate.js b/services/cubejs/src/routes/smartGenerate.js index ed084fea..0ae5ef31 100644 --- a/services/cubejs/src/routes/smartGenerate.js +++ b/services/cubejs/src/routes/smartGenerate.js @@ -1,4 +1,5 @@ import { + authorizeBranchAccess, 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 authorizeBranchAccess({ + userId: securityContext.userId, + dataSourceId: securityContext.userScope?.dataSource?.dataSourceId, + branchId, + readOnly: dryRun, + }); + if (denied) { + return res + .status(denied.status) + .json({ code: denied.code, message: denied.message }); + } + let driver; try { diff --git a/services/cubejs/src/routes/updateDataschema.js b/services/cubejs/src/routes/updateDataschema.js new file mode 100644 index 00000000..7d32db94 --- /dev/null +++ b/services/cubejs/src/routes/updateDataschema.js @@ -0,0 +1,158 @@ +import { + commitVersionFiles, + findVersionDataschemas, +} from "../utils/dataSourceHelpers.js"; +import { resolveMutableDataschema } from "../utils/modelWriteGuards.js"; +import { hasuraTokenForUser } from "../utils/mintHasuraToken.js"; +import { mapHasuraErrorCode } from "../utils/mapHasuraErrorCode.js"; +import { ErrorCode, respondError } 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 hasuraTokenForUser(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, + // `Model Saved`, emitted by commitVersionFiles like every other save. + emit: { + accountId: payload?.accountId ?? null, + partition: payload?.partition ?? null, + userId, + }, + }); + } 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, + }); + + 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/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..92409491 100644 --- a/services/cubejs/src/routes/versionDiff.js +++ b/services/cubejs/src/routes/versionDiff.js @@ -1,16 +1,13 @@ import { verifyAndProvision } from "../utils/directVerifyAuth.js"; import { + authorizeBranchAccess, findUser, findVersionBranch, findVersionDataschemas, } 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 @@ -96,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/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__/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/__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/__tests__/versionWrites.test.js b/services/cubejs/src/utils/__tests__/versionWrites.test.js new file mode 100644 index 00000000..52c8879f --- /dev/null +++ b/services/cubejs/src/utils/__tests__/versionWrites.test.js @@ -0,0 +1,190 @@ +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 { + authorizeBranchAccess, + commitVersionFiles, + createDataSchema, + 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()); + + 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 }); + // 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, [ + { 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("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") + ? { 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"); + }); +}); + +describe("authorizeBranchAccess", () => { + 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 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 read", async () => { + as("member"); + const write = await authorizeBranchAccess({ userId: "u-1", dataSourceId: DS, branchId: BRANCH }); + assert.equal(write.status, 403); + const dry = await authorizeBranchAccess({ + userId: "u-1", + dataSourceId: DS, + branchId: BRANCH, + readOnly: 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 authorizeBranchAccess({ userId: "u-1", dataSourceId: DS, branchId: BRANCH }), + null + ); + } + }); +}); 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/dataSourceHelpers.js b/services/cubejs/src/utils/dataSourceHelpers.js index ae887eea..90071483 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(); @@ -176,12 +177,18 @@ const branchSchemasQuery = ` } `; -const upsertVersionMutation = ` +// 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 - ) { + insert_versions_one(object: $object) { id + dataschemas { + id + name + checksum + version_id + } } } `; @@ -270,6 +277,66 @@ export const findUser = async ({ userId }) => { return result; }; +/** + * 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 authorizeBranchAccess = async ({ + userId, + dataSourceId, + branchId, + readOnly = 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 (!readOnly && !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]; @@ -284,38 +351,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 }) => { @@ -363,29 +429,38 @@ 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. + * 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. * - * Returns `{newVersionId, clonedDataschemaCount}` on success, - * `{errors}` on any Hasura permission/constraint failure so the caller can - * map the extensions.code via mapHasuraErrorCode(). + * `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 + * mapHasuraErrorCode(). */ -export const rollbackVersion = async ({ +export const commitVersionFiles = async ({ branchId, - toVersionId, userId, datasourceId, + files, + origin, + checksum, + sourceVersionId = null, authToken, + emit, }) => { - 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, @@ -394,24 +469,26 @@ 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 || ""}`) - .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: "rollback", checksum: versionChecksum, - dataschemas: { data: clonedData }, + dataschemas: { data }, }; + if (origin) object.origin = origin; + if (sourceVersionId) object.source_version_id = sourceVersionId; const res = await fetchGraphQL( - upsertVersionMutation, + insertVersionMutation, { object }, authToken, { preserveErrors: true } @@ -421,8 +498,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 +507,59 @@ export const rollbackVersion = async ({ }; } - return { newVersionId, clonedDataschemaCount: clonedData.length }; + 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 || [] }; +}; + +/** + * 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) { diff --git a/services/cubejs/src/utils/errorCodes.js b/services/cubejs/src/utils/errorCodes.js index 6bd08422..988a35da 100644 --- a/services/cubejs/src/utils/errorCodes.js +++ b/services/cubejs/src/utils/errorCodes.js @@ -22,10 +22,17 @@ 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))); +/** 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/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/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/modelWriteGuards.js b/services/cubejs/src/utils/modelWriteGuards.js new file mode 100644 index 00000000..7790a290 --- /dev/null +++ b/services/cubejs/src/utils/modelWriteGuards.js @@ -0,0 +1,177 @@ +import { verifyAndProvision } from "./directVerifyAuth.js"; +import { findUser } from "./dataSourceHelpers.js"; +import { fetchGraphQL } from "./graphql.js"; +import { requireOwnerOrAdmin } from "./requireOwnerOrAdmin.js"; +import { resolvePartitionTeamIds } from "../routes/discover.js"; +import { writeAuditLog } from "./auditWriter.js"; +import { ErrorCode, respondError } 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 + } + } + } + } + } +`; + +/** + * 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; +} + +/** + * 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 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). + 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, + }; +} 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/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')); 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..434a9b7a 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] @@ -100,6 +102,7 @@ components: properties: cubeName: { type: string } file: { type: string } + kind: { type: string, enum: [cube, view] } removedCubes: type: array items: @@ -108,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: @@ -116,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: 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]