Skip to content

Commit b3d022a

Browse files
authored
[Analyze 2845] cdm sql audit logging (#2914)
* feat(analytics): add file-based SQL audit logging * feat(analytics): cover direct CDM SQL paths * fix(analytics): allow audit writes in Trex runtime * refactor(analytics): fix audit log directory * feat(analytics): support configurable audit output * refactor(trex): remove audit directory chmod * refactor(analytics-svc): remove SQL audit hash * Include dataset config metadata in SQL audit logs
1 parent e99f6b9 commit b3d022a

16 files changed

Lines changed: 1627 additions & 56 deletions

docker-compose.yml

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -17,6 +17,7 @@ volumes:
1717
hana-data:
1818
supabase-storage-data:
1919
hades-data:
20+
audit-logs:
2021

2122
x-envs:
2223
env_converter: &x-converter
@@ -95,6 +96,7 @@ services:
9596
- trex:/usr/src/data
9697
- cdw-config-cachedb-data-1:/usr/src/cdw_data/dynamically_generated
9798
- hades-data:/data/hades
99+
- audit-logs:/var/log/d2e/audit
98100
ports:
99101
- "8090:8080"
100102
healthcheck:
@@ -445,6 +447,8 @@ services:
445447
CACHE_TASK_TIMEOUT: ${CACHE_TASK_TIMEOUT:-10800}
446448
ANALYTICS_STREAMING_CHUNK_SIZE_BY_DIALECT: '${ANALYTICS_STREAMING_CHUNK_SIZE_BY_DIALECT:-{"hana": 10000}}'
447449
IS_AUDIT_LOG_ENABLED: ${IS_AUDIT_LOG_ENABLED:-false}
450+
IS_CDM_SQL_AUDIT_LOG_ENABLED: ${IS_CDM_SQL_AUDIT_LOG_ENABLED:-false}
451+
AUDIT_LOG_TO_CONSOLE: ${AUDIT_LOG_TO_CONSOLE:-false}
448452
ANALYTICS_HANA_STREAMING_ENABLED: ${ANALYTICS_HANA_STREAMING_ENABLED:-false}
449453
# TREX_DEBUG_GC: ${TREX_DEBUG_GC:-0}
450454
# RUST_LOG: ${RUST_LOG:-info}

plugins/functions/analytics-svc/src/api/controllers/cohort.ts

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -18,6 +18,7 @@ import PortalServerAPI from "../PortalServerAPI";
1818
import { convertIFRToExtCohort } from "../../ifr-to-extcohort/main";
1919
import { dataflowRequest } from "../../utils/DataflowMgmtProxy";
2020
import { env } from "../../env";
21+
import { createCdmSqlAuditContext } from "../../utils/CdmSqlAuditLogger.ts";
2122

2223
const language = "en";
2324

@@ -314,6 +315,17 @@ export async function createCohort(req: IMRIRequest, res: Response) {
314315
datasetId: req.selectedstudyDbMetadata.id,
315316
token: req.headers.authorization,
316317
dbCredential: req.dbCredentials.studyAnalyticsCredential,
318+
auditContext: createCdmSqlAuditContext({
319+
request: req,
320+
actorId: getUser(req)?.getUser() ?? "unknown",
321+
databaseCode:
322+
req.dbCredentials.studyAnalyticsCredential.code,
323+
databaseDialect:
324+
req.dbCredentials.studyAnalyticsCredential.dialect,
325+
databaseEngine: "hana",
326+
schemaName:
327+
req.dbCredentials.studyAnalyticsCredential.schema,
328+
}),
317329
}
318330
);
319331
} else {

plugins/functions/analytics-svc/src/api/controllers/datasetFilter.ts

Lines changed: 36 additions & 21 deletions
Original file line numberDiff line numberDiff line change
@@ -11,6 +11,10 @@ import QueryObject = qo.QueryObject;
1111
import { dataflowRequest } from "../../utils/DataflowMgmtProxy";
1212
import { FilterScopeQueryBuilder } from "../../utils/dataset-filter/filter-scope-query-builder";
1313
import { FilterQueryBuilder } from "../../utils/dataset-filter/query-builder";
14+
import {
15+
createCdmSqlAuditConnection,
16+
createCdmSqlAuditContext,
17+
} from "../../utils/CdmSqlAuditLogger.ts";
1418
import {
1519
DatabaseSchemaMap,
1620
IDatasetFilterScopesDto,
@@ -204,6 +208,36 @@ const createDatabaseSchemaMap = (
204208
}, {});
205209
};
206210

211+
const getDatasetFilterDbConnection = async (databaseName, req) => {
212+
const userObj = getUser(req);
213+
const dbCredentials = getDatabaseCredentials()[databaseName];
214+
const connection =
215+
await dbConnectionUtil.DBConnectionUtil.getDBConnection({
216+
credentials: dbCredentials,
217+
schemaName: dbCredentials.schemaName,
218+
vocabSchemaName: dbCredentials.vocabSchemaName,
219+
userObj,
220+
});
221+
const databaseEngine =
222+
dbCredentials.dialect === "hana"
223+
? "hana"
224+
: dbCredentials.dialect === "postgresql"
225+
? "postgresql"
226+
: "duckdb";
227+
228+
return createCdmSqlAuditConnection(
229+
connection,
230+
createCdmSqlAuditContext({
231+
request: req,
232+
actorId: userObj?.getUser() ?? "unknown",
233+
databaseCode: databaseName,
234+
databaseDialect: dbCredentials.dialect,
235+
databaseEngine,
236+
schemaName: dbCredentials.schema,
237+
})
238+
);
239+
};
240+
207241
const queryFilterScopes = async (
208242
databaseName,
209243
schemas,
@@ -213,16 +247,7 @@ const queryFilterScopes = async (
213247
const builder = new FilterScopeQueryBuilder(dialect, schemas);
214248
const query = builder.build();
215249

216-
let userObj;
217-
userObj = getUser(req);
218-
const dbCredentials = getDatabaseCredentials()[databaseName];
219-
const dbConnection =
220-
await dbConnectionUtil.DBConnectionUtil.getDBConnection({
221-
credentials: dbCredentials,
222-
schemaName: dbCredentials.schemaName,
223-
vocabSchemaName: dbCredentials.vocabSchemaName,
224-
userObj,
225-
});
250+
const dbConnection = await getDatasetFilterDbConnection(databaseName, req);
226251

227252
const queryobject = QueryObject.format(query);
228253
var rangeResults = queryobject.executeQuery(dbConnection);
@@ -261,17 +286,7 @@ const filter = async (databaseName, schemas, dialect, filterParams, req) => {
261286
const builder = new FilterQueryBuilder(dialect, schemas, filterParams);
262287
const query = builder.build();
263288

264-
let userObj;
265-
userObj = getUser(req);
266-
// retrieve the connection
267-
const dbCredentials = getDatabaseCredentials()[databaseName];
268-
const dbConnection =
269-
await dbConnectionUtil.DBConnectionUtil.getDBConnection({
270-
credentials: dbCredentials,
271-
schemaName: dbCredentials.schemaName,
272-
vocabSchemaName: dbCredentials.vocabSchemaName,
273-
userObj,
274-
});
289+
const dbConnection = await getDatasetFilterDbConnection(databaseName, req);
275290

276291
const queryobject = QueryObject.format(query);
277292

plugins/functions/analytics-svc/src/env.ts

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,8 @@ const Env = z
7171
DB_SVC__PATH: z.string().optional(),
7272
DB_SVC__PORT: z.string().optional(),
7373
IS_AUDIT_LOG_ENABLED: z.string().optional(),
74+
IS_CDM_SQL_AUDIT_LOG_ENABLED: z.string().optional(),
75+
AUDIT_LOG_TO_CONSOLE: z.string().optional(),
7476
ANALYTICS_HANA_STREAMING_ENABLED: z.string(),
7577
ANALYTICS_STREAMING_CHUNK_SIZE_BY_DIALECT: z
7678
.string()

plugins/functions/analytics-svc/src/main.ts

Lines changed: 42 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,11 @@ import PortalServerAPI from "./api/PortalServerAPI";
3232
import { env } from "./env";
3333
import addCorrelationIDToHeader from "./middleware/AddCorrelationId.ts";
3434
import { parseValueForPrototypePollutingAssignment } from "./utils/utils";
35+
import { getAuditUserIdFromRequest } from "./utils/AuditLogger.ts";
36+
import {
37+
createCdmSqlAuditConnection,
38+
createCdmSqlAuditContext,
39+
} from "./utils/CdmSqlAuditLogger.ts";
3540
dotenv.config();
3641
const log = console; //Logger.CreateLogger("analytics-log");
3742
const mriConfigConnection = new MriConfigConnection(
@@ -151,22 +156,43 @@ const initRoutes = async (app: express.Application) => {
151156
credentials = req.dbCredentials.studyAnalyticsCredential;
152157
}
153158

154-
if (credentials.dialect === ANALYTICS_DB_DIALECTS.HANA) {
155-
req.dbConnections = await getDBConnections({
156-
analyticsCredentials: credentials,
157-
userObj,
158-
sessionVariables: {
159-
paConfigId: req.paConfigId,
160-
paConfigVersion: req.paConfigVersion,
161-
cdmConfigId: req.cdmConfigId,
162-
cdmConfigVersion: req.cdmConfigVersion,
163-
},
164-
});
165-
} else {
166-
req.dbConnections = getTrexDbConnection({
167-
analyticsCredentials: credentials,
168-
});
169-
}
159+
const dbConnections =
160+
credentials.dialect === ANALYTICS_DB_DIALECTS.HANA
161+
? await getDBConnections({
162+
analyticsCredentials: credentials,
163+
userObj,
164+
sessionVariables: {
165+
paConfigId: req.paConfigId,
166+
paConfigVersion: req.paConfigVersion,
167+
cdmConfigId: req.cdmConfigId,
168+
cdmConfigVersion: req.cdmConfigVersion,
169+
},
170+
})
171+
: getTrexDbConnection({
172+
analyticsCredentials: credentials,
173+
});
174+
175+
req.dbConnections = {
176+
...dbConnections,
177+
analyticsConnection: createCdmSqlAuditConnection(
178+
dbConnections.analyticsConnection,
179+
createCdmSqlAuditContext({
180+
request: req,
181+
actorId:
182+
userObj?.getUser() ??
183+
getAuditUserIdFromRequest(req) ??
184+
"unknown",
185+
databaseCode: credentials.code,
186+
databaseDialect: credentials.dialect,
187+
databaseEngine:
188+
credentials.dialect ===
189+
ANALYTICS_DB_DIALECTS.HANA
190+
? "hana"
191+
: "duckdb",
192+
schemaName: credentials.schema,
193+
})
194+
),
195+
};
170196
}
171197

172198
next();

plugins/functions/analytics-svc/src/mri/endpoint/CohortEndpoint.ts

Lines changed: 20 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,10 @@ import CreateLogger = Logger.CreateLogger;
99
import QueryObject = qo.QueryObject;
1010
import { Connection as connLib } from "@alp/alp-base-utils";
1111
import ConnectionInterface = connLib.ConnectionInterface;
12+
import {
13+
executeWithCdmSqlAudit,
14+
type CdmSqlAuditContext,
15+
} from "../../utils/CdmSqlAuditLogger.ts";
1216
const logger = CreateLogger("analytics-log");
1317

1418
declare const Trex: any;
@@ -562,7 +566,12 @@ export class CohortEndpoint {
562566
cohortDefinitionId: number,
563567
cohort: CohortType,
564568
queryObject: QueryObjectType,
565-
metadata: { datasetId: string; token: string; dbCredential: any }
569+
metadata: {
570+
datasetId: string;
571+
token: string;
572+
dbCredential: any;
573+
auditContext: CdmSqlAuditContext;
574+
}
566575
) {
567576
try {
568577
const partialInsertQuery = QueryObject.formatDict(
@@ -608,9 +617,16 @@ export class CohortEndpoint {
608617
Number(cohortDefinitionId),
609618
JSON.stringify(sessionVars)
610619
);
611-
const result = await materializeQuery.executeQuery<{ processed_rows: number }>(
612-
memConn
613-
);
620+
const result = await executeWithCdmSqlAudit({
621+
context: metadata.auditContext,
622+
operation: "executeQuery",
623+
sql: translatedSql,
624+
parameters: preparedQuery.placeholders ?? [],
625+
execute: () =>
626+
materializeQuery.executeQuery<{
627+
processed_rows: number;
628+
}>(memConn),
629+
});
614630
return result?.data?.[0]?.processed_rows;
615631
} finally {
616632
if (typeof memConn?.close === "function") {
Lines changed: 112 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,112 @@
1+
import { env } from "../env.ts";
2+
3+
export const AUDIT_LOG_DIRECTORY = "/var/log/d2e/audit";
4+
export const PATIENT_ACCESS_AUDIT_FILE = "patient-access.ndjson";
5+
export const CDM_SQL_AUDIT_FILE = "cdm-sql-access.ndjson";
6+
7+
export type AuditEvent = Record<string, unknown>;
8+
9+
export type AuditTransport = {
10+
audit(message: unknown, user: string): void | Promise<void>;
11+
};
12+
13+
export interface AuditEventWriter {
14+
append(fileName: string, event: AuditEvent): Promise<void>;
15+
}
16+
17+
const encoder = new TextEncoder();
18+
19+
function getAuditFilePath(directory: string, fileName: string): string {
20+
if (!/^[a-z0-9][a-z0-9.-]*\.ndjson$/.test(fileName)) {
21+
throw new Error("Invalid audit file name");
22+
}
23+
24+
return `${directory.replace(/\/+$/, "")}/${fileName}`;
25+
}
26+
27+
async function chmodWhenAvailable(path: string, mode: number): Promise<void> {
28+
try {
29+
await Deno.chmod(path, mode);
30+
} catch (error) {
31+
if (
32+
error instanceof Error &&
33+
error.name === "PermissionDenied" &&
34+
error.message === "Deno.chmod is blocklisted"
35+
) {
36+
return;
37+
}
38+
throw error;
39+
}
40+
}
41+
42+
export class NdjsonAuditEventWriter implements AuditEventWriter {
43+
public constructor(private readonly directory = AUDIT_LOG_DIRECTORY) {}
44+
45+
public async append(fileName: string, event: AuditEvent): Promise<void> {
46+
await Deno.mkdir(this.directory, { recursive: true, mode: 0o750 });
47+
await chmodWhenAvailable(this.directory, 0o750);
48+
49+
const path = getAuditFilePath(this.directory, fileName);
50+
const bytes = encoder.encode(`${JSON.stringify(event)}\n`);
51+
const file = await Deno.open(path, {
52+
append: true,
53+
create: true,
54+
mode: 0o640,
55+
write: true,
56+
});
57+
58+
try {
59+
await chmodWhenAvailable(path, 0o640);
60+
// Keep one event in one append operation so concurrent Trex workers
61+
// cannot interleave portions of separate NDJSON records.
62+
const bytesWritten = await file.write(bytes);
63+
if (bytesWritten !== bytes.length) {
64+
throw new Error("Incomplete audit event append");
65+
}
66+
} finally {
67+
file.close();
68+
}
69+
}
70+
}
71+
72+
export class ConsoleAuditEventWriter implements AuditEventWriter {
73+
public append(_fileName: string, event: AuditEvent): Promise<void> {
74+
console.info(JSON.stringify(event));
75+
return Promise.resolve();
76+
}
77+
}
78+
79+
export function createAuditEventWriter(): AuditEventWriter {
80+
return env.AUDIT_LOG_TO_CONSOLE?.toLowerCase() === "true"
81+
? new ConsoleAuditEventWriter()
82+
: new NdjsonAuditEventWriter();
83+
}
84+
85+
export function createPatientAccessAuditTransport(
86+
writer: AuditEventWriter = createAuditEventWriter()
87+
): AuditTransport {
88+
return {
89+
async audit(message: unknown, user: string): Promise<void> {
90+
const eventData =
91+
message && typeof message === "object"
92+
? (message as AuditEvent)
93+
: { message: String(message) };
94+
95+
try {
96+
await writer.append(PATIENT_ACCESS_AUDIT_FILE, {
97+
...eventData,
98+
schemaVersion: 1,
99+
eventType: "patient.access",
100+
actor: {
101+
type: "user",
102+
id: user,
103+
},
104+
});
105+
} catch (_error) {
106+
// Audit persistence is fail-open. Never echo the event payload
107+
// or patient identifier into operational container logs.
108+
console.error("Patient audit event could not be written");
109+
}
110+
},
111+
};
112+
}

0 commit comments

Comments
 (0)