id: presidio-residual-pii-load-gate
namespace: company.team
inputs:
- id: spacy_model
type: SELECT
displayName: spaCy model
description: NER model Presidio loads. en_core_web_sm is 12.8 MB and keeps the
demo small. en_core_web_lg is 400.7 MB and has much better PERSON recall,
so use it in production.
defaults: en_core_web_sm
values:
- en_core_web_sm
- en_core_web_lg
- id: score_threshold
type: FLOAT
displayName: Analyzer score threshold
description: Minimum Presidio confidence for a span to be redacted. Lower
catches more and redacts more false positives.
defaults: 0.4
- id: max_residual_entities
type: INT
displayName: Max residual entities
description: How many residual findings the independent scan may report before
the load is blocked. 0 means any survivor blocks the batch.
defaults: 0
- id: inject_unsupported_ids
type: BOOL
displayName: Inject unsupported identifiers
description: Put identifiers Presidio ships no recognizer for into the sample
batch. True shows the gate blocking. Set it to false to walk the publish
branch.
defaults: true
- id: org_id_patterns
type: STRING
displayName: Org id patterns
description: Comma separated regular expressions for identifiers only your
organisation knows about. The residual scan looks for these after the
scrub.
defaults: EMP-\d{4},ACCT-\d{5}
- id: phone_region
type: STRING
displayName: Default phone region
description: ISO 3166-1 alpha-2 region used to resolve phone numbers written
without a country code.
defaults: US
- id: warehouse_path
type: STRING
displayName: Warehouse path
description: DuckDB file the clean batch is loaded into. /tmp is wiped with the
container, so point this at a mounted volume to keep history.
defaults: /tmp/pii_gate.duckdb
triggers:
- id: nightly_batch
type: io.kestra.plugin.core.trigger.Schedule
description: Scrub and gate yesterday's free-text tickets every night at 03:00.
cron: "0 3 * * *"
tasks:
- id: generate_batch
type: io.kestra.plugin.scripts.python.Script
description: Write a synthetic batch of free-text support tickets. Every value
is invented for the demo. Replace this task with your own extract.
containerImage: python:3.12-slim
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
env:
INJECT_UNSUPPORTED: "{{ inputs.inject_unsupported_ids }}"
outputFiles:
- tickets.csv
script: |
import csv
import json
import os
inject = os.environ["INJECT_UNSUPPORTED"].lower() == "true"
# Synthetic tickets. The card number is the Visa test value, the IBAN is the
# documentation example, every name, address and mailbox is invented.
rows = [
(
"TCK-1001",
"email",
"Caller said her name is Priya Raghunathan and asked us to mail the "
"refund confirmation to priya.raghunathan@example.invalid. Best number "
"is +1 212 555 0147 before 6pm.",
),
(
"TCK-1002",
"chat",
"Customer Dmitri Alvarez-Okonkwo read out the card on the account, "
"4111 1111 1111 1111, expiry 04/29. He wants the duplicate charge "
"reversed to the same card.",
),
(
"TCK-1003",
"phone",
"Payment bounced. Account holder gave IBAN GB82WEST12345698765432 and "
"a second contact on 0161 496 0215. Escalating to billing.",
),
(
"TCK-1004",
"email",
"Thanks for the update. Please copy marisol.betancourt@example.invalid "
"on the next reply, she is handling the account while I am away.",
),
]
if inject:
rows += [
(
"TCK-1005",
"chat",
"Handled by EMP-0442 on the retention desk. Customer reference "
"ACCT-77301 is flagged for a goodwill credit, approval from EMP-0913.",
),
(
"TCK-1006",
"email",
"Card on file reads 4111.1111.1111.1111 in the legacy export, which "
"is why the reconciliation job keeps rejecting ACCT-77301.",
),
]
with open("tickets.csv", "w", newline="", encoding="utf-8") as fh:
writer = csv.writer(fh)
writer.writerow(["ticket_id", "channel", "body"])
writer.writerows(rows)
print("wrote", len(rows), "synthetic tickets, unsupported identifiers injected:", inject)
print("::" + json.dumps({"outputs": {"ticket_count": len(rows)}}) + "::")
- id: scrub_and_rescan
type: io.kestra.plugin.scripts.python.Script
description: >
Scrub the free-text column with Presidio, count what each recognizer
removed, then re-scan the scrubbed text with an independent deterministic
battery. The same AnalyzerEngine is also re-run so the report can show how
little a same-detector re-scan proves.
containerImage: python:3.12-slim
taskRunner:
type: io.kestra.plugin.scripts.runner.docker.Docker
beforeCommands:
- pip install --quiet --no-cache-dir presidio-analyzer==2.2.364
presidio-anonymizer==2.2.364
- python -m spacy download {{ inputs.spacy_model }}
inputFiles:
tickets.csv: "{{ outputs.generate_batch.outputFiles['tickets.csv'] }}"
env:
SPACY_MODEL: "{{ inputs.spacy_model }}"
SCORE_THRESHOLD: "{{ inputs.score_threshold }}"
ORG_ID_PATTERNS: "{{ inputs.org_id_patterns }}"
PHONE_REGION: "{{ inputs.phone_region }}"
outputFiles:
- scrubbed.csv
- redaction-report.json
- residual-findings.json
script: |
import csv
import json
import os
import re
from collections import Counter
import phonenumbers
from presidio_analyzer import AnalyzerEngine
from presidio_analyzer.nlp_engine import NlpEngineProvider
from presidio_anonymizer import AnonymizerEngine
MODEL = os.environ["SPACY_MODEL"]
THRESHOLD = float(os.environ["SCORE_THRESHOLD"])
REGION = os.environ["PHONE_REGION"]
ORG_PATTERNS = [p for p in os.environ["ORG_ID_PATTERNS"].split(",") if p.strip()]
# ---------------------------------------------------------------- detector 2
# Independent of Presidio on purpose. Re-running the analyzer over its own
# output returns almost nothing by construction, so it cannot be the proof.
# Each check below is a closed-form validator, not a model.
CARD_CANDIDATE = re.compile(r"(?<![0-9])(?:[0-9][ .-]?){12,18}[0-9](?![0-9])")
IBAN_CANDIDATE = re.compile(r"\b[A-Z]{2}[0-9]{2}[A-Z0-9]{11,30}\b")
EMAIL_STRICT = re.compile(
r"[A-Za-z0-9!#$%&'*+/=?^_`{|}~-]+"
r"(?:\.[A-Za-z0-9!#$%&'*+/=?^_`{|}~-]+)*"
r"@[A-Za-z0-9](?:[A-Za-z0-9-]{0,61}[A-Za-z0-9])?"
r"(?:\.[A-Za-z0-9](?:[A-Za-z0-9-]{0,61}[A-Za-z0-9])?)+"
)
E164_MAX_DIGITS = 15 # ITU-T E.164, maximum digits excluding the international prefix
def luhn_valid(digits: str) -> bool:
# Double every second digit from the right, subtract 9 from any result
# over 9, sum, and a valid number is divisible by 10.
total = 0
for index, char in enumerate(reversed(digits)):
value = int(char)
if index % 2 == 1:
value *= 2
if value > 9:
value -= 9
total += value
return total % 10 == 0
def iban_valid(candidate: str) -> bool:
# Move the first four characters to the end, map A=10 through Z=35,
# then the resulting integer mod 97 must equal 1.
rotated = candidate[4:] + candidate[:4]
remainder = 0
for char in rotated:
if char.isdigit():
remainder = (remainder * 10 + int(char)) % 97
elif "A" <= char <= "Z":
remainder = (remainder * 100 + ord(char) - 55) % 97
else:
return False
return remainder == 1
def mask(value: str) -> str:
tail = value[-2:] if len(value) > 2 else ""
return "*" * max(len(value) - 2, 1) + tail
def residual_scan(text: str, ticket_id: str):
findings = []
for match in CARD_CANDIDATE.finditer(text):
digits = re.sub(r"[^0-9]", "", match.group(0))
if 13 <= len(digits) <= 19 and luhn_valid(digits):
findings.append(("CARD_NUMBER_LUHN", match.group(0)))
for match in IBAN_CANDIDATE.finditer(text):
if iban_valid(match.group(0)):
findings.append(("IBAN_MOD97", match.group(0)))
for match in EMAIL_STRICT.finditer(text):
findings.append(("EMAIL_PATTERN", match.group(0)))
for match in phonenumbers.PhoneNumberMatcher(
text, REGION, leniency=phonenumbers.Leniency.VALID
):
formatted = phonenumbers.format_number(
match.number, phonenumbers.PhoneNumberFormat.E164
)
if len(formatted.lstrip("+")) <= E164_MAX_DIGITS:
findings.append(("PHONE_E164", formatted))
for pattern in ORG_PATTERNS:
for match in re.finditer(pattern, text):
findings.append(("ORG_IDENTIFIER", match.group(0)))
return [
{
"ticket_id": ticket_id,
"detector": "independent_battery",
"entity_type": kind,
"masked_value": mask(value),
"length": len(value),
}
for kind, value in findings
]
# --------------------------------------------------------------- detector 1
nlp_engine = NlpEngineProvider(
nlp_configuration={
"nlp_engine_name": "spacy",
"models": [{"lang_code": "en", "model_name": MODEL}],
}
).create_engine()
analyzer = AnalyzerEngine(nlp_engine=nlp_engine, supported_languages=["en"])
anonymizer = AnonymizerEngine()
with open("tickets.csv", newline="", encoding="utf-8") as fh:
tickets = list(csv.DictReader(fh))
redacted_counts = Counter()
first_pass_types = Counter()
residual = []
naive_hits = 0
scrubbed_rows = []
for ticket in tickets:
body = ticket["body"]
results = analyzer.analyze(text=body, language="en", score_threshold=THRESHOLD)
for result in results:
first_pass_types[result.entity_type] += 1
scrubbed = anonymizer.anonymize(text=body, analyzer_results=results)
for item in scrubbed.items:
redacted_counts[item.entity_type] += 1
# The hollow check, measured so the report can say what it is worth.
naive_hits += len(
analyzer.analyze(
text=scrubbed.text, language="en", score_threshold=THRESHOLD
)
)
residual.extend(residual_scan(scrubbed.text, ticket["ticket_id"]))
scrubbed_rows.append(
{
"ticket_id": ticket["ticket_id"],
"channel": ticket["channel"],
"body_scrubbed": scrubbed.text,
"entities_removed": len(scrubbed.items),
}
)
with open("scrubbed.csv", "w", newline="", encoding="utf-8") as fh:
writer = csv.DictWriter(
fh, fieldnames=["ticket_id", "channel", "body_scrubbed", "entities_removed"]
)
writer.writeheader()
writer.writerows(scrubbed_rows)
residual_types = sorted({item["entity_type"] for item in residual})
missed_by_presidio = sorted(set(residual_types) - set(first_pass_types))
report = {
"model": MODEL,
"score_threshold": THRESHOLD,
"tickets": len(tickets),
"redacted_per_entity": dict(sorted(redacted_counts.items())),
"redacted_total": sum(redacted_counts.values()),
"presidio_first_pass_entities": dict(sorted(first_pass_types.items())),
"same_detector_rescan_hits": naive_hits,
"independent_residual_count": len(residual),
"independent_residual_types": residual_types,
"types_presidio_never_flagged": missed_by_presidio,
}
with open("redaction-report.json", "w", encoding="utf-8") as fh:
json.dump(report, fh, indent=2, sort_keys=True)
with open("residual-findings.json", "w", encoding="utf-8") as fh:
json.dump(residual, fh, indent=2)
summary = ", ".join(f"{k}={v}" for k, v in sorted(redacted_counts.items())) or "none"
print(json.dumps(report, indent=2, sort_keys=True))
print(
"::"
+ json.dumps(
{
"outputs": {
"residual_entities": len(residual),
"residual_types": ", ".join(residual_types) or "none",
"missed_by_presidio": ", ".join(missed_by_presidio) or "none",
"same_detector_rescan_hits": naive_hits,
"redacted_total": sum(redacted_counts.values()),
"redaction_summary": summary,
"first_pass_summary": ", ".join(
f"{k}={v}" for k, v in sorted(first_pass_types.items())
)
or "none",
"rows": len(scrubbed_rows),
}
}
)
+ "::"
)
- id: log_measurement
type: io.kestra.plugin.core.log.Log
description: One line a reviewer can read without opening an artifact.
message: >-
Scrub measured on {{ outputs.scrub_and_rescan.vars.rows }} tickets with {{
inputs.spacy_model }}: removed {{
outputs.scrub_and_rescan.vars.redacted_total }} spans ({{
outputs.scrub_and_rescan.vars.redaction_summary }}). Presidio first pass
found {{ outputs.scrub_and_rescan.vars.first_pass_summary }}. Re-running
the same analyzer over the scrubbed text finds {{
outputs.scrub_and_rescan.vars.same_detector_rescan_hits }} spans. The
independent battery finds {{
outputs.scrub_and_rescan.vars.residual_entities }} residual ({{
outputs.scrub_and_rescan.vars.residual_types }}), of which Presidio never
flagged {{ outputs.scrub_and_rescan.vars.missed_by_presidio }}. Threshold
is {{ inputs.max_residual_entities }}.
- id: residual_gate
type: io.kestra.plugin.core.flow.If
description: The warehouse load lives behind this condition, so residual PII has
no path into it.
condition: "{{ outputs.scrub_and_rescan.vars.residual_entities <=
inputs.max_residual_entities }}"
then:
- id: load_warehouse
type: io.kestra.plugin.jdbc.duckdb.Queries
description: Load the scrubbed rows, stamped with the execution that proved them
clean.
url: "jdbc:duckdb:{{ inputs.warehouse_path }}"
inputFiles:
scrubbed.csv: "{{ outputs.scrub_and_rescan.outputFiles['scrubbed.csv'] }}"
sql: |
CREATE TABLE IF NOT EXISTS support_tickets_scrubbed (
ticket_id VARCHAR,
channel VARCHAR,
body_scrubbed VARCHAR,
entities_removed INTEGER,
execution_id VARCHAR,
loaded_at TIMESTAMP
);
INSERT INTO support_tickets_scrubbed
SELECT ticket_id,
channel,
body_scrubbed,
entities_removed,
'{{ execution.id }}' AS execution_id,
now() AS loaded_at
FROM read_csv_auto('{{ workingDir }}/scrubbed.csv', header=True);
- id: verify_load
type: io.kestra.plugin.jdbc.duckdb.Query
description: Read the rows back out of the warehouse so the load is proven
rather than assumed.
url: "jdbc:duckdb:{{ inputs.warehouse_path }}"
fetchType: FETCH_ONE
sql: |
SELECT count(*) AS loaded_rows,
sum(entities_removed) AS spans_removed
FROM support_tickets_scrubbed
WHERE execution_id = '{{ execution.id }}';
- id: publish_report
type: io.kestra.plugin.core.namespace.UploadFiles
description: Keep the measurement next to the batch it cleared, so an auditor
can check both.
namespace: "{{ flow.namespace }}"
filesMap:
"pii-gate/passed/{{ execution.id }}/scrubbed.csv": "{{ outputs.scrub_and_rescan.outputFiles['scrubbed.csv'] }}"
"pii-gate/passed/{{ execution.id }}/redaction-report.json": "{{ outputs.scrub_and_rescan.outputFiles['redaction-report.json'] }}"
- id: log_published
type: io.kestra.plugin.core.log.Log
message: >-
Load allowed. {{ outputs.verify_load.row.loaded_rows }} rows read back
from {{ inputs.warehouse_path }} carrying {{
outputs.verify_load.row.spans_removed }} removed spans. Residual scan
found {{ outputs.scrub_and_rescan.vars.residual_entities }} findings
against a threshold of {{ inputs.max_residual_entities }}. Report
published under pii-gate/passed/{{ execution.id }}/.
else:
- id: quarantine_batch
type: io.kestra.plugin.core.namespace.UploadFiles
description: Park the scrubbed batch and the residual findings where a reviewer
can work on them. Nothing reaches the warehouse.
namespace: "{{ flow.namespace }}"
filesMap:
"pii-gate/blocked/{{ execution.id }}/scrubbed.csv": "{{ outputs.scrub_and_rescan.outputFiles['scrubbed.csv'] }}"
"pii-gate/blocked/{{ execution.id }}/redaction-report.json": "{{ outputs.scrub_and_rescan.outputFiles['redaction-report.json'] }}"
"pii-gate/blocked/{{ execution.id }}/residual-findings.json":
"{{ outputs.scrub_and_rescan.outputFiles['residual-findings.json']
}}"
- id: block_load
type: io.kestra.plugin.core.execution.Fail
description: End the run as FAILED. The warehouse is untouched.
errorMessage: >-
Residual PII survived the scrub: {{
outputs.scrub_and_rescan.vars.residual_entities }} findings against a
threshold of {{ inputs.max_residual_entities }}. Surviving types: {{
outputs.scrub_and_rescan.vars.residual_types }}. Presidio never
flagged {{ outputs.scrub_and_rescan.vars.missed_by_presidio }} in the
first pass, which is why the independent scan exists. The batch is
quarantined under pii-gate/blocked/{{ execution.id }}/ and nothing was
loaded into {{ inputs.warehouse_path }}.