From 8f312d299b096223ea9ed1a9ba1bd8aaf91f347d Mon Sep 17 00:00:00 2001 From: Theodore Li Date: Tue, 23 Jun 2026 02:29:01 -0700 Subject: [PATCH] feat(guardrails): PII redaction via Presidio sidecar (native VIN, per-rule language) (#5174) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit * fix(logs): run PII redaction over HTTP and fix Presidio provisioning - resolve the guardrails venv via candidate paths and fail fast instead of silently falling back to system python3 (the misleading "Presidio not installed" that broke redaction and the guardrails block in deployed runtimes) - install the en_core_web_lg spaCy model in setup.sh and app.Dockerfile - route log redaction through an internal /api/guardrails/mask-batch endpoint so Presidio always runs in the app container, including async executions that persist inside the trigger.dev runtime * fix(guardrails): chunk + time-bound internal PII mask requests - chunk maskPIIBatchViaHttp by count (2000) and bytes (256KB) so large executions split across requests and never hit the contract's 100k cap - add AbortSignal.timeout(45s) per request so a slow/unreachable app container aborts and the caller scrubs, instead of hanging the trigger.dev job - catch maskPIIBatch failures in the route: log and return a structured 500 (broken venv fails loudly server-side; caller still scrubs, no leak) - add mask-client tests (order across chunks, count split, non-2xx, empty) * fix(guardrails): mint internal token per mask request A single token (5min TTL) could expire mid-batch when a large execution fans out into many sequential chunk requests; mint one per request instead. * feat(guardrails): run PII via Presidio sidecars + TS recognizer registry - replace the per-call python3 subprocess (cold spaCy load every call) with two long-lived Presidio sidecars (analyzer + anonymizer) reached over HTTP; the app image no longer carries Python/Presidio/venv - add PRESIDIO_ANALYZER_URL / PRESIDIO_ANONYMIZER_URL - move VIN out of Python into a TS recognizer (check-digit validated) behind a CUSTOM_RECOGNIZERS registry so new custom detectors are one entry; masking is handled uniformly by the anonymizer - drive the guardrails block's PII type picker from the shared pii-entities catalog (adds VIN, fixes drift) so block + Data Retention never diverge - delete validate_pii.py, requirements.txt, setup.sh and the Dockerfile venv step * fix(guardrails): bound-parallelize mask batch; refresh stale comments - maskPIIBatch runs per-string sidecar calls with bounded concurrency (8) via mapWithConcurrency, so a chunk of many small leaves finishes within the 45s request timeout instead of aborting and scrubbing; order + fail-on-error kept - drop stale comments referencing the deleted Python venv / 30s subprocess timeout * refactor(guardrails): single Presidio image, native VIN, per-rule redaction language - collapse the analyzer/anonymizer URLs into one PRESIDIO_URL (combined image serves /analyze + /anonymize) - remove the TS VIN recognizer (vin.ts, recognizers.ts) — VIN is now native + multi-language in the image; validate_pii is a thin analyze→anonymize client - trim KR_RRN/TH_TNIN from the catalog (no Korean/Thai model in the image) - add per-rule redaction language: PII_LANGUAGES catalog drives the contract enum, the Data Retention rule modal, and the guardrails block dropdown; resolver + logger thread it through to maskPIIBatch (default en), so non-English entity rules (e.g. ES_NIF) actually fire instead of silently no-op'ing under en * fix(guardrails): correct sidecar port (5001) + README for combined image The combined Presidio image (docker/pii.Dockerfile) serves /analyze + /anonymize on a single port 5001 with native VIN + multi-language recognizers. Fix the PRESIDIO_URL default (was 5002) and rewrite the README, which still described two stock containers and a TS VIN recognizer. * fix(guardrails): coerce stored redaction language in the resolver The persist-path resolver accepted any stored language string, so a stale/invalid code (e.g. a dropped locale) would reach Presidio and scrub the log even though the admin UI shows English. Coerce against the supported set via a shared coercePiiLanguage helper (now reused by the data-retention route too), falling back to en for unknown values. * fix(guardrails): rename PRESIDIO_URL env var to PII_URL Match the infra taskdef, which sets PII_URL on the app container for the combined Presidio sidecar. --- .../api/guardrails/mask-batch/route.test.ts | 64 +++ .../app/api/guardrails/mask-batch/route.ts | 45 +++ .../[id]/data-retention/route.ts | 10 +- apps/sim/blocks/blocks/guardrails.ts | 77 +--- .../components/data-retention-settings.tsx | 38 +- apps/sim/lib/api/contracts/hotspots.ts | 28 ++ apps/sim/lib/api/contracts/primitives.ts | 3 + apps/sim/lib/billing/retention.test.ts | 35 +- apps/sim/lib/billing/retention.ts | 7 +- apps/sim/lib/core/config/env.ts | 1 + apps/sim/lib/guardrails/.gitignore | 13 - apps/sim/lib/guardrails/README.md | 35 +- apps/sim/lib/guardrails/mask-client.test.ts | 68 ++++ apps/sim/lib/guardrails/mask-client.ts | 99 +++++ apps/sim/lib/guardrails/pii-entities.ts | 38 +- apps/sim/lib/guardrails/requirements.txt | 4 - apps/sim/lib/guardrails/setup.sh | 37 -- apps/sim/lib/guardrails/validate_pii.py | 260 ------------ apps/sim/lib/guardrails/validate_pii.test.ts | 118 ++++++ apps/sim/lib/guardrails/validate_pii.ts | 370 ++++++------------ apps/sim/lib/logs/execution/logger.ts | 5 +- .../lib/logs/execution/pii-redaction.test.ts | 4 +- apps/sim/lib/logs/execution/pii-redaction.ts | 9 +- docker/app.Dockerfile | 12 +- packages/db/schema.ts | 2 + scripts/check-api-validation-contracts.ts | 4 +- 26 files changed, 702 insertions(+), 684 deletions(-) create mode 100644 apps/sim/app/api/guardrails/mask-batch/route.test.ts create mode 100644 apps/sim/app/api/guardrails/mask-batch/route.ts delete mode 100644 apps/sim/lib/guardrails/.gitignore create mode 100644 apps/sim/lib/guardrails/mask-client.test.ts create mode 100644 apps/sim/lib/guardrails/mask-client.ts delete mode 100644 apps/sim/lib/guardrails/requirements.txt delete mode 100755 apps/sim/lib/guardrails/setup.sh delete mode 100644 apps/sim/lib/guardrails/validate_pii.py create mode 100644 apps/sim/lib/guardrails/validate_pii.test.ts diff --git a/apps/sim/app/api/guardrails/mask-batch/route.test.ts b/apps/sim/app/api/guardrails/mask-batch/route.test.ts new file mode 100644 index 0000000000..cbb5b12265 --- /dev/null +++ b/apps/sim/app/api/guardrails/mask-batch/route.test.ts @@ -0,0 +1,64 @@ +/** + * @vitest-environment node + */ +import { createMockRequest } from '@sim/testing' +import { beforeEach, describe, expect, it, vi } from 'vitest' + +const { mockCheckInternalAuth, mockMaskPIIBatch } = vi.hoisted(() => ({ + mockCheckInternalAuth: vi.fn(), + mockMaskPIIBatch: vi.fn(), +})) + +vi.mock('@/lib/auth/hybrid', () => ({ + checkInternalAuth: mockCheckInternalAuth, +})) + +vi.mock('@/lib/guardrails/validate_pii', () => ({ + maskPIIBatch: mockMaskPIIBatch, +})) + +import { POST } from '@/app/api/guardrails/mask-batch/route' + +describe('POST /api/guardrails/mask-batch', () => { + beforeEach(() => { + vi.clearAllMocks() + mockCheckInternalAuth.mockResolvedValue({ success: true }) + mockMaskPIIBatch.mockImplementation(async (texts: string[]) => texts.map((t) => `M(${t})`)) + }) + + it('returns 401 without internal auth', async () => { + mockCheckInternalAuth.mockResolvedValue({ + success: false, + error: 'Internal authentication required', + }) + + const res = await POST( + createMockRequest('POST', { texts: ['a@b.com'], entityTypes: ['EMAIL_ADDRESS'] }) + ) + + expect(res.status).toBe(401) + expect(mockMaskPIIBatch).not.toHaveBeenCalled() + }) + + it('masks the batch in-process and preserves order', async () => { + const res = await POST( + createMockRequest('POST', { + texts: ['a@b.com', 'hello'], + entityTypes: ['EMAIL_ADDRESS'], + language: 'en', + }) + ) + + expect(res.status).toBe(200) + const json = await res.json() + expect(json.masked).toEqual(['M(a@b.com)', 'M(hello)']) + expect(mockMaskPIIBatch).toHaveBeenCalledWith(['a@b.com', 'hello'], ['EMAIL_ADDRESS'], 'en') + }) + + it('rejects an invalid body with 400', async () => { + const res = await POST(createMockRequest('POST', { texts: 'not-an-array', entityTypes: [] })) + + expect(res.status).toBe(400) + expect(mockMaskPIIBatch).not.toHaveBeenCalled() + }) +}) diff --git a/apps/sim/app/api/guardrails/mask-batch/route.ts b/apps/sim/app/api/guardrails/mask-batch/route.ts new file mode 100644 index 0000000000..696b69e749 --- /dev/null +++ b/apps/sim/app/api/guardrails/mask-batch/route.ts @@ -0,0 +1,45 @@ +import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' +import { type NextRequest, NextResponse } from 'next/server' +import { guardrailsMaskBatchContract } from '@/lib/api/contracts' +import { parseRequest } from '@/lib/api/server' +import { checkInternalAuth } from '@/lib/auth/hybrid' +import { withRouteHandler } from '@/lib/core/utils/with-route-handler' +import { maskPIIBatch } from '@/lib/guardrails/validate_pii' + +const logger = createLogger('GuardrailsMaskBatchAPI') + +/** + * Internal batch PII masking. The log-redaction persist path runs in both the + * Next.js server and the trigger.dev runtime, but the Presidio sidecars live only + * in the app task — so redaction calls this endpoint server-to-server (internal + * JWT) to keep Presidio centralized here. + */ +export const POST = withRouteHandler(async (request: NextRequest) => { + const auth = await checkInternalAuth(request, { requireWorkflowId: false }) + if (!auth.success) { + return NextResponse.json({ error: 'Unauthorized' }, { status: 401 }) + } + + const parsed = await parseRequest(guardrailsMaskBatchContract, request, {}) + if (!parsed.success) return parsed.response + + const { texts, entityTypes, language } = parsed.data.body + + try { + const masked = await maskPIIBatch(texts, entityTypes, language) + logger.info('Masked PII batch', { count: texts.length }) + return NextResponse.json({ masked }) + } catch (error) { + // An unreachable/misconfigured Presidio sidecar makes maskPIIBatch throw; fail + // loudly here (the caller scrubs to REDACTION_FAILED, so PII is never leaked). + logger.error('PII batch masking failed', { + error: getErrorMessage(error), + count: texts.length, + }) + return NextResponse.json( + { error: getErrorMessage(error, 'PII masking failed') }, + { status: 500 } + ) + } +}) diff --git a/apps/sim/app/api/organizations/[id]/data-retention/route.ts b/apps/sim/app/api/organizations/[id]/data-retention/route.ts index 37fbbaabb9..7d7052a392 100644 --- a/apps/sim/app/api/organizations/[id]/data-retention/route.ts +++ b/apps/sim/app/api/organizations/[id]/data-retention/route.ts @@ -16,6 +16,7 @@ import { isOrganizationOnEnterprisePlan } from '@/lib/billing/core/subscription' import { isBillingEnabled } from '@/lib/core/config/env-flags' import { isFeatureEnabled } from '@/lib/core/config/feature-flags' import { withRouteHandler } from '@/lib/core/utils/with-route-handler' +import { coercePiiLanguage } from '@/lib/guardrails/pii-entities' const logger = createLogger('DataRetentionAPI') @@ -35,7 +36,14 @@ function normalizeConfigured( logRetentionHours: settings?.logRetentionHours ?? null, softDeleteRetentionHours: settings?.softDeleteRetentionHours ?? null, taskCleanupHours: settings?.taskCleanupHours ?? null, - piiRedaction: settings?.piiRedaction?.rules ? { rules: settings.piiRedaction.rules } : null, + piiRedaction: settings?.piiRedaction?.rules + ? { + rules: settings.piiRedaction.rules.map((rule) => ({ + ...rule, + language: coercePiiLanguage(rule.language), + })), + } + : null, } } diff --git a/apps/sim/blocks/blocks/guardrails.ts b/apps/sim/blocks/blocks/guardrails.ts index 42fefcda81..7acd5a8901 100644 --- a/apps/sim/blocks/blocks/guardrails.ts +++ b/apps/sim/blocks/blocks/guardrails.ts @@ -1,4 +1,5 @@ import { ShieldCheckIcon } from '@/components/icons' +import { PII_ENTITY_GROUPS, PII_LANGUAGES } from '@/lib/guardrails/pii-entities' import type { BlockConfig } from '@/blocks/types' import { getModelOptions, @@ -170,65 +171,15 @@ Return ONLY the regex pattern - no explanations, no quotes, no forward slashes, title: 'PII Types to Detect', type: 'grouped-checkbox-list', maxHeight: 400, - options: [ - // Common PII types - { label: 'Person name', id: 'PERSON', group: 'Common' }, - { label: 'Email address', id: 'EMAIL_ADDRESS', group: 'Common' }, - { label: 'Phone number', id: 'PHONE_NUMBER', group: 'Common' }, - { label: 'Location', id: 'LOCATION', group: 'Common' }, - { label: 'Date or time', id: 'DATE_TIME', group: 'Common' }, - { label: 'IP address', id: 'IP_ADDRESS', group: 'Common' }, - { label: 'URL', id: 'URL', group: 'Common' }, - { label: 'Credit card number', id: 'CREDIT_CARD', group: 'Common' }, - { label: 'International bank account number (IBAN)', id: 'IBAN_CODE', group: 'Common' }, - { label: 'Cryptocurrency wallet address', id: 'CRYPTO', group: 'Common' }, - { label: 'Medical license number', id: 'MEDICAL_LICENSE', group: 'Common' }, - { label: 'Nationality / religion / political group', id: 'NRP', group: 'Common' }, - - // USA - { label: 'US bank account number', id: 'US_BANK_NUMBER', group: 'USA' }, - { label: 'US driver license number', id: 'US_DRIVER_LICENSE', group: 'USA' }, - { - label: 'US individual taxpayer identification number (ITIN)', - id: 'US_ITIN', - group: 'USA', - }, - { label: 'US passport number', id: 'US_PASSPORT', group: 'USA' }, - { label: 'US Social Security number', id: 'US_SSN', group: 'USA' }, - - // UK - { label: 'UK National Insurance number', id: 'UK_NINO', group: 'UK' }, - { label: 'UK NHS number', id: 'UK_NHS', group: 'UK' }, - - // Spain - { label: 'Spanish NIF number', id: 'ES_NIF', group: 'Spain' }, - { label: 'Spanish NIE number', id: 'ES_NIE', group: 'Spain' }, - - // Italy - { label: 'Italian fiscal code', id: 'IT_FISCAL_CODE', group: 'Italy' }, - { label: 'Italian driver license', id: 'IT_DRIVER_LICENSE', group: 'Italy' }, - { label: 'Italian identity card', id: 'IT_IDENTITY_CARD', group: 'Italy' }, - { label: 'Italian passport', id: 'IT_PASSPORT', group: 'Italy' }, - - // Poland - { label: 'Polish PESEL', id: 'PL_PESEL', group: 'Poland' }, - - // Singapore - { label: 'Singapore NRIC/FIN', id: 'SG_NRIC_FIN', group: 'Singapore' }, - - // Australia - { label: 'Australian business number (ABN)', id: 'AU_ABN', group: 'Australia' }, - { label: 'Australian company number (ACN)', id: 'AU_ACN', group: 'Australia' }, - { label: 'Australian tax file number (TFN)', id: 'AU_TFN', group: 'Australia' }, - { label: 'Australian Medicare number', id: 'AU_MEDICARE', group: 'Australia' }, - - // India - { label: 'Indian Aadhaar', id: 'IN_AADHAAR', group: 'India' }, - { label: 'Indian PAN', id: 'IN_PAN', group: 'India' }, - { label: 'Indian vehicle registration', id: 'IN_VEHICLE_REGISTRATION', group: 'India' }, - { label: 'Indian voter number', id: 'IN_VOTER', group: 'India' }, - { label: 'Indian passport', id: 'IN_PASSPORT', group: 'India' }, - ], + // Driven by the shared catalog (includes VIN and custom recognizers) so the + // block and the Data Retention settings never drift. + options: PII_ENTITY_GROUPS.flatMap((group) => + group.entities.map((entity) => ({ + label: entity.label, + id: entity.value, + group: group.label, + })) + ), condition: { field: 'validationType', value: ['pii'], @@ -255,13 +206,7 @@ Return ONLY the regex pattern - no explanations, no quotes, no forward slashes, id: 'piiLanguage', title: 'Language', type: 'dropdown', - options: [ - { label: 'English', id: 'en' }, - { label: 'Spanish', id: 'es' }, - { label: 'Italian', id: 'it' }, - { label: 'Polish', id: 'pl' }, - { label: 'Finnish', id: 'fi' }, - ], + options: PII_LANGUAGES.map((language) => ({ label: language.label, id: language.value })), defaultValue: 'en', condition: { field: 'validationType', diff --git a/apps/sim/ee/data-retention/components/data-retention-settings.tsx b/apps/sim/ee/data-retention/components/data-retention-settings.tsx index bca54112d5..1c39594d1e 100644 --- a/apps/sim/ee/data-retention/components/data-retention-settings.tsx +++ b/apps/sim/ee/data-retention/components/data-retention-settings.tsx @@ -21,7 +21,13 @@ import { } from '@/components/emcn' import { useSession } from '@/lib/auth/auth-client' import { isBillingEnabled } from '@/lib/core/config/env-flags' -import { PII_ENTITY_GROUPS, SUPPORTED_PII_ENTITIES } from '@/lib/guardrails/pii-entities' +import { + DEFAULT_PII_LANGUAGE, + PII_ENTITY_GROUPS, + PII_LANGUAGES, + type PIILanguage, + SUPPORTED_PII_ENTITIES, +} from '@/lib/guardrails/pii-entities' import { getUserRole } from '@/lib/workspaces/organization/utils' import { SettingsSection } from '@/app/workspace/[workspaceId]/settings/components/settings-section/settings-section' import { InfoNote } from '@/ee/components/info-note' @@ -59,6 +65,7 @@ interface RuleDraft { id: string entityTypes: string[] workspaceId: string | null + language: PIILanguage } function hoursToDisplayDays(hours: number | null): string { @@ -75,6 +82,7 @@ function normalizeRule(rule: RuleDraft): string { return JSON.stringify({ entityTypes: [...rule.entityTypes].sort(), workspaceId: rule.workspaceId, + language: rule.language, }) } @@ -227,6 +235,18 @@ function RuleModal({ onChange={(entityTypes) => onChange({ ...draft, entityTypes })} /> + + onChange({ ...draft, language: language as PIILanguage })} + options={PII_LANGUAGES.map((l) => ({ value: l.value, label: l.label }))} + align='start' + /> + +export type GuardrailsMaskBatchResult = z.output + const chatMessageSchema = z.object({ role: z.enum(['user', 'assistant', 'system']), content: z.string(), diff --git a/apps/sim/lib/api/contracts/primitives.ts b/apps/sim/lib/api/contracts/primitives.ts index e3e6148402..2b0d598d1a 100644 --- a/apps/sim/lib/api/contracts/primitives.ts +++ b/apps/sim/lib/api/contracts/primitives.ts @@ -1,4 +1,5 @@ import { z } from 'zod' +import { PII_LANGUAGE_CODES } from '@/lib/guardrails/pii-entities' export const unknownRecordSchema = z.record(z.string(), z.unknown()) @@ -93,6 +94,8 @@ export const piiRedactionRuleSchema = z.object({ entityTypes: z.array(z.string().min(1, 'Entity type cannot be empty')).max(100), /** null = all workspaces; otherwise the single targeted workspace. */ workspaceId: z.string().min(1).nullable(), + /** Language whose Presidio recognizers apply; defaults to English. */ + language: z.enum(PII_LANGUAGE_CODES).optional(), }) export type PiiRedactionRule = z.output diff --git a/apps/sim/lib/billing/retention.test.ts b/apps/sim/lib/billing/retention.test.ts index 2852cb6a64..15714bc046 100644 --- a/apps/sim/lib/billing/retention.test.ts +++ b/apps/sim/lib/billing/retention.test.ts @@ -21,7 +21,11 @@ describe('resolveEffectivePiiRedaction', () => { orgSettings: settings([allRule]), workspaceId: 'ws-1', }) - expect(result).toEqual({ enabled: true, entityTypes: ['EMAIL_ADDRESS', 'PHONE_NUMBER'] }) + expect(result).toEqual({ + enabled: true, + entityTypes: ['EMAIL_ADDRESS', 'PHONE_NUMBER'], + language: 'en', + }) }) it('lets a workspace-specific rule override the all rule', () => { @@ -29,7 +33,27 @@ describe('resolveEffectivePiiRedaction', () => { orgSettings: settings([allRule, { id: 'r-1', entityTypes: ['US_SSN'], workspaceId: 'ws-1' }]), workspaceId: 'ws-1', }) - expect(result).toEqual({ enabled: true, entityTypes: ['US_SSN'] }) + expect(result).toEqual({ enabled: true, entityTypes: ['US_SSN'], language: 'en' }) + }) + + it('carries the rule language through (defaults to en)', () => { + const result = resolveEffectivePiiRedaction({ + orgSettings: settings([ + { id: 'r-es', entityTypes: ['ES_NIF'], workspaceId: 'ws-1', language: 'es' }, + ]), + workspaceId: 'ws-1', + }) + expect(result).toEqual({ enabled: true, entityTypes: ['ES_NIF'], language: 'es' }) + }) + + it('falls back to en when a stored language is unsupported/stale', () => { + const result = resolveEffectivePiiRedaction({ + orgSettings: settings([ + { id: 'r-de', entityTypes: ['EMAIL_ADDRESS'], workspaceId: 'ws-1', language: 'de' }, + ]), + workspaceId: 'ws-1', + }) + expect(result).toEqual({ enabled: true, entityTypes: ['EMAIL_ADDRESS'], language: 'en' }) }) it('exempts a workspace when its specific rule has no entity types', () => { @@ -37,7 +61,7 @@ describe('resolveEffectivePiiRedaction', () => { orgSettings: settings([allRule, { id: 'r-1', entityTypes: [], workspaceId: 'ws-1' }]), workspaceId: 'ws-1', }) - expect(result).toEqual({ enabled: false, entityTypes: [] }) + expect(result).toEqual({ enabled: false, entityTypes: [], language: 'en' }) }) it('is disabled when no rule matches and there is no all rule', () => { @@ -45,16 +69,17 @@ describe('resolveEffectivePiiRedaction', () => { orgSettings: settings([{ id: 'r-1', entityTypes: ['US_SSN'], workspaceId: 'ws-2' }]), workspaceId: 'ws-1', }) - expect(result).toEqual({ enabled: false, entityTypes: [] }) + expect(result).toEqual({ enabled: false, entityTypes: [], language: 'en' }) }) it('is disabled when there are no rules', () => { expect( resolveEffectivePiiRedaction({ orgSettings: settings([]), workspaceId: 'ws-1' }) - ).toEqual({ enabled: false, entityTypes: [] }) + ).toEqual({ enabled: false, entityTypes: [], language: 'en' }) expect(resolveEffectivePiiRedaction({ orgSettings: null, workspaceId: 'ws-1' })).toEqual({ enabled: false, entityTypes: [], + language: 'en', }) }) }) diff --git a/apps/sim/lib/billing/retention.ts b/apps/sim/lib/billing/retention.ts index 183dbb280e..dafb9e3a78 100644 --- a/apps/sim/lib/billing/retention.ts +++ b/apps/sim/lib/billing/retention.ts @@ -1,14 +1,18 @@ import type { DataRetentionSettings } from '@sim/db/schema' +import { coercePiiLanguage, DEFAULT_PII_LANGUAGE } from '@/lib/guardrails/pii-entities' export interface EffectivePiiRedaction { enabled: boolean /** Presidio entity types to mask. Empty = redact all detected PII. */ entityTypes: string[] + /** Language whose Presidio recognizers apply when masking. */ + language: string } export const DEFAULT_PII_REDACTION: EffectivePiiRedaction = { enabled: false, entityTypes: [], + language: DEFAULT_PII_LANGUAGE, } /** @@ -34,5 +38,6 @@ export function resolveEffectivePiiRedaction(params: { ? rule.entityTypes.filter((t): t is string => typeof t === 'string') : [] if (types.length === 0) return DEFAULT_PII_REDACTION - return { enabled: true, entityTypes: types } + const language = coercePiiLanguage(rule?.language) ?? DEFAULT_PII_LANGUAGE + return { enabled: true, entityTypes: types, language } } diff --git a/apps/sim/lib/core/config/env.ts b/apps/sim/lib/core/config/env.ts index 5ba9df0775..89924eb568 100644 --- a/apps/sim/lib/core/config/env.ts +++ b/apps/sim/lib/core/config/env.ts @@ -312,6 +312,7 @@ export const env = createEnv({ PORT: z.number().optional(), // Main application port INTERNAL_API_BASE_URL: z.string().optional(), // Optional internal base URL for server-side self-calls; must include protocol if set (e.g., http://sim-app.namespace.svc.cluster.local:3000) ALLOWED_ORIGINS: z.string().optional(), // CORS allowed origins + PII_URL: z.string().optional(), // Presidio PII sidecar base URL serving /analyze + /anonymize (default http://localhost:5001) // OAuth Integration Credentials - All optional, enables third-party integrations GOOGLE_CLIENT_ID: z.string().optional(), // Google OAuth client ID for Google services diff --git a/apps/sim/lib/guardrails/.gitignore b/apps/sim/lib/guardrails/.gitignore deleted file mode 100644 index 3485e9bdf6..0000000000 --- a/apps/sim/lib/guardrails/.gitignore +++ /dev/null @@ -1,13 +0,0 @@ -# Python virtual environment -venv/ - -# Python cache -__pycache__/ -*.pyc -*.pyo -*.pyd -.Python - -# Presidio cache -.presidio/ - diff --git a/apps/sim/lib/guardrails/README.md b/apps/sim/lib/guardrails/README.md index 6ce7802d22..6c0a5df970 100644 --- a/apps/sim/lib/guardrails/README.md +++ b/apps/sim/lib/guardrails/README.md @@ -19,22 +19,29 @@ For **hallucination detection**, you'll need: - A knowledge base with documents - An LLM provider API key (or use hosted models) -### Python Validators (PII Detection) +### PII Detection (Presidio sidecar) -For **PII detection**, you need to set up a Python virtual environment and install Microsoft Presidio: +PII detection runs against **one** long-lived Presidio sidecar — a combined service (built from +`docker/pii.Dockerfile`, source in `apps/pii/server.py`) that constructs a warm `AnalyzerEngine` + +`AnonymizerEngine` once and exposes both `/analyze` and `/anonymize` (plus `/health`) on a single +port. In deployment it runs alongside the app container in the same ECS task; locally, build and run +it: ```bash -cd apps/sim/lib/guardrails -./setup.sh +docker build -f docker/pii.Dockerfile -t sim-pii . +docker run -d -p 5001:5001 sim-pii ``` -This will: -1. Create a Python virtual environment in `apps/sim/lib/guardrails/venv` -2. Install required dependencies: - - `presidio-analyzer` - PII detection engine - - `presidio-anonymizer` - PII masking/anonymization +Point the app at it (default shown): -The TypeScript wrapper will automatically use the virtual environment's Python interpreter. +```bash +PII_URL=http://localhost:5001 +``` + +The image bakes in the recognizers itself — a check-digit-validated **VIN** recognizer and +multi-language NLP models (en/es/it/pl/fi) — so the app is a thin HTTP client (`validate_pii.ts`) with +no Python or local venv. The redaction language is configured per rule (Data Retention) and defaults +to English. ## Usage @@ -93,10 +100,8 @@ See [Presidio documentation](https://microsoft.github.io/presidio/supported_enti - `validate_json.ts` - JSON validation (TypeScript) - `validate_regex.ts` - Regex validation (TypeScript) - `validate_hallucination.ts` - Hallucination detection with RAG + LLM scoring (TypeScript) -- `validate_pii.ts` - PII detection TypeScript wrapper (TypeScript) -- `validate_pii.py` - PII detection using Microsoft Presidio (Python) +- `validate_pii.ts` - PII detection client: calls the Presidio sidecar's /analyze + /anonymize (TypeScript) +- `pii-entities.ts` - Client-safe PII entity + language catalog (shared by the block and Data Retention) +- `mask-client.ts` - Internal HTTP client for batch PII masking from the log-redaction persist path - `validate.test.ts` - Test suite for JSON and regex validators -- `validate_hallucination.py` - Legacy Python hallucination detector (deprecated) -- `requirements.txt` - Python dependencies for PII detection (and legacy hallucination) -- `setup.sh` - Legacy installation script (deprecated) diff --git a/apps/sim/lib/guardrails/mask-client.test.ts b/apps/sim/lib/guardrails/mask-client.test.ts new file mode 100644 index 0000000000..d1c4ad5b84 --- /dev/null +++ b/apps/sim/lib/guardrails/mask-client.test.ts @@ -0,0 +1,68 @@ +/** + * @vitest-environment node + */ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +const { mockToken, mockBaseUrl } = vi.hoisted(() => ({ + mockToken: vi.fn(), + mockBaseUrl: vi.fn(), +})) + +vi.mock('@/lib/auth/internal', () => ({ generateInternalToken: mockToken })) +vi.mock('@/lib/core/utils/urls', () => ({ getInternalApiBaseUrl: mockBaseUrl })) + +import { maskPIIBatchViaHttp } from '@/lib/guardrails/mask-client' + +describe('maskPIIBatchViaHttp', () => { + let fetchMock: ReturnType + + beforeEach(() => { + vi.clearAllMocks() + mockToken.mockResolvedValue('tok') + mockBaseUrl.mockReturnValue('http://app.internal:3000') + fetchMock = vi.fn(async (_url: string, init: { body: string }) => { + const { texts } = JSON.parse(init.body) as { texts: string[] } + return new Response(JSON.stringify({ masked: texts.map((t) => `M(${t})`) }), { + status: 200, + headers: { 'content-type': 'application/json' }, + }) + }) + vi.stubGlobal('fetch', fetchMock) + }) + + afterEach(() => { + vi.unstubAllGlobals() + }) + + it('masks a small batch in a single request, with an abort timeout', async () => { + const out = await maskPIIBatchViaHttp(['a', 'b', 'c'], ['EMAIL_ADDRESS']) + + expect(out).toEqual(['M(a)', 'M(b)', 'M(c)']) + expect(fetchMock).toHaveBeenCalledTimes(1) + expect(fetchMock.mock.calls[0][1].signal).toBeInstanceOf(AbortSignal) + }) + + it('splits by count into multiple requests, preserving global order', async () => { + const texts = Array.from({ length: 5000 }, (_, i) => `t${i}`) + + const out = await maskPIIBatchViaHttp(texts, []) + + expect(out).toHaveLength(5000) + expect(out[0]).toBe('M(t0)') + expect(out[4999]).toBe('M(t4999)') + expect(fetchMock).toHaveBeenCalledTimes(3) // 2000-per-request cap + }) + + it('throws on a non-2xx response so the caller can scrub', async () => { + fetchMock.mockResolvedValueOnce(new Response('boom', { status: 500 })) + + await expect(maskPIIBatchViaHttp(['a'], [])).rejects.toThrow(/mask-batch request failed/) + }) + + it('returns [] without any request for empty input', async () => { + const out = await maskPIIBatchViaHttp([], []) + + expect(out).toEqual([]) + expect(fetchMock).not.toHaveBeenCalled() + }) +}) diff --git a/apps/sim/lib/guardrails/mask-client.ts b/apps/sim/lib/guardrails/mask-client.ts new file mode 100644 index 0000000000..3fb818a3c7 --- /dev/null +++ b/apps/sim/lib/guardrails/mask-client.ts @@ -0,0 +1,99 @@ +import type { GuardrailsMaskBatchResult } from '@/lib/api/contracts' +import { generateInternalToken } from '@/lib/auth/internal' +import { getInternalApiBaseUrl } from '@/lib/core/utils/urls' + +/** + * Per-request limits. A chunk is flushed when it hits either bound, keeping each + * request small enough for one short Presidio pass under a tight timeout and far + * below the contract's 100k-entry cap — so large executions split across + * requests instead of failing validation. + */ +const REQUEST_MAX_BYTES = 256 * 1024 +const REQUEST_MAX_COUNT = 2_000 +/** Bounds one mask-batch request; an unreachable/stuck Presidio sidecar aborts so the caller scrubs. */ +const REQUEST_TIMEOUT_MS = 45_000 + +/** + * Mask PII across many strings via the internal app-container endpoint. + * + * The Presidio sidecars run only in the app task, but the log-redaction persist + * path also runs inside the trigger.dev runtime — so redaction always routes + * through HTTP, the same way the guardrails tool does. + * Strings are grouped into byte/count-budgeted chunks; order is preserved, so + * the returned array matches `texts` length. + * + * Rejects on any non-2xx, timeout, or shape mismatch so the caller can apply + * its own fail-safe (scrubbing rather than leaking). + */ +export async function maskPIIBatchViaHttp( + texts: string[], + entityTypes: string[], + language?: string +): Promise { + if (texts.length === 0) return [] + + const url = `${getInternalApiBaseUrl()}/api/guardrails/mask-batch` + + const masked: string[] = [] + let batch: string[] = [] + let batchBytes = 0 + + const flush = async () => { + if (batch.length === 0) return + const out = await postChunk(url, batch, entityTypes, language) + if (out.length !== batch.length) { + throw new Error('PII mask-batch returned an unexpected result') + } + for (const item of out) masked.push(item) + batch = [] + batchBytes = 0 + } + + for (const text of texts) { + const bytes = Buffer.byteLength(text, 'utf8') + if ( + batch.length > 0 && + (batch.length >= REQUEST_MAX_COUNT || batchBytes + bytes > REQUEST_MAX_BYTES) + ) { + await flush() + } + batch.push(text) + batchBytes += bytes + } + await flush() + + return masked +} + +async function postChunk( + url: string, + texts: string[], + entityTypes: string[], + language: string | undefined +): Promise { + // Mint per request: a single token (5min TTL) can expire mid-batch when a + // large execution fans out into many sequential chunk requests. + const token = await generateInternalToken() + + // boundary-raw-fetch: internal server-to-server call to the app container (internal JWT auth, configurable base URL) + const response = await fetch(url, { + method: 'POST', + headers: { + 'content-type': 'application/json', + authorization: `Bearer ${token}`, + }, + body: JSON.stringify({ texts, entityTypes, language }), + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), + }) + + if (!response.ok) { + const detail = await response.text().catch(() => '') + throw new Error(`PII mask-batch request failed (${response.status}): ${detail.slice(0, 200)}`) + } + + const data = (await response.json()) as GuardrailsMaskBatchResult + if (!Array.isArray(data.masked)) { + throw new Error('PII mask-batch returned an unexpected result') + } + return data.masked +} diff --git a/apps/sim/lib/guardrails/pii-entities.ts b/apps/sim/lib/guardrails/pii-entities.ts index 0e67fe22ff..c26e7dc0b9 100644 --- a/apps/sim/lib/guardrails/pii-entities.ts +++ b/apps/sim/lib/guardrails/pii-entities.ts @@ -51,8 +51,6 @@ export const SUPPORTED_PII_ENTITIES = { IN_VOTER: 'Indian voter ID', IN_PASSPORT: 'Indian passport', FI_PERSONAL_IDENTITY_CODE: 'Finnish Personal Identity Code', - KR_RRN: 'Korean Resident Registration Number', - TH_TNIN: 'Thai National ID Number', } as const export type PIIEntityType = keyof typeof SUPPORTED_PII_ENTITIES @@ -115,8 +113,6 @@ export const PII_ENTITY_GROUPS: ReadonlyArray<{ 'IN_VOTER', 'IN_PASSPORT', 'FI_PERSONAL_IDENTITY_CODE', - 'KR_RRN', - 'TH_TNIN', ], }, ].map((group) => ({ @@ -126,3 +122,37 @@ export const PII_ENTITY_GROUPS: ReadonlyArray<{ label: SUPPORTED_PII_ENTITIES[value as PIIEntityType], })), })) + +/** + * Languages the Presidio image has NLP models for. The analyzer only recognizes a + * language's entities when its model is loaded, so this set must match the image. + */ +export const PII_LANGUAGES = [ + { value: 'en', label: 'English' }, + { value: 'es', label: 'Spanish' }, + { value: 'it', label: 'Italian' }, + { value: 'pl', label: 'Polish' }, + { value: 'fi', label: 'Finnish' }, +] as const + +export type PIILanguage = (typeof PII_LANGUAGES)[number]['value'] + +/** Non-empty tuple of language codes for schema/enum use. */ +export const PII_LANGUAGE_CODES = PII_LANGUAGES.map((l) => l.value) as [ + PIILanguage, + ...PIILanguage[], +] + +/** Default redaction language when a rule doesn't set one. */ +export const DEFAULT_PII_LANGUAGE: PIILanguage = 'en' + +/** + * Narrow a loosely-typed (stored/legacy) language to a supported code. Unknown or + * stale values (e.g. a dropped locale) return `undefined` so callers fall back to + * the default rather than forwarding an unsupported language to Presidio. + */ +export function coercePiiLanguage(value: string | undefined): PIILanguage | undefined { + return value && (PII_LANGUAGE_CODES as readonly string[]).includes(value) + ? (value as PIILanguage) + : undefined +} diff --git a/apps/sim/lib/guardrails/requirements.txt b/apps/sim/lib/guardrails/requirements.txt deleted file mode 100644 index 135efae05b..0000000000 --- a/apps/sim/lib/guardrails/requirements.txt +++ /dev/null @@ -1,4 +0,0 @@ -# Microsoft Presidio for PII detection -presidio-analyzer>=2.2.0 -presidio-anonymizer>=2.2.0 - diff --git a/apps/sim/lib/guardrails/setup.sh b/apps/sim/lib/guardrails/setup.sh deleted file mode 100755 index 233e9a51a2..0000000000 --- a/apps/sim/lib/guardrails/setup.sh +++ /dev/null @@ -1,37 +0,0 @@ -#!/bin/bash - -# Setup script for guardrails validators -# This creates a virtual environment and installs Python dependencies - -set -e - -SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" -VENV_DIR="$SCRIPT_DIR/venv" - -echo "Setting up Python environment for guardrails..." - -# Check if Python 3 is available -if ! command -v python3 &> /dev/null; then - echo "Error: python3 is not installed. Please install Python 3 first." - exit 1 -fi - -# Create virtual environment if it doesn't exist -if [ ! -d "$VENV_DIR" ]; then - echo "Creating virtual environment..." - python3 -m venv "$VENV_DIR" -else - echo "Virtual environment already exists." -fi - -# Activate virtual environment and install dependencies -echo "Installing Python dependencies..." -source "$VENV_DIR/bin/activate" -pip install --upgrade pip -pip install -r "$SCRIPT_DIR/requirements.txt" - -echo "" -echo "✅ Setup complete! Guardrails validators are ready to use." -echo "" -echo "Virtual environment created at: $VENV_DIR" - diff --git a/apps/sim/lib/guardrails/validate_pii.py b/apps/sim/lib/guardrails/validate_pii.py deleted file mode 100644 index d475b96e23..0000000000 --- a/apps/sim/lib/guardrails/validate_pii.py +++ /dev/null @@ -1,260 +0,0 @@ -#!/usr/bin/env python3 -""" -PII Detection Validator using Microsoft Presidio - -Detects personally identifiable information (PII) in text and either: -- Blocks the request if PII is detected (block mode) -- Masks the PII and returns the masked text (mask mode) -""" - -import sys -import json -from typing import List, Dict, Any - -try: - from presidio_analyzer import AnalyzerEngine, Pattern, PatternRecognizer - from presidio_anonymizer import AnonymizerEngine - from presidio_anonymizer.entities import OperatorConfig -except ImportError: - print(json.dumps({ - "passed": False, - "error": "Presidio not installed. Run: pip install presidio-analyzer presidio-anonymizer", - "detectedEntities": [] - })) - sys.exit(0) - - -class VinRecognizer(PatternRecognizer): - """ - Recognizes Vehicle Identification Numbers (17 chars, A-Z/0-9 excluding - I/O/Q) and validates the ISO 3779 check digit (position 9). Validation makes - accidental matches on arbitrary 17-char codes (request ids, SKUs, tokens) - extremely unlikely. Note: some non-North-American VINs don't use the check - digit and will be skipped — an intentional bias toward precision. - """ - - _TRANSLIT = { - **{str(d): d for d in range(10)}, - "A": 1, "B": 2, "C": 3, "D": 4, "E": 5, "F": 6, "G": 7, "H": 8, - "J": 1, "K": 2, "L": 3, "M": 4, "N": 5, "P": 7, "R": 9, - "S": 2, "T": 3, "U": 4, "V": 5, "W": 6, "X": 7, "Y": 8, "Z": 9, - } - _WEIGHTS = [8, 7, 6, 5, 4, 3, 2, 10, 0, 9, 8, 7, 6, 5, 4, 3, 2] - - def validate_result(self, pattern_text: str): - vin = pattern_text.upper() - if len(vin) != 17: - return False - try: - total = sum(self._TRANSLIT[c] * w for c, w in zip(vin, self._WEIGHTS)) - except KeyError: - return False - check = total % 11 - expected = "X" if check == 10 else str(check) - return vin[8] == expected - - -def build_analyzer() -> "AnalyzerEngine": - """ - AnalyzerEngine with custom recognizers registered on top of the Presidio - defaults. Adds a check-digit-validated VIN recognizer. - """ - analyzer = AnalyzerEngine() - vin_pattern = Pattern(name="vin", regex=r"\b[A-HJ-NPR-Z0-9]{17}\b", score=0.7) - vin_recognizer = VinRecognizer( - supported_entity="VIN", - patterns=[vin_pattern], - context=["vin", "vehicle", "chassis"], - ) - analyzer.registry.add_recognizer(vin_recognizer) - return analyzer - - -def detect_pii( - text: str, - entity_types: List[str], - mode: str = "block", - language: str = "en" -) -> Dict[str, Any]: - """ - Detect PII in text using Presidio - - Args: - text: Input text to analyze - entity_types: List of PII entity types to detect (e.g., ["PERSON", "EMAIL_ADDRESS"]) - mode: "block" to fail validation if PII found, "mask" to return masked text - language: Language code (default: "en") - - Returns: - Dictionary with validation result - """ - try: - # Initialize Presidio engines - analyzer = build_analyzer() - - # Analyze text for PII - results = analyzer.analyze( - text=text, - entities=entity_types if entity_types else None, # None = detect all - language=language - ) - - # Extract detected entities - detected_entities = [] - for result in results: - detected_entities.append({ - "type": result.entity_type, - "start": result.start, - "end": result.end, - "score": result.score, - "text": text[result.start:result.end] - }) - - # If no PII detected, validation passes - if not results: - return { - "passed": True, - "detectedEntities": [], - "maskedText": None - } - - # Block mode: fail validation if PII detected - if mode == "block": - entity_summary = {} - for entity in detected_entities: - entity_type = entity["type"] - entity_summary[entity_type] = entity_summary.get(entity_type, 0) + 1 - - summary_str = ", ".join([f"{count} {etype}" for etype, count in entity_summary.items()]) - - return { - "passed": False, - "error": f"PII detected: {summary_str}", - "detectedEntities": detected_entities, - "maskedText": None - } - - # Mask mode: anonymize PII and return masked text - elif mode == "mask": - anonymizer = AnonymizerEngine() - - # Use as the replacement pattern - operators = {} - for entity_type in set([r.entity_type for r in results]): - operators[entity_type] = OperatorConfig("replace", {"new_value": f"<{entity_type}>"}) - - anonymized_result = anonymizer.anonymize( - text=text, - analyzer_results=results, - operators=operators - ) - - return { - "passed": True, - "detectedEntities": detected_entities, - "maskedText": anonymized_result.text - } - - else: - return { - "passed": False, - "error": f"Invalid mode: {mode}. Must be 'block' or 'mask'", - "detectedEntities": [] - } - - except Exception as e: - return { - "passed": False, - "error": f"PII detection failed: {str(e)}", - "detectedEntities": [] - } - - -def mask_batch( - texts: List[str], - entity_types: List[str], - language: str = "en" -) -> Dict[str, Any]: - """ - Mask PII across many strings in a single process, reusing one analyzer + - anonymizer instance (engine construction loads the spaCy model and is the - dominant cost). Returns masked text per input, in input order; strings with - no detected PII are returned unchanged so callers can substitute directly. - """ - analyzer = build_analyzer() - anonymizer = AnonymizerEngine() - entities = entity_types if entity_types else None - - results = [] - for text in texts: - if not text: - results.append({"maskedText": text}) - continue - analyzer_results = analyzer.analyze(text=text, entities=entities, language=language) - if not analyzer_results: - results.append({"maskedText": text}) - continue - operators = { - entity_type: OperatorConfig("replace", {"new_value": f"<{entity_type}>"}) - for entity_type in set([r.entity_type for r in analyzer_results]) - } - anonymized = anonymizer.anonymize( - text=text, - analyzer_results=analyzer_results, - operators=operators - ) - results.append({"maskedText": anonymized.text}) - - return {"passed": True, "results": results} - - -def main(): - """Main entry point for CLI usage""" - try: - # Read input from stdin - input_data = sys.stdin.read() - data = json.loads(input_data) - - entity_types = data.get("entityTypes", []) - language = data.get("language", "en") - - # Batch mask mode: an array of texts processed with one warm engine pair. - if "texts" in data: - texts = data.get("texts", []) - result = mask_batch(texts, entity_types, language) - print(f"__SIM_RESULT__={json.dumps(result)}") - return - - text = data.get("text", "") - mode = data.get("mode", "block") - - # Validate inputs - if not text: - result = { - "passed": False, - "error": "No text provided", - "detectedEntities": [] - } - else: - result = detect_pii(text, entity_types, mode, language) - - # Output result with marker for parsing - print(f"__SIM_RESULT__={json.dumps(result)}") - - except json.JSONDecodeError as e: - print(f"__SIM_RESULT__={json.dumps({ - 'passed': False, - 'error': f'Invalid JSON input: {str(e)}', - 'detectedEntities': [] - })}") - except Exception as e: - print(f"__SIM_RESULT__={json.dumps({ - 'passed': False, - 'error': f'Unexpected error: {str(e)}', - 'detectedEntities': [] - })}") - - -if __name__ == "__main__": - main() - diff --git a/apps/sim/lib/guardrails/validate_pii.test.ts b/apps/sim/lib/guardrails/validate_pii.test.ts new file mode 100644 index 0000000000..0ba1c585bc --- /dev/null +++ b/apps/sim/lib/guardrails/validate_pii.test.ts @@ -0,0 +1,118 @@ +/** + * @vitest-environment node + */ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { maskPIIBatch, validatePII } from '@/lib/guardrails/validate_pii' + +interface Span { + entity_type: string + start: number + end: number + score: number +} + +/** Mimic the Presidio anonymizer's default `replace`: each span → ``. */ +function applyReplace(text: string, results: Span[]): string { + let out = text + for (const s of [...results].sort((a, b) => b.start - a.start)) { + out = `${out.slice(0, s.start)}<${s.entity_type}>${out.slice(s.end)}` + } + return out +} + +/** Analyzer mock: flags `a@b.com` as EMAIL_ADDRESS when that entity is in scope. */ +function emailSpans(text: string, entities: string[] | undefined): Span[] { + if (entities && !entities.includes('EMAIL_ADDRESS')) return [] + const idx = text.indexOf('a@b.com') + return idx === -1 ? [] : [{ entity_type: 'EMAIL_ADDRESS', start: idx, end: idx + 7, score: 0.9 }] +} + +describe('validate_pii (Presidio sidecar)', () => { + let analyzeBodies: Array<{ text: string; language: string; entities?: string[] }> + let fetchMock: ReturnType + + beforeEach(() => { + analyzeBodies = [] + fetchMock = vi.fn(async (url: string, init: { body: string }) => { + const body = JSON.parse(init.body) + if (url.includes('/analyze')) { + analyzeBodies.push({ text: body.text, language: body.language, entities: body.entities }) + return new Response(JSON.stringify(emailSpans(body.text, body.entities)), { status: 200 }) + } + // /anonymize + return new Response( + JSON.stringify({ text: applyReplace(body.text, body.analyzer_results) }), + { + status: 200, + } + ) + }) + vi.stubGlobal('fetch', fetchMock) + }) + + afterEach(() => vi.unstubAllGlobals()) + + describe('maskPIIBatch', () => { + it('masks detected entities, preserving input order', async () => { + const out = await maskPIIBatch(['email a@b.com', 'nothing here'], []) + expect(out[0]).toBe('email ') + expect(out[1]).toBe('nothing here') + }) + + it('forwards entityTypes (and language) to the analyzer; empty ⇒ omitted (all)', async () => { + await maskPIIBatch(['mail a@b.com'], ['EMAIL_ADDRESS', 'PERSON'], 'es') + expect(analyzeBodies[0].entities).toEqual(['EMAIL_ADDRESS', 'PERSON']) + expect(analyzeBodies[0].language).toBe('es') + + analyzeBodies.length = 0 + await maskPIIBatch(['mail a@b.com'], []) + expect(analyzeBodies[0].entities).toBeUndefined() + }) + + it('returns [] for empty input and leaves empty strings untouched', async () => { + expect(await maskPIIBatch([], [])).toEqual([]) + expect(await maskPIIBatch([''], [])).toEqual(['']) + }) + + it('throws on a sidecar failure so the caller can scrub', async () => { + fetchMock.mockResolvedValueOnce(new Response('boom', { status: 500 })) + await expect(maskPIIBatch(['email a@b.com'], [])).rejects.toThrow(/Presidio analyze failed/) + }) + }) + + describe('validatePII', () => { + it('block mode fails with a summary when PII is detected', async () => { + const res = await validatePII({ + text: 'reach me at a@b.com', + entityTypes: [], + mode: 'block', + requestId: 'r1', + }) + expect(res.passed).toBe(false) + expect(res.error).toContain('EMAIL_ADDRESS') + expect(res.detectedEntities).toHaveLength(1) + }) + + it('mask mode returns masked text', async () => { + const res = await validatePII({ + text: 'mail a@b.com', + entityTypes: [], + mode: 'mask', + requestId: 'r2', + }) + expect(res.passed).toBe(true) + expect(res.maskedText).toBe('mail ') + }) + + it('passes clean text', async () => { + const res = await validatePII({ + text: 'nothing to see', + entityTypes: [], + mode: 'block', + requestId: 'r3', + }) + expect(res.passed).toBe(true) + expect(res.detectedEntities).toHaveLength(0) + }) + }) +}) diff --git a/apps/sim/lib/guardrails/validate_pii.ts b/apps/sim/lib/guardrails/validate_pii.ts index ba6886bb92..a24c8f880e 100644 --- a/apps/sim/lib/guardrails/validate_pii.ts +++ b/apps/sim/lib/guardrails/validate_pii.ts @@ -1,17 +1,18 @@ -import { spawn } from 'child_process' -import fs from 'fs' -import path from 'path' import { createLogger } from '@sim/logger' +import { getErrorMessage } from '@sim/utils/errors' +import { env } from '@/lib/core/config/env' +import { mapWithConcurrency } from '@/lib/core/utils/concurrency' const logger = createLogger('PIIValidator') -const DEFAULT_TIMEOUT = 30000 // 30 seconds -/** - * Max total bytes of text sent to a single Presidio subprocess. spaCy NER is the - * bottleneck, so large payloads are split into multiple short calls instead of - * one that risks the 30s timeout. - */ -const PII_CHUNK_MAX_BYTES = 256 * 1024 +/** Just above the analyzer's spaCy NER budget so a stuck sidecar aborts gracefully. */ +const REQUEST_TIMEOUT_MS = 45_000 + +/** Concurrent per-string sidecar calls within one batch; the warm model handles parallelism. */ +const MASK_CONCURRENCY = 8 + +/** Single Presidio sidecar serving both /analyze and /anonymize (VIN is native there). */ +const PII_URL = env.PII_URL || 'http://localhost:5001' export interface PIIValidationInput { text: string @@ -36,12 +37,65 @@ export interface PIIValidationResult { maskedText?: string } +interface AnalyzerSpan { + entity_type: string + start: number + end: number + score: number +} + /** - * Validate text for PII using Microsoft Presidio + * Detect PII spans via the Presidio analyzer. An empty `entityTypes` ⇒ detect all. + * Throws on transport/HTTP failure so callers can apply their own fail-safe. + */ +async function analyze( + text: string, + entityTypes: string[], + language: string +): Promise { + const entities = entityTypes.length > 0 ? entityTypes : undefined + + // boundary-raw-fetch: internal call to the Presidio analyzer sidecar over localhost + const response = await fetch(`${PII_URL}/analyze`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ text, language, ...(entities ? { entities } : {}) }), + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), + }) + if (!response.ok) { + const detail = await response.text().catch(() => '') + throw new Error(`Presidio analyze failed (${response.status}): ${detail.slice(0, 200)}`) + } + return (await response.json()) as AnalyzerSpan[] +} + +/** + * Mask spans via the Presidio anonymizer sidecar. Omitting `anonymizers` uses the + * default `replace` operator, which yields ``. Throws on failure. + */ +async function anonymize(text: string, spans: AnalyzerSpan[]): Promise { + if (spans.length === 0) return text + + // boundary-raw-fetch: internal call to the Presidio anonymizer sidecar over localhost + const response = await fetch(`${PII_URL}/anonymize`, { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ text, analyzer_results: spans }), + signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), + }) + if (!response.ok) { + const detail = await response.text().catch(() => '') + throw new Error(`Presidio anonymize failed (${response.status}): ${detail.slice(0, 200)}`) + } + const data = (await response.json()) as { text: string } + return data.text +} + +/** + * Validate text for PII using the Presidio sidecar. * - * Supports two modes: - * - block: Fails validation if any PII is detected - * - mask: Passes validation and returns masked text with PII replaced + * - block: fails validation if any PII is detected + * - mask: passes and returns masked text with PII replaced by `` */ export async function validatePII(input: PIIValidationInput): Promise { const { text, entityTypes, mode, language = 'en', requestId } = input @@ -54,41 +108,60 @@ export async function validatePII(input: PIIValidationInput): Promise ({ + type: s.entity_type, + start: s.start, + end: s.end, + score: s.score, + text: text.slice(s.start, s.end), + })) + + if (spans.length === 0) { + logger.info(`[${requestId}] PII validation completed`, { passed: true, detectedCount: 0 }) + return { passed: true, detectedEntities: [], maskedText: mode === 'mask' ? text : undefined } + } + + if (mode === 'block') { + const counts = new Map() + for (const e of detectedEntities) counts.set(e.type, (counts.get(e.type) ?? 0) + 1) + const summary = Array.from(counts.entries()) + .map(([type, count]) => `${count} ${type}`) + .join(', ') + logger.info(`[${requestId}] PII validation completed`, { + passed: false, + detectedCount: detectedEntities.length, + }) + return { passed: false, error: `PII detected: ${summary}`, detectedEntities } + } + + // mask mode: the anonymizer replaces every span with ``. + const maskedText = await anonymize(text, spans) logger.info(`[${requestId}] PII validation completed`, { - passed: result.passed, - detectedCount: result.detectedEntities.length, - hasMaskedText: !!result.maskedText, + passed: true, + detectedCount: detectedEntities.length, + hasMaskedText: true, }) - - return result - } catch (error: any) { - logger.error(`[${requestId}] PII validation failed`, { - error: error.message, - }) - + return { passed: true, detectedEntities, maskedText } + } catch (error) { + logger.error(`[${requestId}] PII validation failed`, { error: getErrorMessage(error) }) return { passed: false, - error: `PII validation failed: ${error.message}`, + error: `PII validation failed: ${getErrorMessage(error)}`, detectedEntities: [], } } } -interface PIIMaskBatchResult { - passed: boolean - error?: string - results?: { maskedText: string }[] -} - /** - * Mask PII across many strings, preserving input order. Strings are grouped into - * byte-budgeted chunks so no single subprocess exceeds {@link PII_CHUNK_MAX_BYTES} - * (keeping each call well under the 30s timeout). One Presidio engine pair is - * reused per subprocess invocation. Rejects on any subprocess failure so callers - * can apply their own fail-safe. + * Mask PII across many strings via the Presidio sidecar, preserving input order. + * Each string runs analyze → anonymize; strings with no detected PII are returned + * unchanged. Calls run with bounded concurrency: the sidecar's model is warm, so + * the bottleneck is round-trip latency, and a batch of thousands of small leaves + * would otherwise exceed the caller's request timeout if run strictly sequentially. + * Rejects on any sidecar failure (which fails the whole batch) so callers can apply + * their own fail-safe (scrub). */ export async function maskPIIBatch( texts: string[], @@ -97,223 +170,10 @@ export async function maskPIIBatch( ): Promise { if (texts.length === 0) return [] - const chunks: string[][] = [] - let current: string[] = [] - let currentBytes = 0 - for (const text of texts) { - const bytes = Buffer.byteLength(text, 'utf8') - if (current.length > 0 && currentBytes + bytes > PII_CHUNK_MAX_BYTES) { - chunks.push(current) - current = [] - currentBytes = 0 - } - current.push(text) - currentBytes += bytes - } - if (current.length > 0) chunks.push(current) - - const masked: string[] = [] - for (const chunk of chunks) { - const result = await runPythonScript({ - texts: chunk, - entityTypes, - mode: 'mask', - language, - }) - if (!result.passed || !result.results || result.results.length !== chunk.length) { - throw new Error(result.error || 'PII batch masking returned an unexpected result') - } - for (const item of result.results) masked.push(item.maskedText) - } - - return masked -} - -/** - * Spawn the Presidio Python script, write the payload to stdin as JSON, and parse - * the `__SIM_RESULT__=` marker from stdout. Rejects on non-zero exit, timeout, - * spawn failure, or a missing/unparseable marker. - */ -function runPythonScript(payload: Record): Promise { - return new Promise((resolve, reject) => { - const guardrailsDir = path.join(process.cwd(), 'lib/guardrails') - const scriptPath = path.join(guardrailsDir, 'validate_pii.py') - const venvPython = path.join(guardrailsDir, 'venv/bin/python3') - const pythonCmd = fs.existsSync(venvPython) ? venvPython : 'python3' - - const python = spawn(pythonCmd, [scriptPath]) - let stdout = '' - let stderr = '' - - const timeout = setTimeout(() => { - python.kill() - reject(new Error('PII processing timeout')) - }, DEFAULT_TIMEOUT) - - // stdin errors (e.g. EPIPE when the child exits before draining the payload — - // chunks can exceed the OS pipe buffer) emit on stdin, not the process. Without - // a listener Node throws an unhandled 'error' and crashes; funnel it into the - // promise so the caller's fail-safe scrub path handles it. - python.stdin.on('error', (error: Error) => { - clearTimeout(timeout) - reject(new Error(`PII script stdin error: ${error.message}`)) - }) - python.stdin.write(JSON.stringify(payload)) - python.stdin.end() - python.stdout.on('data', (data) => { - stdout += data.toString() - }) - python.stderr.on('data', (data) => { - stderr += data.toString() - }) - - python.on('close', (code) => { - clearTimeout(timeout) - if (code !== 0) { - reject(new Error(stderr || `PII script exited with code ${code}`)) - return - } - const prefix = '__SIM_RESULT__=' - const marker = stdout.split('\n').find((l) => l.startsWith(prefix)) - if (!marker) { - reject(new Error(`No result marker in PII script output: ${stdout.substring(0, 200)}`)) - return - } - try { - resolve(JSON.parse(marker.slice(prefix.length)) as T) - } catch (error: any) { - reject(new Error(`Failed to parse PII script result: ${error.message}`)) - } - }) - - python.on('error', (error) => { - clearTimeout(timeout) - reject( - new Error( - `Failed to execute Python: ${error.message}. Make sure Python 3 and Presidio are installed.` - ) - ) - }) - }) -} - -/** - * Execute Python PII detection script - */ -async function executePythonPIIDetection( - text: string, - entityTypes: string[], - mode: string, - language: string, - requestId: string -): Promise { - return new Promise((resolve, reject) => { - // Use path relative to project root - // In Next.js, process.cwd() returns the project root - const guardrailsDir = path.join(process.cwd(), 'lib/guardrails') - const scriptPath = path.join(guardrailsDir, 'validate_pii.py') - const venvPython = path.join(guardrailsDir, 'venv/bin/python3') - - // Use venv Python if it exists, otherwise fall back to system python3 - const pythonCmd = fs.existsSync(venvPython) ? venvPython : 'python3' - - const python = spawn(pythonCmd, [scriptPath]) - - let stdout = '' - let stderr = '' - - const timeout = setTimeout(() => { - python.kill() - reject(new Error('PII validation timeout')) - }, DEFAULT_TIMEOUT) - - // Write input to stdin as JSON - const inputData = JSON.stringify({ - text, - entityTypes, - mode, - language, - }) - // See runPythonScript: stdin errors (EPIPE on early child exit) must be - // caught here or Node throws an unhandled 'error' and crashes the process. - python.stdin.on('error', (error: Error) => { - clearTimeout(timeout) - reject(new Error(`Failed to write to Python: ${error.message}`)) - }) - python.stdin.write(inputData) - python.stdin.end() - - python.stdout.on('data', (data) => { - stdout += data.toString() - }) - - python.stderr.on('data', (data) => { - stderr += data.toString() - }) - - python.on('close', (code) => { - clearTimeout(timeout) - - if (code !== 0) { - logger.error(`[${requestId}] Python PII detection failed`, { - code, - stderr, - }) - resolve({ - passed: false, - error: stderr || 'PII detection failed', - detectedEntities: [], - }) - return - } - - // Parse result from stdout - try { - const prefix = '__SIM_RESULT__=' - const lines = stdout.split('\n') - const marker = lines.find((l) => l.startsWith(prefix)) - - if (marker) { - const jsonPart = marker.slice(prefix.length) - const result = JSON.parse(jsonPart) - resolve(result) - } else { - logger.error(`[${requestId}] No result marker found`, { - stdout, - stderr, - stdoutLines: lines, - }) - resolve({ - passed: false, - error: `No result marker found in output. stdout: ${stdout.substring(0, 200)}, stderr: ${stderr.substring(0, 200)}`, - detectedEntities: [], - }) - } - } catch (error: any) { - logger.error(`[${requestId}] Failed to parse Python result`, { - error: error.message, - stdout, - stderr, - }) - resolve({ - passed: false, - error: `Failed to parse result: ${error.message}. stdout: ${stdout.substring(0, 200)}`, - detectedEntities: [], - }) - } - }) - - python.on('error', (error) => { - clearTimeout(timeout) - logger.error(`[${requestId}] Failed to spawn Python process`, { - error: error.message, - }) - reject( - new Error( - `Failed to execute Python: ${error.message}. Make sure Python 3 and Presidio are installed.` - ) - ) - }) + return mapWithConcurrency(texts, MASK_CONCURRENCY, async (text) => { + if (!text) return text + const spans = await analyze(text, entityTypes, language) + return anonymize(text, spans) }) } diff --git a/apps/sim/lib/logs/execution/logger.ts b/apps/sim/lib/logs/execution/logger.ts index 9a531b7293..e68ad3100f 100644 --- a/apps/sim/lib/logs/execution/logger.ts +++ b/apps/sim/lib/logs/execution/logger.ts @@ -620,7 +620,10 @@ export class ExecutionLogger implements IExecutionLoggerService { const config = resolveEffectivePiiRedaction({ orgSettings: row.orgSettings, workspaceId }) if (!config.enabled) return payload - return redactPIIFromExecution(payload, { entityTypes: config.entityTypes }) + return redactPIIFromExecution(payload, { + entityTypes: config.entityTypes, + language: config.language, + }) } async completeWorkflowExecution(params: { diff --git a/apps/sim/lib/logs/execution/pii-redaction.test.ts b/apps/sim/lib/logs/execution/pii-redaction.test.ts index dccbc59cc3..5a2da7a599 100644 --- a/apps/sim/lib/logs/execution/pii-redaction.test.ts +++ b/apps/sim/lib/logs/execution/pii-redaction.test.ts @@ -7,8 +7,8 @@ const { mockMaskPIIBatch } = vi.hoisted(() => ({ mockMaskPIIBatch: vi.fn(), })) -vi.mock('@/lib/guardrails/validate_pii', () => ({ - maskPIIBatch: mockMaskPIIBatch, +vi.mock('@/lib/guardrails/mask-client', () => ({ + maskPIIBatchViaHttp: mockMaskPIIBatch, })) import { REDACTION_FAILED_MARKER, redactPIIFromExecution } from '@/lib/logs/execution/pii-redaction' diff --git a/apps/sim/lib/logs/execution/pii-redaction.ts b/apps/sim/lib/logs/execution/pii-redaction.ts index 7b4794fd48..8cd0fac532 100644 --- a/apps/sim/lib/logs/execution/pii-redaction.ts +++ b/apps/sim/lib/logs/execution/pii-redaction.ts @@ -1,5 +1,6 @@ import { createLogger } from '@sim/logger' import { getErrorMessage } from '@sim/utils/errors' +import { maskPIIBatchViaHttp } from '@/lib/guardrails/mask-client' const logger = createLogger('PiiRedaction') @@ -158,11 +159,9 @@ export async function redactPIIFromExecution( masked = collected.map(() => REDACTION_FAILED_MARKER) } else { try { - // Lazy import keeps the Python-spawning guardrails module (child_process + - // a `lib/guardrails` dir reference) out of the static middleware/RSC graph; - // it's only loaded at runtime on the Node log-persist path. - const { maskPIIBatch } = await import('@/lib/guardrails/validate_pii') - masked = await maskPIIBatch(collected, entityTypes, language) + // Presidio runs only in the app container; the persist path also runs in + // the trigger.dev runtime, so masking always goes over HTTP to the app. + masked = await maskPIIBatchViaHttp(collected, entityTypes, language) } catch (error) { logger.error('PII masking failed; scrubbing text to avoid leaking PII', { error: getErrorMessage(error), diff --git a/docker/app.Dockerfile b/docker/app.Dockerfile index 67eb5f02c7..ff0ea1ccc2 100644 --- a/docker/app.Dockerfile +++ b/docker/app.Dockerfile @@ -114,16 +114,8 @@ COPY --from=builder --chown=nextjs:nodejs /app/apps/sim/lib/execution/isolated-v # apps/sim/lib/execution/sandbox/bundles/build.ts to regenerate. COPY --from=builder --chown=nextjs:nodejs /app/apps/sim/lib/execution/sandbox/bundles ./apps/sim/lib/execution/sandbox/bundles -# Guardrails setup with pip caching -COPY --from=builder --chown=nextjs:nodejs /app/apps/sim/lib/guardrails/requirements.txt ./apps/sim/lib/guardrails/requirements.txt -COPY --from=builder --chown=nextjs:nodejs /app/apps/sim/lib/guardrails/validate_pii.py ./apps/sim/lib/guardrails/validate_pii.py - -# Install Python dependencies with pip cache mount for faster rebuilds -RUN --mount=type=cache,target=/root/.cache/pip \ - python3 -m venv ./apps/sim/lib/guardrails/venv && \ - ./apps/sim/lib/guardrails/venv/bin/pip install --upgrade pip && \ - ./apps/sim/lib/guardrails/venv/bin/pip install -r ./apps/sim/lib/guardrails/requirements.txt && \ - chown -R nextjs:nodejs /app/apps/sim/lib/guardrails +# Guardrails PII runs in dedicated Presidio sidecar containers (analyzer + +# anonymizer), reached over localhost — no Python/Presidio in this image. # Create .next/cache directory with correct ownership RUN mkdir -p apps/sim/.next/cache && \ diff --git a/packages/db/schema.ts b/packages/db/schema.ts index 8718883751..f066c19ad4 100644 --- a/packages/db/schema.ts +++ b/packages/db/schema.ts @@ -1077,6 +1077,8 @@ export interface PiiRedactionRule { entityTypes: string[] /** `null` = all workspaces; otherwise the single targeted workspace. */ workspaceId: string | null + /** Language whose Presidio recognizers apply (e.g. 'en', 'es'); defaults to English. */ + language?: string } /** diff --git a/scripts/check-api-validation-contracts.ts b/scripts/check-api-validation-contracts.ts index 09744c629b..17f0a25fa2 100644 --- a/scripts/check-api-validation-contracts.ts +++ b/scripts/check-api-validation-contracts.ts @@ -9,8 +9,8 @@ const QUERY_HOOKS_DIR = path.join(ROOT, 'apps/sim/hooks/queries') const SELECTOR_HOOKS_DIR = path.join(ROOT, 'apps/sim/hooks/selectors') const BASELINE = { - totalRoutes: 859, - zodRoutes: 859, + totalRoutes: 860, + zodRoutes: 860, nonZodRoutes: 0, } as const