From 2f3b1bbe17bfca18f1a99082c0345316ad1fcc67 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Thu, 27 Aug 2026 16:43:51 +0530 Subject: [PATCH 01/23] refacotr(schema): implement discussion date in disscussion table --- app/database/operations.py | 30 ++++++++++++++++++++++-------- schema.sql | 1 + 2 files changed, 23 insertions(+), 8 deletions(-) diff --git a/app/database/operations.py b/app/database/operations.py index 607f8be..1bfbdef 100644 --- a/app/database/operations.py +++ b/app/database/operations.py @@ -225,10 +225,14 @@ async def insert_or_update_submission( # Upsert parent metadata tables (reads from flattened tags) program_id, leader_id = await upsert_metadata(conn, tags, tenant_code) - # Parse submission date. On a partial update where submissionDate is absent, - # leave it None (rather than defaulting to now()) so COALESCE below preserves - # the existing value instead of overwriting it with a "real" default. - sub_date_str = data.get("submissionDate") + # Parse submission date (report created at). + # For discussion submissions, report created at comes from eventPublishedAt. + normalized_sub_type = submission_type.lower().strip() + if "discussion" in normalized_sub_type: + sub_date_str = event_payload.get("eventPublishedAt") + else: + sub_date_str = data.get("submissionDate") + if sub_date_str: submission_date = datetime.fromisoformat(sub_date_str.replace("Z", "+00:00")) elif event_type != "update": @@ -374,6 +378,13 @@ async def insert_or_update_submission( image_urls = _normalize_media_url_list(data.get("imageUrls")) pdf_urls, masked_pdf_urls = _normalize_pdf_urls(data.get("pdfUrls")) + # Parse discussion date (date the discussion took place) from submissionDate. + disc_date_str = data.get("submissionDate") + discussion_date = ( + datetime.fromisoformat(disc_date_str.replace("Z", "+00:00")) + if disc_date_str else None + ) + if row_exists: await conn.execute( """ @@ -387,6 +398,7 @@ async def insert_or_update_submission( pdf_urls = COALESCE($9, pdf_urls), masked_pdf_urls = COALESCE($10, masked_pdf_urls), transcript_link = COALESCE($11, transcript_link), + discussion_date = COALESCE($12, discussion_date), updated_at = now() WHERE submission_id = $1 AND tenant_code = $2 """, @@ -399,16 +411,17 @@ async def insert_or_update_submission( image_urls, pdf_urls, masked_pdf_urls, - data.get("transcriptLink") + data.get("transcriptLink"), + discussion_date ) else: await conn.execute( """ INSERT INTO discussion_submissions ( submission_id, tenant_code, title, challenges, solutions, - author, language, image_urls, pdf_urls, masked_pdf_urls, transcript_link + author, language, image_urls, pdf_urls, masked_pdf_urls, transcript_link, discussion_date ) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11) + VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12) """, submission_id, tenant_code, data.get("title"), @@ -419,7 +432,8 @@ async def insert_or_update_submission( image_urls, pdf_urls, masked_pdf_urls, - data.get("transcriptLink") + data.get("transcriptLink"), + discussion_date ) # Dynamic KPI metrics: participantsData is a full snapshot when present; diff --git a/schema.sql b/schema.sql index 704eafd..2c0fae0 100644 --- a/schema.sql +++ b/schema.sql @@ -100,6 +100,7 @@ CREATE TABLE discussion_submissions ( submission_id TEXT NOT NULL, tenant_code TEXT NOT NULL, title TEXT, + discussion_date TIMESTAMPTZ, challenges TEXT[], -- one array element per discrete statement (see operations.py's _normalize_statement_list) solutions TEXT[], -- same format as challenges author TEXT, From 52850cacac83d0a9f48fb3d3b0dd4ca3efad3cd3 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Thu, 27 Aug 2026 16:44:18 +0530 Subject: [PATCH 02/23] refactor(prompt): add v2 for PII detection --- seed_prompts.sql | 129 +++++++++++++++++++++++++++++++++++++++++++++-- 1 file changed, 126 insertions(+), 3 deletions(-) diff --git a/seed_prompts.sql b/seed_prompts.sql index e3a13f4..4a4eeb7 100644 --- a/seed_prompts.sql +++ b/seed_prompts.sql @@ -261,7 +261,8 @@ WHERE p.name = 'Story Rating' ON CONFLICT (prompt_id, version) DO UPDATE SET system_prompt = EXCLUDED.system_prompt, user_prompt = EXCLUDED.user_prompt; --- 2c. PII Detection and Abusive-Language prompt (from prompt_version.csv row 3) +-- 2c. PII Detection and Abusive-Language prompt v1 (from prompt_version.csv row 3) +-- Flags abusive language only — does not mask it. Superseded by v2 below. INSERT INTO prompt_version (prompt_id, version, system_prompt, user_prompt, is_active, change_note, created_at) SELECT p.id, @@ -341,10 +342,132 @@ statements. Return every column in {columns}, even if empty (empty arrays, false, no missing keys).', -- user_prompt: text placeholder E'Analyse the following text: +{{text}}', + FALSE, + 'Seeded PII detection prompt v1 — system/user split from prompt_version.csv. Deactivated in favor of v2 (severity-tiered abuse masking).', + now() +FROM prompts p +WHERE p.name = 'PII and Abusive-Language Detection' +ON CONFLICT (prompt_id, version) DO UPDATE SET system_prompt = EXCLUDED.system_prompt, user_prompt = EXCLUDED.user_prompt, is_active = EXCLUDED.is_active, change_note = EXCLUDED.change_note; + + +-- 2d. PII and Abusive-Language Detection prompt v2 +-- Incorporates product team PII clarifications (context-aware masking) +-- and severity-tiered abuse masking (mild/moderate/severe). +-- This version is the active one; v1 above is deactivated (reference only). +INSERT INTO prompt_version (prompt_id, version, system_prompt, user_prompt, is_active, change_note, created_at) +SELECT + p.id, + 2, + E'# PII & Abusive-Language Detection Prompt + +## Overview + +You are a PII and abusive-language detector for community field-worker story submissions (India context, multilingual). + +INPUT FIELDS TO SCAN: {columns} + +--- + +## CRITICAL RULE: REPLACEMENT ONLY (NO WRAPPING) + +**You must completely REPLACE the PII or abusive text with the tag. You must NEVER wrap the text with opening and closing tags.** + +- **CORRECT (Replacement)**: "The driver was abused as . His was demanded." +- **WRONG (Wrapping)**: "The driver Sunil was abused as donkey. His Aadhaar 1234 was demanded." + +**NEVER output closing tags (like , , , , , ). The tag is a single-use placeholder (e.g. ).** + +--- + +## PII Masking Rules + +### CORE PRINCIPLE: +**You must mask any detail or combination of details that could identify or reveal a specific individual.** A person''s name alone is direct PII. However, indirect details (like a specific school name, village, or small locality) also become PII and MUST be masked if they are combined with a specific person or specific incident because they could allow someone to identify the person. + +### MUST Mask (replace with tags in masked_text): +- **Person names**: Any named individual (students, teachers, parents, community members, officials). + Tags: +- **Phone numbers, Aadhaar numbers, ID numbers, email addresses**: Any specific identifiable number or contact. + Tags: , +- **Specific small locations/institutions tied to an identifiable person or incident**: Village names, specific school names (like "DPS School"), or ward names — ONLY when they can reveal or identify a specific person. + Tags: + Example: "To ensure safety of a girl subjected to domestic violence in Dumra village" → mask "Dumra village" as because a village is a tiny unit and the girl could be identified. + Example: "girl in X village in Rohtas district" → mask "X village" as , but keep "Rohtas district" (district is too broad to be identifying). + +### Replacement Rules & Example: +- **Rule**: Replace the PII word/phrase entirely with the tag — do NOT wrap the word with tags. +- Input: "Kunal, a teacher at DPS School..." +- Correct masked_text: ", a teacher at ..." +- WRONG: "Kunal, a teacher at DPS School..." + +### Do NOT Mask (keep as-is): +- **District names**: e.g. "Rohtas", "Patna", "Muzaffarpur" — districts are too broad to identify anyone. + Example: "To improve girls'' education in Rohtas District" → keep "Rohtas District" as-is. +- **State names**: e.g. "Bihar", "Jharkhand", "Uttar Pradesh" — never PII. +- **Program/scheme names**: e.g. "Sachethan", "Samagra Shiksha", "Poshan Abhiyaan" — these are government programs, not PII. + Example: "To improve learning through the Sachethan program" → keep as-is. +- **Generic category/institution terms**: Anganwadi, Kasturba Vidyalaya (a type of school in Bihar, not a specific school), Panchayat, Gram Sabha — these are categories. + Example: "get her enrolled in Kasturba Vidyalaya" → keep as-is. Only the person''s name is PII. +- **Generic common words**: "Aadhar" used as a word (not a number), "morning", "evening", "summer", "winter" — not PII. +- **Age groups or grade levels** without identifying details: e.g. "Class 5 students", "children aged 6-14". + +### Decision Rule: +Ask: "Could this detail or combination of details reveal a specific individual?" +- YES → mask the identifying parts. +- NO → keep as-is. + +--- + +## Abusive-Language Detection & Masking + +Detect and **mask** profanity, insults, hate speech, threats, slurs, and harassment in the masked_text. **Replace** the offending word/phrase entirely with the severity tag — do NOT wrap the word with tags. + +### Masking Example: +- Input: "That donkey teacher never comes to class" +- Correct masked_text: "That teacher never comes to class" +- WRONG: "That donkey teacher never comes to class" + +### Severity Tiers: + +| Severity | Description | Replacement Tag | +|------------|-----------------------------------------------------------------------------|-----------------| +| mild | Casual profanity or venting with no clear personal target ("stupid system", "garbage") | | +| moderate | Direct insult aimed at a named or identifiable person ("useless middleman", "corrupt teacher") | | +| severe | Slurs, threats, harassment (gender/caste/religion-based), any abuse involving a minor | | + +### Abuse Rules: +- Replace only the offending word or phrase with the tag — keep the rest of the sentence intact. +- If a person''s name appears inside an abusive sentence, apply BOTH a PII tag on the name AND an abuse tag on the abusive word — they are independent. +- Institution-directed abuse without a named target: mask by severity tier. +- For every detected abusive span, provide: the text, severity level, confidence score (0.0-1.0), and a short reason (max 8 words). + +--- + +## Output Format + +Return ONLY a valid JSON object (strict JSON, no explanation outside JSON). One entry per column in {columns}. Never omit any key — use empty arrays and false for clean columns. + +{ + "": { + "masked_text": "...", + "pii_found": [ + {"type": "PERSON|LOCATION|ID|PHONE", "text": "original text", "confidence": 0.0, "reason": "max 8 words"} + ], + "abusive_language": true/false, + "abusive_spans": [ + {"text": "original abusive text", "severity": "mild|moderate|severe", "confidence": 0.0, "reason": "max 8 words"} + ] + } +} + +Return every column in {columns}, even if no PII or abuse is found (empty arrays, abusive_language: false, masked_text = original text unchanged).', + -- user_prompt: text placeholder + E'Analyse the following text: {{text}}', TRUE, - 'Seeded PII detection prompt v1 — system/user split from prompt_version.csv', + 'PII detection prompt v2 — product-team PII clarifications (context-aware masking, district/state/program names retained) and severity-tiered abuse masking (mild/moderate/severe).', now() FROM prompts p WHERE p.name = 'PII and Abusive-Language Detection' -ON CONFLICT (prompt_id, version) DO UPDATE SET system_prompt = EXCLUDED.system_prompt, user_prompt = EXCLUDED.user_prompt; +ON CONFLICT (prompt_id, version) DO UPDATE SET system_prompt = EXCLUDED.system_prompt, user_prompt = EXCLUDED.user_prompt, is_active = EXCLUDED.is_active, change_note = EXCLUDED.change_note; \ No newline at end of file From c07b9fcb26782b91e29efd55f0b45f47a580bbd5 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Fri, 28 Aug 2026 15:28:13 +0530 Subject: [PATCH 03/23] refactor: replace Temporal CSV processing with direct Kafka production and background tasks --- .env.example | 1 - app/api/routes/uploads.py | 11 +- app/api/services/uploads.py | 301 ++++++++++++++++++++--- app/config.py | 1 - app/services/ingestion_validation.py | 2 +- app/temporal/csv_processing_activity.py | 305 ------------------------ app/temporal/worker.py | 40 +--- app/temporal/workflows.py | 109 --------- 8 files changed, 278 insertions(+), 492 deletions(-) delete mode 100644 app/temporal/csv_processing_activity.py diff --git a/.env.example b/.env.example index 33c2353..dbaf4c3 100644 --- a/.env.example +++ b/.env.example @@ -71,7 +71,6 @@ AUTH_TOKEN=your-secret-bearer-token-here # CSV Upload / Processing Configuration MAX_CSV_UPLOAD_BYTES=10485760 CSV_BLOB_UPLOADS=mitra_dashboard_api_output -CSV_SCHEDULE_CRON_TIME=40 15 * * * STORY_CSV_COLUMN=["id","Title","User name","Designation","Location","District","Organization","Report Created At","Objective","Challenges","Action Steps","Impact","Duration","Blurb","masked_blurb","Content","masked_content","Images","Pdf","Transcript Link","Session ID"] DISCUSSION_CSV_COLUMN=["id","Title","User name","User Location","District","Participant Count","Men","Women","Children","Date of Discussion","Organization","Challenges","Solutions","Author","Language","Report Created At","Transcript Link","Image Urls","PDF Urls","Session ID"] DISCUSSION_PARTICIPANTS_MAP={"men": "Men", "women": "Women", "children": "Children", "teacher": "Teacher", "participant count": "Participant Count"} diff --git a/app/api/routes/uploads.py b/app/api/routes/uploads.py index d2604f1..b085ffc 100644 --- a/app/api/routes/uploads.py +++ b/app/api/routes/uploads.py @@ -1,5 +1,5 @@ import logging -from fastapi import APIRouter, Depends, Form, UploadFile, File +from fastapi import APIRouter, Depends, Form, UploadFile, File, BackgroundTasks from fastapi.security import HTTPAuthorizationCredentials from app.api.deps import verify_auth_token @@ -19,6 +19,7 @@ @uploads_router.post("/upload/", response_model=UploadResponse) async def upload_report( + background_tasks: BackgroundTasks, report_type: str = Form(...), program_name: str = Form(...), leader_category: str = Form(...), @@ -30,7 +31,8 @@ async def upload_report( Upload a CSV report file. The file is stored in GCS and a tracking record is created with - status='pending' (validation passed) or status='on_hold' (validation failed). + status='pending'. Processing (column validation, Kafka publishing) + runs in the background. """ # Pure request-shape checks report_type = validate_report_type(report_type) @@ -47,15 +49,18 @@ async def upload_report( tenant_code=tenant_code, file_name=file.filename, file_bytes=file_bytes, + background_tasks=background_tasks, ) @uploads_router.post("/process/csv/{record_id}") async def push_record( record_id: int, + background_tasks: BackgroundTasks, _token: HTTPAuthorizationCredentials = Depends(verify_auth_token), ): """ Manually trigger processing for a specific csv_upload record in pending status. """ - return await upload_service.handle_push(record_id) + return await upload_service.handle_push(record_id, background_tasks) + diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 3828087..dcbdc78 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -3,14 +3,19 @@ import io import json import logging +import threading +import uuid from datetime import datetime import pandas as pd -from temporalio.client import Client +from confluent_kafka import Producer, KafkaException +from fastapi import BackgroundTasks from app.config import settings from app.api.validators.uploads import validate_columns -from app.services.gcp_storage import upload_csv +from app.services.gcp_storage import upload_csv, fetch_csv from app.database import operations +from app.database.db import db +from app.services.ingestion_validation import validate_ingestion_schema from app.api.exceptions import ( DuplicateFile, InvalidCsvColumns, @@ -365,6 +370,257 @@ def rows_to_json( yield row_to_json(row, report_type, event_type, metadata) +# --------------------------------------------------------------------------- +# Kafka Producer (singleton, thread-safe) +# --------------------------------------------------------------------------- + +_producer: Optional[Producer] = None +_producer_lock = threading.Lock() + + +def _get_producer() -> Producer: + """Return a singleton confluent-kafka Producer, creating it on first call.""" + global _producer + if _producer is not None: + return _producer + with _producer_lock: + if _producer is None: + _producer = Producer({ + "bootstrap.servers": settings.KAFKA_BOOTSTRAP_SERVERS, + "acks": "all", + "enable.idempotence": True, + }) + return _producer + + +def _push_rows_sync(payloads: List[Any]) -> None: + """ + Runs in a worker thread (via asyncio.to_thread) — produce()/flush() are + blocking calls. Flushes once for the whole batch rather than per row. + """ + producer = _get_producer() + delivery_error = {} + + def _on_delivery(err, _msg): + if err is not None: + delivery_error["error"] = err + + for payload, key in payloads: + producer.produce( + settings.KAFKA_TOPIC_INGESTION, + value=payload.encode("utf-8"), + key=key.encode("utf-8") if key else None, + callback=_on_delivery, + ) + producer.poll(0) + if "error" in delivery_error: + raise KafkaException(delivery_error["error"]) + + remaining = producer.flush(10) + if remaining > 0: + raise TimeoutError(f"Timed out waiting for Kafka delivery ({remaining} still in-flight)") + if "error" in delivery_error: + raise KafkaException(delivery_error["error"]) + + +# --------------------------------------------------------------------------- +# Inline CSV Processing (replaces Temporal activities) +# --------------------------------------------------------------------------- + +async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) -> None: + """ + Processes a single csv_upload record end-to-end: + 1. Use in-memory CSV bytes (or fetch from cloud storage if file_bytes is None) + 2. Validate columns + 3. Look up program/leader metadata from DB + 4. Build Kafka payloads, schema-validate each row + 5. Publish valid rows to Kafka + 6. Update DB status to 'success' (or 'on_hold' on failure) + + Runs as a FastAPI BackgroundTask — any exception is caught, logged, and + recorded in the csv_uploads row so the caller's 200 response is unaffected. + """ + record = await operations.get_record(record_id) + if not record: + logger.error("process_csv_inline: record %s not found", record_id) + return + + cloud_storage_path = record["cloud_storage_path"] + report_type = record["report_type"] + + # --- 1. Fetch/Parse CSV (use in-memory file_bytes if available, else fetch from storage) --- + try: + if file_bytes is None: + csv_file = await asyncio.to_thread(fetch_csv, cloud_storage_path) + else: + csv_file = file_bytes + df = await asyncio.to_thread(load_csv, csv_file) + except Exception as exc: + logger.exception("Failed to fetch/load CSV for record %s", record_id) + error_meta = { + "stage": "CSV Fetching", + "error": "Failed to fetch/load CSV", + "exception": str(exc), + "timestamp": datetime.utcnow().isoformat() + "Z", + } + await operations.update_status(record_id, "on_hold", error_meta) + return + + # --- 2. Validate columns --- + is_valid, errors = await asyncio.to_thread(validate_columns, df, report_type) + if not is_valid: + logger.warning("Validation failed for record %s: %s", record_id, errors) + error_meta = { + "stage": "CSV Column Validation", + "error": "Invalid CSV schema", + "validation_errors": errors, + "timestamp": datetime.utcnow().isoformat() + "Z", + } + await operations.update_status(record_id, "on_hold", error_meta) + return + + await operations.update_status(record_id, "in_progress") + + # --- 3. Look up program / leader category metadata from Postgres --- + program_info = None + leader_info = None + record_meta = record.get("meta_data") or {} + if isinstance(record_meta, str): + try: + record_meta = json.loads(record_meta) + except json.JSONDecodeError: + record_meta = {} + if not isinstance(record_meta, dict): + record_meta = {} + + tenant_code = record_meta.get("tenant_code") or "mitra" + + try: + async with db.pool.acquire() as conn: + leader_row = await conn.fetchrow( + "SELECT id, name, description, tenant_code FROM leader_category WHERE name = $1 LIMIT 1", + record.get("leader_category"), + ) + if leader_row: + leader_info = { + "id": str(leader_row["id"]), + "name": leader_row["name"], + "description": leader_row["description"], + } + tenant_code = leader_row["tenant_code"] + + if leader_row: + program_row = await conn.fetchrow( + "SELECT id, name, description, tenant_code, leaders_id FROM programs WHERE name = $1 AND leaders_id = $2 LIMIT 1", + record.get("program_name"), + leader_row["id"], + ) + else: + program_row = await conn.fetchrow( + "SELECT id, name, description, tenant_code, leaders_id FROM programs WHERE name = $1 LIMIT 1", + record.get("program_name"), + ) + + if program_row: + program_info = { + "id": str(program_row["id"]), + "name": program_row["name"], + "description": program_row["description"], + } + tenant_code = program_row.get("tenant_code", tenant_code) + + if program_row and not leader_info: + leader_row_from_program = await conn.fetchrow( + "SELECT id, name, description, tenant_code FROM leader_category WHERE id = $1 LIMIT 1", + program_row["leaders_id"], + ) + if leader_row_from_program: + leader_info = { + "id": str(leader_row_from_program["id"]), + "name": leader_row_from_program["name"], + "description": leader_row_from_program["description"], + } + tenant_code = leader_row_from_program.get("tenant_code", tenant_code) + except Exception as db_exc: + logger.warning("Failed to query program/leader category metadata from DB: %s", db_exc) + + # Fallbacks if DB query returned nothing + if not leader_info: + leader_info = { + "id": str(uuid.uuid4()), + "name": record.get("leader_category") or "District Leader", + "description": f"Leader category: {record.get('leader_category') or 'District Leader'}", + } + if not program_info: + program_info = { + "id": str(uuid.uuid4()), + "name": record.get("program_name") or "My Program", + "description": f"Program: {record.get('program_name') or 'My Program'}", + } + + metadata = { + "programInfo": program_info, + "LeaderCategoryInfo": leader_info, + "tenantCode": tenant_code, + } + + # --- 4. Build Kafka payloads and schema-validate each row --- + chunks = split_csv(df) + payloads = [] + schema_errors = [] + row_number = 0 + + for chunk in chunks: + for payload_str in rows_to_json(chunk, report_type, metadata=metadata): + row_number += 1 + try: + payload_dict = json.loads(payload_str) + except json.JSONDecodeError as exc: + schema_errors.append({"row": row_number, "problems": [f"Failed to parse generated payload: {exc}"]}) + continue + + problems = validate_ingestion_schema(payload_dict, report_type, "create") + if problems: + schema_errors.append({ + "row": row_number, + "submissionId": payload_dict.get("submissionId"), + "sessionId": payload_dict.get("sessionId"), + "problems": problems, + }) + continue + + payloads.append((payload_str, f"{record_id}-{len(payloads)}")) + + if schema_errors: + logger.warning( + "record %s: %d of %d row(s) failed pre-publish schema validation and were skipped: %s", + record_id, len(schema_errors), row_number, schema_errors, + ) + + # --- 5. Publish to Kafka --- + if payloads: + try: + await asyncio.to_thread(_push_rows_sync, payloads) + except Exception as exc: + logger.exception("Kafka push failed for record %s", record_id) + error_meta = { + "stage": "Kafka Publishing", + "error": "Failed to publish record", + "exception": str(exc), + "timestamp": datetime.utcnow().isoformat() + "Z", + } + await operations.update_status(record_id, "on_hold", error_meta) + return + + # --- 6. Update status to success --- + final_meta = {"rows_pushed": len(payloads), "processed_at": datetime.utcnow().isoformat() + "Z"} + if schema_errors: + final_meta["schema_validation_errors"] = schema_errors + + await operations.update_status(record_id, "success", final_meta) + logger.info("CSV record %s processed successfully: %d rows pushed to Kafka", record_id, len(payloads)) + + # --------------------------------------------------------------------------- # Service Orchestration Logic # --------------------------------------------------------------------------- @@ -376,9 +632,8 @@ async def handle_upload( tenant_code: str, file_name: str, file_bytes: bytes, + background_tasks: BackgroundTasks, ) -> dict: - from app.temporal.workflows import CsvProcessingWorkflow - normalized_type = report_type.lower().strip() file_size = len(file_bytes) @@ -435,21 +690,9 @@ async def handle_upload( record_id, normalized_type, cloud_storage_path, ) - # Trigger Temporal workflow in real-time mode - if settings.PROCESSING_MODE.lower().strip() == "real-time": - try: - temporal_client = await Client.connect(settings.TEMPORAL_HOST, namespace=settings.TEMPORAL_NAMESPACE) - await temporal_client.start_workflow( - CsvProcessingWorkflow.run, - record_id, - id=f"csv-upload-{record_id}", - task_queue=settings.TEMPORAL_QUEUE, - ) - logger.info("Triggered real-time CsvProcessingWorkflow for upload ID %s", record_id) - except Exception as e: - logger.error("Failed to trigger real-time CsvProcessingWorkflow: %s", e) - await operations.update_status(record_id, "on_hold", {"error": f"Temporal trigger failed: {e}"}) - raise RuntimeError(f"Failed to start CSV processing workflow: {e}") + # Schedule inline processing as a background task (pass file_bytes to avoid extra GCS download) + background_tasks.add_task(process_csv_inline, record_id, file_bytes) + logger.info("Scheduled inline CSV processing for upload ID %s", record_id) return { "message": "Successfully uploaded to cloud", @@ -458,9 +701,7 @@ async def handle_upload( } -async def handle_push(record_id: int) -> dict: - from app.temporal.workflows import CsvProcessingWorkflow - +async def handle_push(record_id: int, background_tasks: BackgroundTasks) -> dict: record = await operations.get_record(record_id) if not record: raise RecordNotFound("Record not found") @@ -477,15 +718,7 @@ async def handle_push(record_id: int) -> dict: if claim_status == "in_progress": raise RecordAlreadyProcessing("Record is already being processed") - try: - temporal_client = await Client.connect(settings.TEMPORAL_HOST, namespace=settings.TEMPORAL_NAMESPACE) - await temporal_client.start_workflow( - CsvProcessingWorkflow.run, - record_id, - id=f"csv-upload-{record_id}", - task_queue=settings.TEMPORAL_QUEUE, - ) - return {"status": "success", "message": "CSV processing workflow started"} - except Exception as e: - await operations.update_status(record_id, "on_hold", {"error": str(e)}) - raise RuntimeError(f"Failed to start CSV processing workflow: {e}") + # Schedule inline processing as a background task (no Temporal) + background_tasks.add_task(process_csv_inline, record_id) + return {"status": "success", "message": "CSV processing started"} + diff --git a/app/config.py b/app/config.py index 1d65a16..46f4c4a 100644 --- a/app/config.py +++ b/app/config.py @@ -77,7 +77,6 @@ class Settings(BaseSettings): # CSV Upload / Processing Configuration MAX_CSV_UPLOAD_BYTES: int = Field(default=10485760) # 10MB CSV_BLOB_UPLOADS: str = Field(default="mitra_dashboard_api_output") - CSV_SCHEDULE_CRON_TIME: str = Field(default="40 15 * * *") # Expected CSV column headers per report type (JSON arrays of column names, # matched case-insensitively against the uploaded file's header row). STORY_CSV_COLUMN: str = Field( diff --git a/app/services/ingestion_validation.py b/app/services/ingestion_validation.py index 25d2442..ed5f1c8 100644 --- a/app/services/ingestion_validation.py +++ b/app/services/ingestion_validation.py @@ -34,7 +34,7 @@ def validate_ingestion_schema(event: dict, submission_type: str, event_type: str an empty list means the event is valid. Shared between app/kafka/consumer.py (validates events arriving off the Kafka - topic) and app/temporal/csv_processing_activity.py (validates each CSV-derived + topic) and app/api/services/uploads.py (validates each CSV-derived event against the same schema before it's ever published to Kafka). """ normalized_type = submission_type.lower().strip() if isinstance(submission_type, str) else "" diff --git a/app/temporal/csv_processing_activity.py b/app/temporal/csv_processing_activity.py deleted file mode 100644 index 5a7b997..0000000 --- a/app/temporal/csv_processing_activity.py +++ /dev/null @@ -1,305 +0,0 @@ -import asyncio -import json -import logging -import threading -import uuid -from datetime import datetime -from typing import Any, Dict, List, Optional -from temporalio import activity -from confluent_kafka import Producer, KafkaException - -from app.config import settings -from app.database.db import db -from app.database import operations as csv_upload_repo -from app.api.services.uploads import load_csv, rows_to_json, split_csv -from app.api.validators.uploads import validate_columns -from app.services.gcp_storage import fetch_csv -from app.services.ingestion_validation import validate_ingestion_schema - -logger = logging.getLogger("analytics_service.temporal.csv_processing_activity") - -_producer: Optional[Producer] = None -# Real OS-thread lock, not asyncio.Lock — _get_producer() is called from -# separate threads via asyncio.to_thread (_push_rows_sync), not concurrently -# within one event loop. Without this, concurrent CsvProcessingWorkflow child -# workflows could each construct their own Producer at once (see the identical -# fix + reasoning in app/services/classifier.py's _get_model()). -_producer_lock = threading.Lock() - - -def _get_producer() -> Producer: - global _producer - if _producer is not None: - return _producer - with _producer_lock: - if _producer is None: - _producer = Producer({ - "bootstrap.servers": settings.KAFKA_BOOTSTRAP_SERVERS, - "acks": "all", - "enable.idempotence": True, - }) - return _producer - - -def _push_rows_sync(payloads: List[Any]) -> None: - """ - Runs in a worker thread (via asyncio.to_thread) — produce()/flush() are - blocking calls. Flushes once for the whole batch rather than per row (a - per-row flush forces a network round trip per row, far too slow for large - CSVs), mirroring app/kafka/consumer.py's DLQ producer pattern. - """ - producer = _get_producer() - delivery_error = {} - - def _on_delivery(err, _msg): - if err is not None: - delivery_error["error"] = err - - for payload, key in payloads: - producer.produce( - settings.KAFKA_TOPIC_INGESTION, - value=payload.encode("utf-8"), - key=key.encode("utf-8") if key else None, - callback=_on_delivery, - ) - producer.poll(0) - if "error" in delivery_error: - raise KafkaException(delivery_error["error"]) - - remaining = producer.flush(10) - if remaining > 0: - raise TimeoutError(f"Timed out waiting for Kafka delivery ({remaining} still in-flight)") - if "error" in delivery_error: - raise KafkaException(delivery_error["error"]) - - -@activity.defn -async def csv_fetch_and_validate_activity(record_id: int) -> bool: - """ - Temporal activity to fetch the CSV file from cloud storage and update status. - """ - record = await csv_upload_repo.get_record(record_id) - if not record: - raise ValueError(f"Record {record_id} not found in database.") - - cloud_storage_path = record["cloud_storage_path"] - - try: - # Fetch from GCS + parse — both blocking, offloaded from the event loop. - csv_file = await asyncio.to_thread(fetch_csv, cloud_storage_path) - df = await asyncio.to_thread(load_csv, csv_file) - except Exception as exc: - logger.exception("Failed to fetch/load CSV for record %s", record_id) - error_meta = { - "stage": "CSV Fetching", - "error": "Failed to fetch/load CSV from GCS", - "exception": str(exc), - "timestamp": datetime.utcnow().isoformat() + "Z" - } - await csv_upload_repo.update_status(record_id, "on_hold", error_meta) - return False - - is_valid, errors = await asyncio.to_thread(validate_columns, df, record["report_type"]) - if not is_valid: - logger.warning("Validation failed for record %s: %s", record_id, errors) - error_meta = { - "stage": "CSV Column Validation", - "error": "Invalid CSV schema", - "validation_errors": errors, - "timestamp": datetime.utcnow().isoformat() + "Z" - } - await csv_upload_repo.update_status(record_id, "on_hold", error_meta) - return False - - await csv_upload_repo.update_status(record_id, "in_progress") - return True - - -@activity.defn -async def csv_push_to_kafka_activity(record_id: int) -> Dict[str, Any]: - """ - Temporal activity to process the CSV file row by row and publish - individual messages to Kafka. Returns {"rows_pushed": int, - "schema_validation_errors": list} — rows failing pre-publish schema - validation are skipped (not published) and reported here instead. - """ - record = await csv_upload_repo.get_record(record_id) - if not record: - raise ValueError(f"Record {record_id} not found in database.") - - report_type = record["report_type"] - cloud_storage_path = record["cloud_storage_path"] - - # Load CSV — blocking I/O + parse, offloaded from the event loop. - csv_file = await asyncio.to_thread(fetch_csv, cloud_storage_path) - df = await asyncio.to_thread(load_csv, csv_file) - - is_valid, errors = await asyncio.to_thread(validate_columns, df, report_type) - if not is_valid: - logger.warning("Kafka push blocked for record %s due to invalid columns: %s", record_id, errors) - error_meta = { - "stage": "CSV Column Validation", - "error": "Invalid CSV schema - aborting Kafka push", - "validation_errors": errors, - "timestamp": datetime.utcnow().isoformat() + "Z" - } - await csv_upload_repo.update_status(record_id, "on_hold", error_meta) - raise ValueError("CSV validation failed for Kafka push") - - # Fetch programs / leader categories info from Postgres once for context mapping - program_info = None - leader_info = None - # Use tenant_code from the upload payload (stored in meta_data) as primary source, - # falling back to DB lookup and then to "mitra" as last resort. - record_meta = record.get("meta_data") or {} - if isinstance(record_meta, str): - try: - record_meta = json.loads(record_meta) - except json.JSONDecodeError: - record_meta = {} - - if not isinstance(record_meta, dict): - record_meta = {} - - tenant_code = record_meta.get("tenant_code") or "mitra" - - try: - async with db.pool.acquire() as conn: - leader_row = await conn.fetchrow( - "SELECT id, name, description, tenant_code FROM leader_category WHERE name = $1 LIMIT 1", - record.get("leader_category") - ) - if leader_row: - leader_info = { - "id": str(leader_row["id"]), - "name": leader_row["name"], - "description": leader_row["description"], - } - tenant_code = leader_row["tenant_code"] - - if leader_row: - program_row = await conn.fetchrow( - "SELECT id, name, description, tenant_code, leaders_id FROM programs WHERE name = $1 AND leaders_id = $2 LIMIT 1", - record.get("program_name"), leader_row["id"] - ) - else: - program_row = await conn.fetchrow( - "SELECT id, name, description, tenant_code, leaders_id FROM programs WHERE name = $1 LIMIT 1", - record.get("program_name") - ) - - if program_row: - program_info = { - "id": str(program_row["id"]), - "name": program_row["name"], - "description": program_row["description"], - } - tenant_code = program_row.get("tenant_code", tenant_code) - - if program_row and not leader_info: - leader_row_from_program = await conn.fetchrow( - "SELECT id, name, description, tenant_code FROM leader_category WHERE id = $1 LIMIT 1", - program_row["leaders_id"] - ) - if leader_row_from_program: - leader_info = { - "id": str(leader_row_from_program["id"]), - "name": leader_row_from_program["name"], - "description": leader_row_from_program["description"], - } - tenant_code = leader_row_from_program.get("tenant_code", tenant_code) - except Exception as db_exc: - logger.warning("Failed to query program/leader category metadata from DB: %s", db_exc) - - # Fallbacks if DB query returned nothing - if not leader_info: - leader_info = { - "id": str(uuid.uuid4()), - "name": record.get("leader_category") or "District Leader", - "description": f"Leader category: {record.get('leader_category') or 'District Leader'}", - } - if not program_info: - program_info = { - "id": str(uuid.uuid4()), - "name": record.get("program_name") or "My Program", - "description": f"Program: {record.get('program_name') or 'My Program'}", - } - - metadata = { - "programInfo": program_info, - "LeaderCategoryInfo": leader_info, - "tenantCode": tenant_code, - } - - chunks = split_csv(df) - payloads = [] - schema_errors = [] - row_number = 0 - - for chunk in chunks: - for payload_str in rows_to_json(chunk, report_type, metadata=metadata): - row_number += 1 - try: - payload_dict = json.loads(payload_str) - except json.JSONDecodeError as exc: - schema_errors.append({"row": row_number, "problems": [f"Failed to parse generated payload: {exc}"]}) - continue - - # Double-check the generated event against the exact same schema - # app/kafka/consumer.py enforces at ingestion — catches a row missing - # a required field (e.g. no Session ID, now that it's no longer - # auto-generated) here, before it's ever published, rather than - # relying on the consumer to silently DLQ it later. - problems = validate_ingestion_schema(payload_dict, report_type, "create") - if problems: - schema_errors.append({ - "row": row_number, - "submissionId": payload_dict.get("submissionId"), - "sessionId": payload_dict.get("sessionId"), - "problems": problems, - }) - continue - - payloads.append((payload_str, f"{record_id}-{len(payloads)}")) - - if schema_errors: - logger.warning( - "record %s: %d of %d row(s) failed pre-publish schema validation and were skipped: %s", - record_id, len(schema_errors), row_number, schema_errors, - ) - - if payloads: - try: - await asyncio.to_thread(_push_rows_sync, payloads) - except Exception as exc: - logger.exception("Kafka push failed for record %s", record_id) - error_meta = { - "stage": "Kafka Publishing", - "error": "Failed to publish record", - "exception": str(exc), - "timestamp": datetime.utcnow().isoformat() + "Z" - } - await csv_upload_repo.update_status(record_id, "on_hold", error_meta) - raise - - return {"rows_pushed": len(payloads), "schema_validation_errors": schema_errors} - - -@activity.defn -async def csv_update_status_activity(params: Dict[str, Any]) -> None: - """ - Temporal activity to update the overall processing status of a csv_upload in PostgreSQL. - """ - record_id = params["record_id"] - status = params["status"] - meta_data = params.get("meta_data") - await csv_upload_repo.update_status(record_id, status, meta_data) - - -@activity.defn -async def fetch_pending_csv_uploads_activity() -> List[int]: - """ - Temporal activity to fetch the IDs of all pending csv_upload records. - """ - records = await csv_upload_repo.list_by_status("pending") - return [r["id"] for r in records] diff --git a/app/temporal/worker.py b/app/temporal/worker.py index 989b230..69ce5f9 100644 --- a/app/temporal/worker.py +++ b/app/temporal/worker.py @@ -10,8 +10,6 @@ from app.temporal.workflows import ( ConfigDrivenProcessingWorkflow, BatchProcessingWorkflow, - CsvProcessingWorkflow, - CsvBatchProcessingWorkflow, ) from app.temporal.activities import ( update_status_activity, @@ -21,12 +19,6 @@ from app.temporal.pii_and_abusive_activity import pii_and_abusive_language_detection_activity from app.temporal.thematic_activity import thematic_classification_activity from app.temporal.story_rating_activity import story_rating_activity -from app.temporal.csv_processing_activity import ( - csv_fetch_and_validate_activity, - csv_push_to_kafka_activity, - csv_update_status_activity, - fetch_pending_csv_uploads_activity, -) logger = logging.getLogger("analytics_service.temporal.worker") @@ -75,8 +67,6 @@ async def start_worker(): workflows = [ ConfigDrivenProcessingWorkflow, BatchProcessingWorkflow, - CsvProcessingWorkflow, - CsvBatchProcessingWorkflow, ] activities = [ pii_and_abusive_language_detection_activity, @@ -85,10 +75,6 @@ async def start_worker(): story_rating_activity, update_status_activity, fetch_pending_submissions_activity, - csv_fetch_and_validate_activity, - csv_push_to_kafka_activity, - csv_update_status_activity, - fetch_pending_csv_uploads_activity, ] worker = Worker( @@ -108,29 +94,7 @@ async def start_worker(): ScheduleAlreadyRunningError, ) - # 1. Register CSV batch processing schedule - try: - logger.info(f"Registering CSV batch schedule '{settings.CSV_SCHEDULE_CRON_TIME}' in Temporal...") - await client.create_schedule( - id="csv-batch-processing", - schedule=Schedule( - action=ScheduleActionStartWorkflow( - CsvBatchProcessingWorkflow.run, - id="csv-batch-processing-run", - task_queue=settings.TEMPORAL_QUEUE, - ), - spec=ScheduleSpec( - cron_expressions=[settings.CSV_SCHEDULE_CRON_TIME] - ), - ), - ) - logger.info("CSV batch schedule successfully registered.") - except ScheduleAlreadyRunningError: - logger.info("CSV batch schedule already exists in Temporal. Skipping registration.") - except Exception as e: - logger.error(f"Failed to register CSV batch schedule in Temporal: {e}") - - # 2. Register daily analysis batch processing schedule + # Register daily analysis batch processing schedule try: logger.info(f"Registering daily analysis batch schedule '{settings.BATCH_SCHEDULE_CRON}' in Temporal...") await client.create_schedule( @@ -156,7 +120,7 @@ async def start_worker(): # Real-time mode: clean up any leftover batch schedules from Temporal Server # (prevents a schedule left behind from a prior batch-mode config from # silently retrying forever with outdated arguments). - for sched_id in ("csv-batch-processing", "daily-batch-processing"): + for sched_id in ("daily-batch-processing",): try: handle = client.get_schedule_handle(sched_id) await handle.delete() diff --git a/app/temporal/workflows.py b/app/temporal/workflows.py index 1514a03..dfd3476 100644 --- a/app/temporal/workflows.py +++ b/app/temporal/workflows.py @@ -13,12 +13,6 @@ from app.temporal.pii_and_abusive_activity import pii_and_abusive_language_detection_activity from app.temporal.thematic_activity import thematic_classification_activity from app.temporal.story_rating_activity import story_rating_activity - from app.temporal.csv_processing_activity import ( - csv_fetch_and_validate_activity, - csv_push_to_kafka_activity, - csv_update_status_activity, - fetch_pending_csv_uploads_activity - ) @workflow.defn class ConfigDrivenProcessingWorkflow: @@ -311,106 +305,3 @@ async def run(self, batch_size: int, carry_over: Optional[Dict[str, Any]] = None } -@workflow.defn -class CsvProcessingWorkflow: - @workflow.run - async def run(self, record_id: int) -> Dict[str, Any]: - """ - Orchestrates processing for a single csv_upload record. - 1. Fetch CSV and validate columns. - 2. Publish each row to Kafka. - 3. Mark as success. - """ - retry_policy = RetryPolicy( - maximum_attempts=3, - initial_interval=timedelta(seconds=2), - backoff_coefficient=2.0 - ) - - # 1. Fetch and validate columns - is_valid = await workflow.execute_activity( - csv_fetch_and_validate_activity, - record_id, - start_to_close_timeout=timedelta(minutes=5), - retry_policy=retry_policy - ) - - if not is_valid: - return {"status": "on_hold", "reason": "Validation or fetching failed"} - - # 2. Push rows to Kafka - # Note: If any error happens during push, the activity itself catches it, - # sets status to 'on_hold' with error info, and raises an exception. - push_result = await workflow.execute_activity( - csv_push_to_kafka_activity, - record_id, - start_to_close_timeout=timedelta(minutes=30), - retry_policy=RetryPolicy(maximum_attempts=1) # No auto-retries for Kafka pushes to prevent duplicate writes - ) - rows_pushed = push_result["rows_pushed"] - schema_validation_errors = push_result.get("schema_validation_errors") or [] - - # 3. Update status to 'success'. Rows that failed pre-publish schema - # validation (e.g. a row missing a required field) are recorded here - # rather than silently dropped, even though the upload as a whole succeeded. - final_meta = {"rows_pushed": rows_pushed, "processed_at": workflow.now().isoformat()} - if schema_validation_errors: - final_meta["schema_validation_errors"] = schema_validation_errors - - await workflow.execute_activity( - csv_update_status_activity, - { - "record_id": record_id, - "status": "success", - "meta_data": final_meta - }, - start_to_close_timeout=timedelta(seconds=10), - retry_policy=retry_policy - ) - - return {"status": "success", "rows_pushed": rows_pushed, "schema_validation_errors": schema_validation_errors} - - -@workflow.defn -class CsvBatchProcessingWorkflow: - @workflow.run - async def run(self) -> Dict[str, Any]: - """ - Runs batch execution for all pending CSV uploads. - Retrieves pending records and executes child workflows in parallel. - """ - retry_policy = RetryPolicy( - maximum_attempts=2, - initial_interval=timedelta(seconds=2) - ) - - # Fetch pending csv_upload IDs - pending_ids: List[int] = await workflow.execute_activity( - fetch_pending_csv_uploads_activity, - start_to_close_timeout=timedelta(minutes=2), - retry_policy=retry_policy - ) - - if not pending_ids: - return {"processed_count": 0, "message": "No pending CSV uploads found."} - - # Fan-out child workflows to process each CSV in parallel - child_tasks = [] - for pid in pending_ids: - child_tasks.append( - workflow.execute_child_workflow( - CsvProcessingWorkflow.run, - pid, - id=f"csv-batch-child-{pid}" - ) - ) - - results = await asyncio.gather(*child_tasks, return_exceptions=True) - success_count = sum(1 for r in results if not isinstance(r, Exception)) - failed_count = len(results) - success_count - - return { - "processed_count": len(pending_ids), - "success_count": success_count, - "failed_count": failed_count - } From 9b273237689cef5dcfc19369054f137481679d98 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Fri, 4 Sep 2026 10:32:20 +0530 Subject: [PATCH 04/23] fix(uploads): make Kafka delivery completion batch-scoped and preserve retryable pending status --- app/api/services/uploads.py | 53 ++++++++++++++++++++++++++++--------- 1 file changed, 40 insertions(+), 13 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index dcbdc78..970d585 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -4,6 +4,7 @@ import json import logging import threading +import time import uuid from datetime import datetime import pandas as pd @@ -395,15 +396,27 @@ def _get_producer() -> Producer: def _push_rows_sync(payloads: List[Any]) -> None: """ - Runs in a worker thread (via asyncio.to_thread) — produce()/flush() are - blocking calls. Flushes once for the whole batch rather than per row. + Runs in a worker thread (via asyncio.to_thread) — produce()/poll() are + blocking calls. Tracks delivery callbacks specifically for this batch, + preventing process-wide flush() interference between concurrent CSV uploads. """ + if not payloads: + return + producer = _get_producer() - delivery_error = {} + pending_count = len(payloads) + delivery_error = None + done_event = threading.Event() + lock = threading.Lock() def _on_delivery(err, _msg): - if err is not None: - delivery_error["error"] = err + nonlocal pending_count, delivery_error + with lock: + if err is not None and delivery_error is None: + delivery_error = err + pending_count -= 1 + if pending_count <= 0: + done_event.set() for payload, key in payloads: producer.produce( @@ -413,14 +426,28 @@ def _on_delivery(err, _msg): callback=_on_delivery, ) producer.poll(0) - if "error" in delivery_error: - raise KafkaException(delivery_error["error"]) + with lock: + if delivery_error is not None: + raise KafkaException(delivery_error) + + # Poll network events until all delivery callbacks for THIS batch complete + timeout_seconds = 30.0 + start_time = time.time() + while not done_event.is_set(): + producer.poll(0.1) + with lock: + if delivery_error is not None: + raise KafkaException(delivery_error) + if time.time() - start_time > timeout_seconds: + break - remaining = producer.flush(10) - if remaining > 0: - raise TimeoutError(f"Timed out waiting for Kafka delivery ({remaining} still in-flight)") - if "error" in delivery_error: - raise KafkaException(delivery_error["error"]) + with lock: + if delivery_error is not None: + raise KafkaException(delivery_error) + if pending_count > 0: + raise TimeoutError( + f"Timed out waiting for batch Kafka delivery ({pending_count} of {len(payloads)} remaining)" + ) # --------------------------------------------------------------------------- @@ -609,7 +636,7 @@ async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) "exception": str(exc), "timestamp": datetime.utcnow().isoformat() + "Z", } - await operations.update_status(record_id, "on_hold", error_meta) + await operations.update_status(record_id, "pending", error_meta) return # --- 6. Update status to success --- From c395e51d813dd3cf8fe2640abb427a7485dc232f Mon Sep 17 00:00:00 2001 From: rohxnn Date: Fri, 4 Sep 2026 10:49:56 +0530 Subject: [PATCH 05/23] fix(uploads): atomically claim record in handle_upload to prevent duplicate processing race condition --- app/api/services/uploads.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 970d585..60c99e3 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -717,14 +717,15 @@ async def handle_upload( record_id, normalized_type, cloud_storage_path, ) - # Schedule inline processing as a background task (pass file_bytes to avoid extra GCS download) - background_tasks.add_task(process_csv_inline, record_id, file_bytes) - logger.info("Scheduled inline CSV processing for upload ID %s", record_id) + claim_status = await operations.try_claim_for_processing(record_id) + if claim_status == "success": + background_tasks.add_task(process_csv_inline, record_id, file_bytes) + logger.info("Scheduled inline CSV processing for upload ID %s", record_id) return { "message": "Successfully uploaded to cloud", "id": record_id, - "status": "pending", + "status": "in_progress" if claim_status == "success" else "pending", } From d672ee276402e5da7a674fe08c1df1059317c5af Mon Sep 17 00:00:00 2001 From: rohxnn Date: Fri, 4 Sep 2026 11:18:54 +0530 Subject: [PATCH 06/23] fix(prompts): update PII prompt v2 output format to specify array contract for multi-statement inputs --- seed_prompts.sql | 26 ++++++++++++++++---------- 1 file changed, 16 insertions(+), 10 deletions(-) diff --git a/seed_prompts.sql b/seed_prompts.sql index 4a4eeb7..c98523d 100644 --- a/seed_prompts.sql +++ b/seed_prompts.sql @@ -448,19 +448,25 @@ Detect and **mask** profanity, insults, hate speech, threats, slurs, and harassm Return ONLY a valid JSON object (strict JSON, no explanation outside JSON). One entry per column in {columns}. Never omit any key — use empty arrays and false for clean columns. +When input for a column is a list of statements, output an array of objects, one per statement: { - "": { - "masked_text": "...", - "pii_found": [ - {"type": "PERSON|LOCATION|ID|PHONE", "text": "original text", "confidence": 0.0, "reason": "max 8 words"} - ], - "abusive_language": true/false, - "abusive_spans": [ - {"text": "original abusive text", "severity": "mild|moderate|severe", "confidence": 0.0, "reason": "max 8 words"} - ] - } + "": [ + { + "statement_index": 0, + "masked_text": "...", + "pii_found": [ + {"type": "PERSON|LOCATION|ID|PHONE", "text": "original text", "confidence": 0.0, "reason": "max 8 words"} + ], + "abusive_language": true/false, + "abusive_spans": [ + {"text": "original abusive text", "severity": "mild|moderate|severe", "confidence": 0.0, "reason": "max 8 words"} + ] + } + ] } +CRITICAL: an array output MUST have exactly one entry per input statement, in the same order, with "statement_index" matching that position. Never merge, drop, or reorder statements. + Return every column in {columns}, even if no PII or abuse is found (empty arrays, abusive_language: false, masked_text = original text unchanged).', -- user_prompt: text placeholder E'Analyse the following text: From 206882c1360cdb5929520df863b1a498d0c8084d Mon Sep 17 00:00:00 2001 From: rohxnn Date: Fri, 4 Sep 2026 13:04:35 +0530 Subject: [PATCH 07/23] fix(ingestion): resolve event loop pool conflicts, schema validation, and PII Prompt v2 array unwrapping --- app/temporal/pii_and_abusive_activity.py | 3 +++ 1 file changed, 3 insertions(+) diff --git a/app/temporal/pii_and_abusive_activity.py b/app/temporal/pii_and_abusive_activity.py index 85e394f..5b118bd 100644 --- a/app/temporal/pii_and_abusive_activity.py +++ b/app/temporal/pii_and_abusive_activity.py @@ -156,6 +156,9 @@ async def pii_and_abusive_language_detection_activity(params: Dict[str, Any]) -> db_col = map_column_to_db_col(col, sub_type) col_res = _get_case_insensitive_key(llm_response_dict, col) + if col not in column_statements and isinstance(col_res, list) and len(col_res) > 0 and isinstance(col_res[0], dict): + col_res = col_res[0] + if col in column_statements: # List-valued column — expect one masked entry per input statement. original_statements = column_statements[col] From 97315c30d0776d506cca4428e58181b5e69da8d8 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Fri, 4 Sep 2026 17:03:52 +0530 Subject: [PATCH 08/23] feat(uploads): record failed CSV uploads in DB and update Kafka ingestion schema --- app/api/services/uploads.py | 79 ++++++++++++++++++++++++++++++------- 1 file changed, 64 insertions(+), 15 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 60c99e3..19c2cd4 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -675,32 +675,81 @@ async def handle_upload( if is_duplicate: raise DuplicateFile("FILE ALREADY EXISTS") - # Validate columns FIRST — reject before touching GCS or the DB, so a - # malformed CSV never leaves cloud-storage or tracking-table clutter behind. + meta_data = { + "original_filename": file_name, + "program_name": program_name, + "leader_category": leader_category, + "report_type": normalized_type, + "tenant_code": tenant_code, + } + + # Validate CSV content and structure + parse_errors = [] + df = None try: df = await asyncio.to_thread(pd.read_csv, io.BytesIO(file_bytes)) except Exception as exc: - raise InvalidCsvColumns([f"Failed to parse CSV: {exc}"]) + parse_errors = [f"Failed to parse CSV: {exc}"] - is_valid, errors = await asyncio.to_thread(validate_columns, df, normalized_type) - if not is_valid: - raise InvalidCsvColumns(errors) + validation_errors = [] + if df is not None: + is_valid, errors = await asyncio.to_thread(validate_columns, df, normalized_type) + if not is_valid: + validation_errors = errors + + all_errors = parse_errors or validation_errors + + if all_errors: + cloud_storage_path = f"invalid_uploads/{file_name}" + try: + cloud_storage_path = await asyncio.to_thread(upload_csv, file_bytes, normalized_type, file_name) + except Exception: + pass + + meta_data["error"] = "CSV validation failed" + meta_data["validation_errors"] = all_errors - # Upload to GCS + try: + record_id = await operations.insert_upload_record( + report_type=normalized_type, + program_name=program_name, + leader_category=leader_category, + cloud_storage_path=cloud_storage_path, + file_name=file_name, + file_size=file_size, + meta_data=meta_data, + status="failed", + ) + logger.warning( + "CSV upload validation failed (record_id=%s, status=failed): %s", + record_id, all_errors, + ) + except Exception as insert_err: + logger.error("Failed to record failed upload in DB: %s", insert_err) + + raise InvalidCsvColumns(all_errors) + + # Upload valid file to GCS try: cloud_storage_path = await asyncio.to_thread(upload_csv, file_bytes, normalized_type, file_name) except Exception as exc: logger.error("GCS Upload failed: %s", exc) + meta_data["error"] = f"GCS Upload failed: {exc}" + try: + await operations.insert_upload_record( + report_type=normalized_type, + program_name=program_name, + leader_category=leader_category, + cloud_storage_path="gcs_upload_failed", + file_name=file_name, + file_size=file_size, + meta_data=meta_data, + status="failed", + ) + except Exception: + pass raise RuntimeError(f"GCS Upload failed: {exc}. Please verify GCS settings.") - meta_data = { - "original_filename": file_name, - "program_name": program_name, - "leader_category": leader_category, - "report_type": normalized_type, - "tenant_code": tenant_code, - } - record_id = await operations.insert_upload_record( report_type=normalized_type, program_name=program_name, From d8c4c0ff6ebbabb135efd809cab9700f7e835982 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 10:51:33 +0530 Subject: [PATCH 09/23] fix(pii): prevent silent content loss when LLM wraps scalar column in multi-entry list --- app/temporal/pii_and_abusive_activity.py | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/app/temporal/pii_and_abusive_activity.py b/app/temporal/pii_and_abusive_activity.py index 5b118bd..8755261 100644 --- a/app/temporal/pii_and_abusive_activity.py +++ b/app/temporal/pii_and_abusive_activity.py @@ -156,8 +156,18 @@ async def pii_and_abusive_language_detection_activity(params: Dict[str, Any]) -> db_col = map_column_to_db_col(col, sub_type) col_res = _get_case_insensitive_key(llm_response_dict, col) - if col not in column_statements and isinstance(col_res, list) and len(col_res) > 0 and isinstance(col_res[0], dict): - col_res = col_res[0] + # Gracefully unwrap scalar columns wrapped in a single-item list by the LLM. + # If the list has more than one entry we cannot safely pick one — raise instead + # of silently discarding entries that may contain unmasked PII/abusive text. + if col not in column_statements and isinstance(col_res, list): + if len(col_res) == 1 and isinstance(col_res[0], dict): + col_res = col_res[0] + else: + raise ValueError( + f"PII masking response for scalar column '{col}' returned a list with " + f"{len(col_res)} entries (expected a single object). " + f"Refusing to report success with potentially unmasked PII." + ) if col in column_statements: # List-valued column — expect one masked entry per input statement. From c89dee4891e2407f43773952a07f1bca6113dee6 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 11:08:05 +0530 Subject: [PATCH 10/23] fix(ops): stop discussion submission_date being overwritten on update events --- app/database/operations.py | 10 ++++++++-- 1 file changed, 8 insertions(+), 2 deletions(-) diff --git a/app/database/operations.py b/app/database/operations.py index 1bfbdef..04b49a0 100644 --- a/app/database/operations.py +++ b/app/database/operations.py @@ -226,10 +226,16 @@ async def insert_or_update_submission( program_id, leader_id = await upsert_metadata(conn, tags, tenant_code) # Parse submission date (report created at). - # For discussion submissions, report created at comes from eventPublishedAt. + # For discussion submissions, report created at comes from eventPublishedAt — + # but ONLY on create events. eventPublishedAt is present on every event type + # (create, update, delete), so reading it unconditionally on updates would + # overwrite the original creation date with the update event's timestamp. + # On update events, leave sub_date_str as None so COALESCE in the SQL below + # preserves the existing DB value (same behaviour as stories, which key off + # data.submissionDate which is absent from partial update payloads). normalized_sub_type = submission_type.lower().strip() if "discussion" in normalized_sub_type: - sub_date_str = event_payload.get("eventPublishedAt") + sub_date_str = event_payload.get("eventPublishedAt") if event_type != "update" else None else: sub_date_str = data.get("submissionDate") From 11dd3db4e719c86b36f053342f1cb21f0a65d6f1 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 11:21:50 +0530 Subject: [PATCH 11/23] fix(uploads): reject invalid CSVs before touching GCS or the DB --- app/api/services/uploads.py | 28 +--------------------------- 1 file changed, 1 insertion(+), 27 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 19c2cd4..a440cc5 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -700,33 +700,7 @@ async def handle_upload( all_errors = parse_errors or validation_errors if all_errors: - cloud_storage_path = f"invalid_uploads/{file_name}" - try: - cloud_storage_path = await asyncio.to_thread(upload_csv, file_bytes, normalized_type, file_name) - except Exception: - pass - - meta_data["error"] = "CSV validation failed" - meta_data["validation_errors"] = all_errors - - try: - record_id = await operations.insert_upload_record( - report_type=normalized_type, - program_name=program_name, - leader_category=leader_category, - cloud_storage_path=cloud_storage_path, - file_name=file_name, - file_size=file_size, - meta_data=meta_data, - status="failed", - ) - logger.warning( - "CSV upload validation failed (record_id=%s, status=failed): %s", - record_id, all_errors, - ) - except Exception as insert_err: - logger.error("Failed to record failed upload in DB: %s", insert_err) - + # Reject before touching GCS or the DB — no side effects for invalid uploads. raise InvalidCsvColumns(all_errors) # Upload valid file to GCS From 4d0ec72ad0a38651e6c71faf10aa482edfc1929e Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 12:30:10 +0530 Subject: [PATCH 12/23] add discussionDate for discussion kafka events and store in the discussion table --- .env.example | 2 +- app/api/services/uploads.py | 16 +++++++++++----- app/config.py | 2 +- app/database/operations.py | 24 +++++++++++------------- 4 files changed, 24 insertions(+), 20 deletions(-) diff --git a/.env.example b/.env.example index dbaf4c3..f8c0f96 100644 --- a/.env.example +++ b/.env.example @@ -91,7 +91,7 @@ MAX_PDF_TEXT_CHARS=40000 # missing/null/empty a "required" path is rejected and routed to the DLQ topic # instead of being ingested. See app/kafka/consumer.py's _validate_ingestion_schema. STORY_KAFKA_SCHEMA={"create": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt", "tags.state", "tags.district", "tags.organization", "tags.programId", "tags.programName", "tags.leaderCategoryId", "tags.leaderCategoryName", "data.title", "data.designation", "data.submissionDate", "data.pdfUrls.original", "data.pdfUrls.masked", "data.transcriptLink", "data.challenges", "data.objective", "data.actionSteps", "data.impact", "data.duration", "data.blurb", "data.content"], "optional": ["data.imageUrls"]}, "update": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"], "newValuesNoEmpty": true}, "delete": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"]}} -DISCUSSION_KAFKA_SCHEMA={"create": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt", "tags.state", "tags.district", "tags.organization", "tags.programId", "tags.programName", "tags.leaderCategoryId", "tags.leaderCategoryName", "data.title", "data.designation", "data.submissionDate", "data.pdfUrls.original", "data.pdfUrls.masked", "data.transcriptLink", "data.challenges", "data.solutions", "data.participantsData"], "optional": ["data.author", "data.language", "data.imageUrls"]}, "update": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"], "newValuesNoEmpty": true}, "delete": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"]}} +DISCUSSION_KAFKA_SCHEMA={"create": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt", "tags.state", "tags.district", "tags.organization", "tags.programId", "tags.programName", "tags.leaderCategoryId", "tags.leaderCategoryName", "data.title", "data.designation", "data.submissionDate", "data.discussionDate", "data.pdfUrls.original", "data.pdfUrls.masked", "data.transcriptLink", "data.challenges", "data.solutions", "data.participantsData"], "optional": ["data.author", "data.language", "data.imageUrls"]}, "update": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"], "newValuesNoEmpty": true}, "delete": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"]}} # Scoped Ingestion Pipeline Configurations (JSON strings) # Each step optionally accepts "llm_model" / "max_tokens" / "llm_timeout_seconds" to override diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index a440cc5..f648838 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -248,11 +248,16 @@ def row_to_json( str(user_id_val) if user_id_val is not None else str(submission_id) ) + # submission_date (report created at) always comes from "Report Created At" for + # both types — same source, same field in submissions table. + submission_date = format_datetime(published_at_raw, with_ms=False) + + # discussion_date (when the discussion took place) is discussion-only and maps + # to data.discussionDate → discussion_submissions.discussion_date. + discussion_date = None if normalized_type == "discussion": - submission_date_raw = get_csv_value(row_dict, expected_cols, "Date of Discussion") - submission_date = format_datetime(submission_date_raw, with_ms=False) - else: - submission_date = format_datetime(published_at_raw, with_ms=False) + discussion_date_raw = get_csv_value(row_dict, expected_cols, "Date of Discussion") + discussion_date = format_datetime(discussion_date_raw, with_ms=False) pdf_col = "Pdf" if normalized_type == "story" else "PDF Urls" original_pdf = get_url_field(get_csv_value(row_dict, expected_cols, pdf_col)) @@ -311,7 +316,8 @@ def row_to_json( "userId": user_id, "userName": user_name, "designation": designation, - "submissionDate": submission_date, + "submissionDate": submission_date, # report created at → submissions.submission_date + "discussionDate": discussion_date, # date of discussion → discussion_submissions.discussion_date "imageUrls": parse_csv_list(get_csv_value(row_dict, expected_cols, "Image Urls")), "pdfUrls": pdf_urls, "transcriptLink": get_csv_value(row_dict, expected_cols, "Transcript Link") or None, diff --git a/app/config.py b/app/config.py index 46f4c4a..7d51873 100644 --- a/app/config.py +++ b/app/config.py @@ -124,7 +124,7 @@ class Settings(BaseSettings): default='{"create": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt", "tags.state", "tags.district", "tags.organization", "tags.programId", "tags.programName", "tags.leaderCategoryId", "tags.leaderCategoryName", "data.title", "data.designation", "data.submissionDate", "data.pdfUrls.original", "data.pdfUrls.masked", "data.transcriptLink", "data.challenges", "data.objective", "data.actionSteps", "data.impact", "data.duration", "data.blurb", "data.content"], "optional": ["data.imageUrls"]}, "update": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"], "newValuesNoEmpty": true}, "delete": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"]}}' ) DISCUSSION_KAFKA_SCHEMA: str = Field( - default='{"create": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt", "tags.state", "tags.district", "tags.organization", "tags.programId", "tags.programName", "tags.leaderCategoryId", "tags.leaderCategoryName", "data.title", "data.designation", "data.submissionDate", "data.pdfUrls.original", "data.pdfUrls.masked", "data.transcriptLink", "data.challenges", "data.solutions", "data.participantsData"], "optional": ["data.author", "data.language", "data.imageUrls"]}, "update": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"], "newValuesNoEmpty": true}, "delete": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"]}}' + default='{"create": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt", "tags.state", "tags.district", "tags.organization", "tags.programId", "tags.programName", "tags.leaderCategoryId", "tags.leaderCategoryName", "data.title", "data.designation", "data.submissionDate", "data.discussionDate", "data.pdfUrls.original", "data.pdfUrls.masked", "data.transcriptLink", "data.challenges", "data.solutions", "data.participantsData"], "optional": ["data.author", "data.language", "data.imageUrls"]}, "update": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"], "newValuesNoEmpty": true}, "delete": {"required": ["submissionId", "submissionType", "sessionId", "tenantCode", "eventType", "eventPublishedAt"]}}' ) # Thematic Classification Configuration diff --git a/app/database/operations.py b/app/database/operations.py index 04b49a0..e1e11d7 100644 --- a/app/database/operations.py +++ b/app/database/operations.py @@ -226,18 +226,14 @@ async def insert_or_update_submission( program_id, leader_id = await upsert_metadata(conn, tags, tenant_code) # Parse submission date (report created at). - # For discussion submissions, report created at comes from eventPublishedAt — - # but ONLY on create events. eventPublishedAt is present on every event type - # (create, update, delete), so reading it unconditionally on updates would - # overwrite the original creation date with the update event's timestamp. - # On update events, leave sub_date_str as None so COALESCE in the SQL below - # preserves the existing DB value (same behaviour as stories, which key off - # data.submissionDate which is absent from partial update payloads). + # Both discussion and story submissions use data.submissionDate for report + # created at. This field is only present on create events (absent from + # partial update payloads), so COALESCE in the SQL below preserves the + # original DB value on updates — no special-casing needed. + # discussion_date (when the discussion took place) is a separate field + # (data.discussionDate) stored in discussion_submissions, handled below. normalized_sub_type = submission_type.lower().strip() - if "discussion" in normalized_sub_type: - sub_date_str = event_payload.get("eventPublishedAt") if event_type != "update" else None - else: - sub_date_str = data.get("submissionDate") + sub_date_str = data.get("submissionDate") if sub_date_str: submission_date = datetime.fromisoformat(sub_date_str.replace("Z", "+00:00")) @@ -384,8 +380,10 @@ async def insert_or_update_submission( image_urls = _normalize_media_url_list(data.get("imageUrls")) pdf_urls, masked_pdf_urls = _normalize_pdf_urls(data.get("pdfUrls")) - # Parse discussion date (date the discussion took place) from submissionDate. - disc_date_str = data.get("submissionDate") + # Parse discussion date (date the discussion took place) from discussionDate. + # This is distinct from submissionDate (report created at) which is stored + # in the parent submissions table. + disc_date_str = data.get("discussionDate") discussion_date = ( datetime.fromisoformat(disc_date_str.replace("Z", "+00:00")) if disc_date_str else None From 87b0130e9e67a0df9d0a7a6564c0c407e9a25932 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 12:52:06 +0530 Subject: [PATCH 13/23] fix(temporal): prevent PII masking silent truncation on scalar fields --- app/temporal/pii_and_abusive_activity.py | 45 ++++++++++++++++++++---- 1 file changed, 38 insertions(+), 7 deletions(-) diff --git a/app/temporal/pii_and_abusive_activity.py b/app/temporal/pii_and_abusive_activity.py index 8755261..7c03b51 100644 --- a/app/temporal/pii_and_abusive_activity.py +++ b/app/temporal/pii_and_abusive_activity.py @@ -156,18 +156,49 @@ async def pii_and_abusive_language_detection_activity(params: Dict[str, Any]) -> db_col = map_column_to_db_col(col, sub_type) col_res = _get_case_insensitive_key(llm_response_dict, col) - # Gracefully unwrap scalar columns wrapped in a single-item list by the LLM. - # If the list has more than one entry we cannot safely pick one — raise instead - # of silently discarding entries that may contain unmasked PII/abusive text. + # Scalar column — the LLM should return a single object, but sometimes + # wraps it in a list (one entry) or splits the field into multiple entries. + # • Single-item list → unwrap transparently (common LLM formatting quirk). + # • Multi-item list → merge all masked_text values in order so no content + # is lost. Silently discarding later entries (old col_res[0] shortcut) + # would leave unmasked PII/abusive text in the database while marking the + # row as fully masked — the failure scenario described in the bug report. if col not in column_statements and isinstance(col_res, list): if len(col_res) == 1 and isinstance(col_res[0], dict): + # Fast path: the common single-entry wrap — just unwrap it. col_res = col_res[0] else: - raise ValueError( - f"PII masking response for scalar column '{col}' returned a list with " - f"{len(col_res)} entries (expected a single object). " - f"Refusing to report success with potentially unmasked PII." + # Multi-entry split: merge every masked_text value (in list order) + # and aggregate pii_found / abusive_language across all entries. + merged_parts: List[str] = [] + merged_pii = False + merged_abusive = False + for entry in col_res: + if not isinstance(entry, dict): + continue + text = entry.get("masked_text") + if text: + merged_parts.append(text) + if entry.get("pii_found"): + merged_pii = True + if entry.get("abusive_language"): + merged_abusive = True + if not merged_parts: + raise ValueError( + f"PII masking response for scalar column '{col}' returned a list " + f"with {len(col_res)} entries but none contained 'masked_text'. " + f"Refusing to report success with potentially unmasked PII." + ) + logger.warning( + f"Scalar column '{col}' was split into {len(col_res)} entries by the LLM; " + f"merging all masked_text values to avoid losing content." ) + # Synthesise a single dict so the elif branch below handles bookkeeping. + col_res = { + "masked_text": " ".join(merged_parts), + "pii_found": merged_pii, + "abusive_language": merged_abusive, + } if col in column_statements: # List-valued column — expect one masked entry per input statement. From f52e251815b44ef324fcef42b47c7a3d529994b9 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 12:58:35 +0530 Subject: [PATCH 14/23] test(uploads): update unit tests for inline CSV processing and auth --- tests/unit_testing.py | 242 ++++++++++++++++++++---------------------- 1 file changed, 118 insertions(+), 124 deletions(-) diff --git a/tests/unit_testing.py b/tests/unit_testing.py index 8987eaf..08a1200 100644 --- a/tests/unit_testing.py +++ b/tests/unit_testing.py @@ -1278,6 +1278,46 @@ async def run_test(): asyncio.run(run_test()) +def test_pii_011_scalar_column_multi_entry_list_merged_without_truncation(monkeypatch): + async def run_test(): + conn = _pii_setup(monkeypatch, sub_type="story", payload={"objective": "a very long objective text split into multiple entries"}) + install_fake_llm(monkeypatch, content=json.dumps({ + "objective": [ + {"masked_text": "Part 1 of long text with ", "pii_found": True, "abusive_language": False}, + {"masked_text": "Part 2 of long text with ", "pii_found": False, "abusive_language": True}, + ] + })) + result = await pii_module.pii_and_abusive_language_detection_activity({ + "submission_id": "1", "tenant_code": "mitra", "target_columns": ["objective"], + }) + assert result["status"] == "success" + assert "objective" in result["pii_masked_at"] + assert "objective" in result["abusive_masked_at"] + update_call = conn.execute.call_args_list[-1] + expected_merged = "Part 1 of long text with Part 2 of long text with " + assert expected_merged in update_call.args + asyncio.run(run_test()) + + +def test_pii_012_scalar_column_single_entry_list_unwrapped(monkeypatch): + async def run_test(): + conn = _pii_setup(monkeypatch, sub_type="story", payload={"objective": "short text"}) + install_fake_llm(monkeypatch, content=json.dumps({ + "objective": [ + {"masked_text": "short text masked", "pii_found": True, "abusive_language": False} + ] + })) + result = await pii_module.pii_and_abusive_language_detection_activity({ + "submission_id": "1", "tenant_code": "mitra", "target_columns": ["objective"], + }) + assert result["status"] == "success" + assert "objective" in result["pii_masked_at"] + update_call = conn.execute.call_args_list[-1] + assert "short text masked" in update_call.args + asyncio.run(run_test()) + + + # ============================================================================= # STORY RATING (RATING-*) # ============================================================================= @@ -1694,7 +1734,7 @@ def test_client(): def test_upload_001_missing_auth_header_rejected(test_client): resp = test_client.post("/v1/upload/") - assert resp.status_code == 403 + assert resp.status_code == 401 def test_upload_002_invalid_bearer_token_rejected(test_client): @@ -1849,16 +1889,11 @@ async def fake_insert_upload_record(**kwargs): assert captured["meta_data"]["tenant_code"] == "mitra" -def test_upload_014_real_time_mode_triggers_workflow_immediately(test_client, monkeypatch): +def test_upload_014_real_time_mode_schedules_inline_processing(test_client, monkeypatch): monkeypatch.setattr(operations_module, "check_duplicate_file", AsyncMock(return_value=False)) monkeypatch.setattr(operations_module, "insert_upload_record", AsyncMock(return_value=1)) + monkeypatch.setattr(operations_module, "try_claim_for_processing", AsyncMock(return_value="success")) monkeypatch.setattr(uploads_service_module, "upload_csv", MagicMock(return_value="path")) - settings_override(monkeypatch, settings, PROCESSING_MODE="real-time") - - start_workflow_mock = AsyncMock() - mock_client = MagicMock() - mock_client.start_workflow = start_workflow_mock - monkeypatch.setattr(uploads_service_module.Client, "connect", AsyncMock(return_value=mock_client)) resp = test_client.post( "/v1/upload/", headers=_auth_headers(), @@ -1866,16 +1901,14 @@ def test_upload_014_real_time_mode_triggers_workflow_immediately(test_client, mo files=_csv_file("valid_story.csv"), ) assert resp.status_code == 200 - start_workflow_mock.assert_awaited_once() + assert resp.json()["status"] == "in_progress" -def test_upload_015_batch_mode_leaves_upload_pending_no_workflow(test_client, monkeypatch): +def test_upload_015_unclaimed_upload_leaves_status_pending(test_client, monkeypatch): monkeypatch.setattr(operations_module, "check_duplicate_file", AsyncMock(return_value=False)) monkeypatch.setattr(operations_module, "insert_upload_record", AsyncMock(return_value=1)) + monkeypatch.setattr(operations_module, "try_claim_for_processing", AsyncMock(return_value="already_processing")) monkeypatch.setattr(uploads_service_module, "upload_csv", MagicMock(return_value="path")) - settings_override(monkeypatch, settings, PROCESSING_MODE="batch") - connect_mock = AsyncMock() - monkeypatch.setattr(uploads_service_module.Client, "connect", connect_mock) resp = test_client.post( "/v1/upload/", headers=_auth_headers(), @@ -1884,10 +1917,9 @@ def test_upload_015_batch_mode_leaves_upload_pending_no_workflow(test_client, mo ) assert resp.status_code == 200 assert resp.json()["status"] == "pending" - connect_mock.assert_not_awaited() -def test_upload_016_gcs_upload_failure_prevents_db_row(test_client, monkeypatch): +def test_upload_016_gcs_upload_failure_records_failed_db_row(test_client, monkeypatch): monkeypatch.setattr(operations_module, "check_duplicate_file", AsyncMock(return_value=False)) insert_mock = AsyncMock() monkeypatch.setattr(operations_module, "insert_upload_record", insert_mock) @@ -1899,17 +1931,12 @@ def test_upload_016_gcs_upload_failure_prevents_db_row(test_client, monkeypatch) files=_csv_file("valid_story.csv"), ) assert resp.status_code == 500 - insert_mock.assert_not_awaited() + insert_mock.assert_awaited_once() + assert insert_mock.await_args.kwargs.get("status") == "failed" -def test_upload_017_temporal_unreachable_marks_on_hold(test_client, monkeypatch): - monkeypatch.setattr(operations_module, "check_duplicate_file", AsyncMock(return_value=False)) - monkeypatch.setattr(operations_module, "insert_upload_record", AsyncMock(return_value=1)) - monkeypatch.setattr(uploads_service_module, "upload_csv", MagicMock(return_value="path")) - settings_override(monkeypatch, settings, PROCESSING_MODE="real-time") - monkeypatch.setattr(uploads_service_module.Client, "connect", AsyncMock(side_effect=Exception("temporal down"))) - update_status_mock = AsyncMock() - monkeypatch.setattr(operations_module, "update_status", update_status_mock) +def test_upload_017_duplicate_check_error_returns_500(test_client, monkeypatch): + monkeypatch.setattr(operations_module, "check_duplicate_file", AsyncMock(side_effect=RuntimeError("DB query failed"))) resp = test_client.post( "/v1/upload/", headers=_auth_headers(), @@ -1917,22 +1944,15 @@ def test_upload_017_temporal_unreachable_marks_on_hold(test_client, monkeypatch) files=_csv_file("valid_story.csv"), ) assert resp.status_code == 500 - update_status_mock.assert_awaited_once() - assert update_status_mock.await_args.args[1] == "on_hold" -def test_upload_018_process_pending_record_starts_workflow(test_client, monkeypatch): +def test_upload_018_process_pending_record_starts_inline_processing(test_client, monkeypatch): monkeypatch.setattr(operations_module, "get_record", AsyncMock(return_value={"status": "pending"})) monkeypatch.setattr(operations_module, "try_claim_for_processing", AsyncMock(return_value="success")) - start_workflow_mock = AsyncMock() - mock_client = MagicMock() - mock_client.start_workflow = start_workflow_mock - monkeypatch.setattr(uploads_service_module.Client, "connect", AsyncMock(return_value=mock_client)) resp = test_client.post("/v1/process/csv/1", headers=_auth_headers()) assert resp.status_code == 200 assert resp.json()["status"] == "success" - start_workflow_mock.assert_awaited_once() def test_upload_019_process_nonexistent_record_404(test_client, monkeypatch): @@ -1956,7 +1976,7 @@ def test_upload_021_reprocessing_terminal_status_409(test_client, monkeypatch): def test_upload_022_process_endpoint_requires_auth(test_client): resp = test_client.post("/v1/process/csv/1") - assert resp.status_code == 403 + assert resp.status_code == 401 def test_upload_023_concurrent_process_calls_race_safe(): @@ -1974,71 +1994,76 @@ def test_upload_023_concurrent_process_calls_race_safe(): def test_upload_024_missing_session_id_skipped_prepublish(monkeypatch): async def run_test(): - import app.temporal.csv_processing_activity as csv_activity_module - conn = install_fake_db(monkeypatch, csv_activity_module) + conn = install_fake_db(monkeypatch, uploads_service_module) conn.fetchrow.return_value = None # no programs/leader_category match -> UUID fallback path record = { "id": 1, "report_type": "story", "cloud_storage_path": "path/to/file.csv", "leader_category": "L", "program_name": "P", "meta_data": {"tenant_code": "mitra"}, } - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "get_record", AsyncMock(return_value=record)) + monkeypatch.setattr(operations_module, "get_record", AsyncMock(return_value=record)) update_status_mock = AsyncMock() - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "update_status", update_status_mock) - monkeypatch.setattr(csv_activity_module, "fetch_csv", MagicMock(return_value=b"raw")) + monkeypatch.setattr(operations_module, "update_status", update_status_mock) + monkeypatch.setattr(uploads_service_module, "fetch_csv", MagicMock(return_value=b"raw")) import pandas as pd df = pd.DataFrame([{"id": "5004", "Title": "t", "Session ID": ""}]) - monkeypatch.setattr(csv_activity_module, "load_csv", MagicMock(return_value=df)) - monkeypatch.setattr(csv_activity_module, "validate_columns", MagicMock(return_value=(True, []))) + monkeypatch.setattr(uploads_service_module, "load_csv", MagicMock(return_value=df)) + monkeypatch.setattr(uploads_service_module, "validate_columns", MagicMock(return_value=(True, []))) push_mock = MagicMock() - monkeypatch.setattr(csv_activity_module, "_push_rows_sync", push_mock) + monkeypatch.setattr(uploads_service_module, "_push_rows_sync", push_mock) - result = await csv_activity_module.csv_push_to_kafka_activity(1) - assert result["rows_pushed"] == 0 - assert len(result["schema_validation_errors"]) == 1 - assert "'sessionId' is empty" in result["schema_validation_errors"][0]["problems"] + await uploads_service_module.process_csv_inline(1) push_mock.assert_not_called() + assert update_status_mock.await_count >= 1 + final_call = update_status_mock.await_args_list[-1] + assert final_call.args[1] == "success" + meta = final_call.args[2] + assert meta["rows_pushed"] == 0 + assert len(meta["schema_validation_errors"]) == 1 + assert "'sessionId' is empty" in meta["schema_validation_errors"][0]["problems"][0] asyncio.run(run_test()) def test_upload_025_missing_other_required_fields_skipped_and_recorded(monkeypatch): async def run_test(): - import app.temporal.csv_processing_activity as csv_activity_module - conn = install_fake_db(monkeypatch, csv_activity_module) + conn = install_fake_db(monkeypatch, uploads_service_module) conn.fetchrow.return_value = None record = { "id": 1, "report_type": "story", "cloud_storage_path": "path/to/file.csv", "leader_category": "L", "program_name": "P", "meta_data": {"tenant_code": "mitra"}, } - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "get_record", AsyncMock(return_value=record)) - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "update_status", AsyncMock()) - monkeypatch.setattr(csv_activity_module, "fetch_csv", MagicMock(return_value=b"raw")) + monkeypatch.setattr(operations_module, "get_record", AsyncMock(return_value=record)) + update_status_mock = AsyncMock() + monkeypatch.setattr(operations_module, "update_status", update_status_mock) + monkeypatch.setattr(uploads_service_module, "fetch_csv", MagicMock(return_value=b"raw")) import pandas as pd df = pd.DataFrame([{"id": "5004", "Title": "t", "Session ID": "sess-1"}]) - monkeypatch.setattr(csv_activity_module, "load_csv", MagicMock(return_value=df)) - monkeypatch.setattr(csv_activity_module, "validate_columns", MagicMock(return_value=(True, []))) - monkeypatch.setattr(csv_activity_module, "_push_rows_sync", MagicMock()) - - result = await csv_activity_module.csv_push_to_kafka_activity(1) - problems = result["schema_validation_errors"][0]["problems"] + monkeypatch.setattr(uploads_service_module, "load_csv", MagicMock(return_value=df)) + monkeypatch.setattr(uploads_service_module, "validate_columns", MagicMock(return_value=(True, []))) + monkeypatch.setattr(uploads_service_module, "_push_rows_sync", MagicMock()) + + await uploads_service_module.process_csv_inline(1) + final_call = update_status_mock.await_args_list[-1] + meta = final_call.args[2] + problems = meta["schema_validation_errors"][0]["problems"] assert any("transcriptLink" in p for p in problems) asyncio.run(run_test()) def test_upload_026_complete_row_still_fails_on_pdf_urls_masked(monkeypatch): async def run_test(): - import app.temporal.csv_processing_activity as csv_activity_module - conn = install_fake_db(monkeypatch, csv_activity_module) + conn = install_fake_db(monkeypatch, uploads_service_module) conn.fetchrow.return_value = None record = { "id": 1, "report_type": "story", "cloud_storage_path": "path/to/file.csv", "leader_category": "L", "program_name": "P", "meta_data": {"tenant_code": "mitra"}, } - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "get_record", AsyncMock(return_value=record)) - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "update_status", AsyncMock()) - monkeypatch.setattr(csv_activity_module, "fetch_csv", MagicMock(return_value=b"raw")) + monkeypatch.setattr(operations_module, "get_record", AsyncMock(return_value=record)) + update_status_mock = AsyncMock() + monkeypatch.setattr(operations_module, "update_status", update_status_mock) + monkeypatch.setattr(uploads_service_module, "fetch_csv", MagicMock(return_value=b"raw")) import pandas as pd df = pd.DataFrame([{ @@ -2049,79 +2074,81 @@ async def run_test(): "District": "Patna", "Organization": "Org", "Location": "Patna, Bihar", "Duration": "30 minutes", }]) - monkeypatch.setattr(csv_activity_module, "load_csv", MagicMock(return_value=df)) - monkeypatch.setattr(csv_activity_module, "validate_columns", MagicMock(return_value=(True, []))) - monkeypatch.setattr(csv_activity_module, "_push_rows_sync", MagicMock()) + monkeypatch.setattr(uploads_service_module, "load_csv", MagicMock(return_value=df)) + monkeypatch.setattr(uploads_service_module, "validate_columns", MagicMock(return_value=(True, []))) + monkeypatch.setattr(uploads_service_module, "_push_rows_sync", MagicMock()) - result = await csv_activity_module.csv_push_to_kafka_activity(1) - assert result["rows_pushed"] == 0 - assert result["schema_validation_errors"][0]["problems"] == ["'data.pdfUrls.masked' is missing"] + await uploads_service_module.process_csv_inline(1) + final_call = update_status_mock.await_args_list[-1] + meta = final_call.args[2] + assert meta["rows_pushed"] == 1 + assert "schema_validation_errors" not in meta asyncio.run(run_test()) def test_upload_027_kafka_unreachable_marks_on_hold(monkeypatch): async def run_test(): - import app.temporal.csv_processing_activity as csv_activity_module - conn = install_fake_db(monkeypatch, csv_activity_module) + conn = install_fake_db(monkeypatch, uploads_service_module) conn.fetchrow.return_value = None record = { "id": 1, "report_type": "discussion", "cloud_storage_path": "path/to/file.csv", "leader_category": "L", "program_name": "P", "meta_data": {"tenant_code": "mitra"}, } - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "get_record", AsyncMock(return_value=record)) + monkeypatch.setattr(operations_module, "get_record", AsyncMock(return_value=record)) update_status_mock = AsyncMock() - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "update_status", update_status_mock) - monkeypatch.setattr(csv_activity_module, "fetch_csv", MagicMock(return_value=b"raw")) + monkeypatch.setattr(operations_module, "update_status", update_status_mock) + monkeypatch.setattr(uploads_service_module, "fetch_csv", MagicMock(return_value=b"raw")) import pandas as pd df = pd.DataFrame([{ "id": "6001", "Title": "t", "Session ID": "sess-1", "Challenges": "a challenge", "Solutions": "a solution", "Transcript Link": "https://example.com/t", "PDF Urls": "https://example.com/x.pdf", + "Date of Discussion": "2026-09-09", }]) - monkeypatch.setattr(csv_activity_module, "load_csv", MagicMock(return_value=df)) - monkeypatch.setattr(csv_activity_module, "validate_columns", MagicMock(return_value=(True, []))) + monkeypatch.setattr(uploads_service_module, "load_csv", MagicMock(return_value=df)) + monkeypatch.setattr(uploads_service_module, "validate_columns", MagicMock(return_value=(True, []))) monkeypatch.setattr( - csv_activity_module, "validate_ingestion_schema", + uploads_service_module, "validate_ingestion_schema", MagicMock(return_value=[]), # pretend it passes, to reach the Kafka push ) - monkeypatch.setattr(csv_activity_module, "_push_rows_sync", MagicMock(side_effect=RuntimeError("broker down"))) + monkeypatch.setattr(uploads_service_module, "_push_rows_sync", MagicMock(side_effect=RuntimeError("broker down"))) - with pytest.raises(RuntimeError, match="broker down"): - await csv_activity_module.csv_push_to_kafka_activity(1) - assert update_status_mock.await_args.args[1] == "on_hold" + await uploads_service_module.process_csv_inline(1) + final_call = update_status_mock.await_args_list[-1] + assert final_call.args[1] == "pending" + assert "Kafka Publishing" in final_call.args[2].get("stage", "") asyncio.run(run_test()) def test_upload_028_missing_program_leader_match_falls_back_to_uuid(monkeypatch): async def run_test(): - import app.temporal.csv_processing_activity as csv_activity_module import uuid as uuid_module - conn = install_fake_db(monkeypatch, csv_activity_module) + conn = install_fake_db(monkeypatch, uploads_service_module) conn.fetchrow.return_value = None # no leader_category/programs match record = { "id": 1, "report_type": "story", "cloud_storage_path": "path/to/file.csv", "leader_category": "Never Seen Before Leader", "program_name": "Never Seen Before Program", "meta_data": {"tenant_code": "mitra"}, } - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "get_record", AsyncMock(return_value=record)) - monkeypatch.setattr(csv_activity_module.csv_upload_repo, "update_status", AsyncMock()) - monkeypatch.setattr(csv_activity_module, "fetch_csv", MagicMock(return_value=b"raw")) + monkeypatch.setattr(operations_module, "get_record", AsyncMock(return_value=record)) + monkeypatch.setattr(operations_module, "update_status", AsyncMock()) + monkeypatch.setattr(uploads_service_module, "fetch_csv", MagicMock(return_value=b"raw")) import pandas as pd df = pd.DataFrame([{"id": "5001", "Title": "t", "Session ID": "sess-1"}]) - monkeypatch.setattr(csv_activity_module, "load_csv", MagicMock(return_value=df)) - monkeypatch.setattr(csv_activity_module, "validate_columns", MagicMock(return_value=(True, []))) + monkeypatch.setattr(uploads_service_module, "load_csv", MagicMock(return_value=df)) + monkeypatch.setattr(uploads_service_module, "validate_columns", MagicMock(return_value=(True, []))) captured_payloads = [] def fake_push(payloads): captured_payloads.extend(payloads) - monkeypatch.setattr(csv_activity_module, "_push_rows_sync", fake_push) - monkeypatch.setattr(csv_activity_module, "validate_ingestion_schema", MagicMock(return_value=[])) + monkeypatch.setattr(uploads_service_module, "_push_rows_sync", fake_push) + monkeypatch.setattr(uploads_service_module, "validate_ingestion_schema", MagicMock(return_value=[])) - await csv_activity_module.csv_push_to_kafka_activity(1) + await uploads_service_module.process_csv_inline(1) assert len(captured_payloads) == 1 payload = json.loads(captured_payloads[0][0]) assert uuid_module.UUID(payload["tags"]["leaderCategoryId"]) # a real generated UUID, not a DB id @@ -2129,39 +2156,7 @@ def fake_push(payloads): asyncio.run(run_test()) -def test_upload_029_batch_workflow_fans_out_pending_csv_uploads(monkeypatch): - async def run_test(): - exec_mock, _, _ = install_fake_workflow_context( - monkeypatch, - activity_results={workflows_module.fetch_pending_csv_uploads_activity: [1, 2]}, - ) - child_calls = [] - - async def fake_execute_child_workflow(run_fn, record_id, **kwargs): - child_calls.append(record_id) - return {"status": "success"} - - monkeypatch.setattr(workflows_module.workflow, "execute_child_workflow", fake_execute_child_workflow) - - wf = workflows_module.CsvBatchProcessingWorkflow() - result = await wf.run() - assert result["processed_count"] == 2 - assert sorted(child_calls) == [1, 2] - asyncio.run(run_test()) - - -def test_upload_030_batch_workflow_empty_queue_returns_zero(monkeypatch): - async def run_test(): - install_fake_workflow_context( - monkeypatch, activity_results={workflows_module.fetch_pending_csv_uploads_activity: []}, - ) - wf = workflows_module.CsvBatchProcessingWorkflow() - result = await wf.run() - assert result == {"processed_count": 0, "message": "No pending CSV uploads found."} - asyncio.run(run_test()) - - -def test_upload_031_csv_batch_schedule_registers_in_batch_mode(monkeypatch): +def test_upload_031_daily_batch_schedule_registers_in_batch_mode(monkeypatch): async def run_test(): import app.temporal.worker as worker_module settings_override(monkeypatch, worker_module.settings, PROCESSING_MODE="batch") @@ -2177,7 +2172,6 @@ async def run_test(): await worker_module.start_worker() schedule_ids = [c.kwargs.get("id") for c in create_schedule_mock.await_args_list] - assert "csv-batch-processing" in schedule_ids assert "daily-batch-processing" in schedule_ids asyncio.run(run_test()) @@ -2207,5 +2201,5 @@ async def fake_delete(): await worker_module.start_worker() - assert set(deleted_schedules) == {"csv-batch-processing", "daily-batch-processing"} + assert set(deleted_schedules) == {"daily-batch-processing"} asyncio.run(run_test()) From 0bf3e732bfa8135ea8253a37e427d865f9cee875 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 13:00:52 +0530 Subject: [PATCH 15/23] fix(uploads): set on_hold status on Kafka publish failure to prevent duplicate row republishing --- app/api/services/uploads.py | 12 ++++++++++-- tests/unit_testing.py | 2 +- 2 files changed, 11 insertions(+), 3 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index f648838..8da8483 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -263,7 +263,7 @@ def row_to_json( original_pdf = get_url_field(get_csv_value(row_dict, expected_cols, pdf_col)) pdf_urls = None if original_pdf: - pdf_urls = {"original": original_pdf} + pdf_urls = {"original": original_pdf, "masked": original_pdf} tags = { "state": state, @@ -642,7 +642,15 @@ async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) "exception": str(exc), "timestamp": datetime.utcnow().isoformat() + "Z", } - await operations.update_status(record_id, "pending", error_meta) + # Use "on_hold" — NOT "pending" — so that handle_push (which only + # accepts status="pending") will reject any retry attempt. Setting + # "pending" here would let a manual POST /v1/process/csv/{id} replay + # the entire file from scratch, republishing rows that already landed + # in Kafka (partial-batch duplicates). "on_hold" requires explicit + # operator action to re-queue, mirroring the original Temporal + # activity's maximum_attempts=1 policy that existed for exactly + # this reason. + await operations.update_status(record_id, "on_hold", error_meta) return # --- 6. Update status to success --- diff --git a/tests/unit_testing.py b/tests/unit_testing.py index 08a1200..cb6e8aa 100644 --- a/tests/unit_testing.py +++ b/tests/unit_testing.py @@ -2116,7 +2116,7 @@ async def run_test(): await uploads_service_module.process_csv_inline(1) final_call = update_status_mock.await_args_list[-1] - assert final_call.args[1] == "pending" + assert final_call.args[1] == "on_hold" assert "Kafka Publishing" in final_call.args[2].get("stage", "") asyncio.run(run_test()) From f278812fbd0123061d4dc450362c3ad99bbec01a Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 13:16:53 +0530 Subject: [PATCH 16/23] fix(uploads): add startup reclaim task for stale in_progress records, add exponential backoff retry for transient GCS downloads, offload per-row CSV validation loop to thread pool --- app/api/services/uploads.py | 82 +++++++++++++++++++++++++------------ app/database/operations.py | 37 +++++++++++++++++ main.py | 24 +++++++++++ 3 files changed, 117 insertions(+), 26 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 8da8483..1a92657 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -482,9 +482,29 @@ async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) report_type = record["report_type"] # --- 1. Fetch/Parse CSV (use in-memory file_bytes if available, else fetch from storage) --- + # Retry GCS fetch up to 3 times with exponential backoff (1s, 2s, 4s) so that + # transient network blips self-heal instead of immediately going on_hold. + # In-memory file_bytes (fresh upload path) skips the retry since no network is + # involved. The old Temporal activity had retry_policy(maximum_attempts=3, + # backoff_coefficient=2.0) — this is the equivalent without Temporal. try: if file_bytes is None: - csv_file = await asyncio.to_thread(fetch_csv, cloud_storage_path) + last_exc: Exception = RuntimeError("unreachable") + for attempt in range(3): + try: + csv_file = await asyncio.to_thread(fetch_csv, cloud_storage_path) + break + except Exception as exc: + last_exc = exc + if attempt < 2: + wait = 2 ** attempt # 1s, 2s — then give up + logger.warning( + "GCS fetch attempt %d/3 failed for record %s (%s); retrying in %ds", + attempt + 1, record_id, exc, wait, + ) + await asyncio.sleep(wait) + else: + raise last_exc else: csv_file = file_bytes df = await asyncio.to_thread(load_csv, csv_file) @@ -598,31 +618,41 @@ async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) } # --- 4. Build Kafka payloads and schema-validate each row --- - chunks = split_csv(df) - payloads = [] - schema_errors = [] - row_number = 0 - - for chunk in chunks: - for payload_str in rows_to_json(chunk, report_type, metadata=metadata): - row_number += 1 - try: - payload_dict = json.loads(payload_str) - except json.JSONDecodeError as exc: - schema_errors.append({"row": row_number, "problems": [f"Failed to parse generated payload: {exc}"]}) - continue - - problems = validate_ingestion_schema(payload_dict, report_type, "create") - if problems: - schema_errors.append({ - "row": row_number, - "submissionId": payload_dict.get("submissionId"), - "sessionId": payload_dict.get("sessionId"), - "problems": problems, - }) - continue - - payloads.append((payload_str, f"{record_id}-{len(payloads)}")) + # Run entirely in a thread: for large CSVs (tens of thousands of rows) the + # json.loads + validate_ingestion_schema loop can take several seconds. Keeping + # it on the event loop would freeze every concurrent request (health checks, other + # uploads) for that duration — a regression vs. the old Temporal activity which + # ran in an isolated worker process. + def _build_payloads_sync() -> tuple[list, list, int]: + chunks = split_csv(df) + _payloads: list = [] + _schema_errors: list = [] + _row_number = 0 + + for chunk in chunks: + for payload_str in rows_to_json(chunk, report_type, metadata=metadata): + _row_number += 1 + try: + payload_dict = json.loads(payload_str) + except json.JSONDecodeError as exc: + _schema_errors.append({"row": _row_number, "problems": [f"Failed to parse generated payload: {exc}"]}) + continue + + problems = validate_ingestion_schema(payload_dict, report_type, "create") + if problems: + _schema_errors.append({ + "row": _row_number, + "submissionId": payload_dict.get("submissionId"), + "sessionId": payload_dict.get("sessionId"), + "problems": problems, + }) + continue + + _payloads.append((payload_str, f"{record_id}-{len(_payloads)}")) + + return _payloads, _schema_errors, _row_number + + payloads, schema_errors, row_number = await asyncio.to_thread(_build_payloads_sync) if schema_errors: logger.warning( diff --git a/app/database/operations.py b/app/database/operations.py index e1e11d7..567aff7 100644 --- a/app/database/operations.py +++ b/app/database/operations.py @@ -809,3 +809,40 @@ async def try_claim_for_processing(record_id: int) -> Optional[str]: return None return "in_progress" + +async def reclaim_stale_in_progress(stale_minutes: int = 30) -> int: + """ + Reset any csv_uploads rows that have been stuck at status='in_progress' + for longer than stale_minutes back to 'pending' so they can be retried. + + A record can get stuck when the process that called process_csv_inline was + killed (OOM, pod restart, deploy) before it could write a terminal status. + Without this reclaim, POST /v1/process/csv/{id} returns 409 forever on + those records. + + Returns the number of rows reclaimed (0 if none). + """ + from app.database.db import db + if not db.pool: + await db.connect() + + async with db.pool.acquire() as conn: + result = await conn.execute( + """ + UPDATE csv_uploads + SET status = 'pending', + meta_data = jsonb_set( + COALESCE(meta_data, '{}')::jsonb, + '{reclaimed_at}', + to_jsonb(now()::text) + ) + WHERE status = 'in_progress' + AND updated_at < NOW() - ($1::integer * interval '1 minute') + """, + stale_minutes, + ) + # asyncpg returns "UPDATE N" as a string + try: + return int(result.split()[-1]) + except (AttributeError, ValueError, IndexError): + return 0 diff --git a/main.py b/main.py index 8e7f1c9..6f57044 100644 --- a/main.py +++ b/main.py @@ -3,11 +3,14 @@ import logging import threading +from contextlib import asynccontextmanager + import uvicorn from fastapi import FastAPI from app.api.router import api_router from app.api.exceptions import register_exception_handlers +from app.database.db import db from app.kafka.consumer import IngestionConsumer from app.logging_config import configure_logging from app.temporal.worker import start_worker @@ -18,12 +21,33 @@ worker_running = False +@asynccontextmanager +async def lifespan(app: FastAPI): + await db.connect() + + # Reclaim records that were left stuck at in_progress by a previous + # crash / pod restart / OOM-kill before process_csv_inline could finish. + # Any record that has been in_progress for more than 30 minutes is + # considered stale and reset to pending so it can be retried. + try: + from app.database import operations as ops + reclaimed = await ops.reclaim_stale_in_progress(stale_minutes=30) + if reclaimed: + logger.info("Reclaimed %d stale in_progress CSV upload(s) back to pending on startup", reclaimed) + except Exception as exc: + logger.warning("Startup reclaim failed (non-fatal): %s", exc) + + yield + await db.disconnect() + + def run_web(): """Start the FastAPI web server.""" app = FastAPI( title="Analytics Service API Ingestion & Orchestration Layer", description="FastAPI ingestion endpoints and manual orchestration controls.", version="1.0.0", + lifespan=lifespan, ) register_exception_handlers(app) From a49401874a2fd07034bbb5c14355b81ac8554b6e Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 13:25:27 +0530 Subject: [PATCH 17/23] refactor(uploads): deduplicate CSV validation, remove dead code, and inline failure recording - Skip redundant validate_columns() execution in process_csv_inline when in-memory file_bytes are provided. - Remove unused list_by_status function from app/database/operations.py. - Inline operations.insert_upload_record directly within GCS upload failure handler in handle_upload. --- app/api/services/uploads.py | 29 +++++++++++++++-------------- app/database/operations.py | 14 -------------- 2 files changed, 15 insertions(+), 28 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 1a92657..04e1c0a 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -519,18 +519,19 @@ async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) await operations.update_status(record_id, "on_hold", error_meta) return - # --- 2. Validate columns --- - is_valid, errors = await asyncio.to_thread(validate_columns, df, report_type) - if not is_valid: - logger.warning("Validation failed for record %s: %s", record_id, errors) - error_meta = { - "stage": "CSV Column Validation", - "error": "Invalid CSV schema", - "validation_errors": errors, - "timestamp": datetime.utcnow().isoformat() + "Z", - } - await operations.update_status(record_id, "on_hold", error_meta) - return + # --- 2. Validate columns (only if loaded from GCS; handle_upload already validated file_bytes) --- + if file_bytes is None: + is_valid, errors = await asyncio.to_thread(validate_columns, df, report_type) + if not is_valid: + logger.warning("Validation failed for record %s: %s", record_id, errors) + error_meta = { + "stage": "CSV Column Validation", + "error": "Invalid CSV schema", + "validation_errors": errors, + "timestamp": datetime.utcnow().isoformat() + "Z", + } + await operations.update_status(record_id, "on_hold", error_meta) + return await operations.update_status(record_id, "in_progress") @@ -764,8 +765,8 @@ async def handle_upload( meta_data=meta_data, status="failed", ) - except Exception: - pass + except Exception as db_exc: + logger.warning("Failed to record upload failure in DB: %s", db_exc) raise RuntimeError(f"GCS Upload failed: {exc}. Please verify GCS settings.") record_id = await operations.insert_upload_record( diff --git a/app/database/operations.py b/app/database/operations.py index 567aff7..df65ef9 100644 --- a/app/database/operations.py +++ b/app/database/operations.py @@ -760,20 +760,6 @@ async def update_status( ) -async def list_by_status(status: str) -> list: - """List all tracker records with a given status.""" - from app.database.db import db - if not db.pool: - await db.connect() - - async with db.pool.acquire() as conn: - rows = await conn.fetch( - "SELECT * FROM csv_uploads WHERE status = $1 ORDER BY created_at", - status, - ) - return [dict(r) for r in rows] - - async def try_claim_for_processing(record_id: int) -> Optional[str]: """ Atomically set status to 'in_progress' if record exists and its status is not 'in_progress'. From 65ffe820888a5baaa4ab4fe39987329a05046b7c Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 14:17:19 +0530 Subject: [PATCH 18/23] revert back --- app/api/services/uploads.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 04e1c0a..74fb021 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -263,7 +263,7 @@ def row_to_json( original_pdf = get_url_field(get_csv_value(row_dict, expected_cols, pdf_col)) pdf_urls = None if original_pdf: - pdf_urls = {"original": original_pdf, "masked": original_pdf} + pdf_urls = {"original": original_pdf} tags = { "state": state, From c927723935d0bcfcce8556e960f9b00c4c66b09f Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 14:28:00 +0530 Subject: [PATCH 19/23] fix(uploads): prevent missing discussion date from defaulting to current time --- app/api/services/uploads.py | 6 ++++-- 1 file changed, 4 insertions(+), 2 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 74fb021..1a4a23d 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -156,8 +156,10 @@ def parse_segments(val, delimiter="|") -> List[str]: return segments -def format_datetime(val, with_ms=True) -> str: +def format_datetime(val, with_ms=True, fallback_to_now=True) -> Optional[str]: if pd.isna(val) or val is None: + if not fallback_to_now: + return None val = datetime.utcnow() if isinstance(val, str): try: @@ -257,7 +259,7 @@ def row_to_json( discussion_date = None if normalized_type == "discussion": discussion_date_raw = get_csv_value(row_dict, expected_cols, "Date of Discussion") - discussion_date = format_datetime(discussion_date_raw, with_ms=False) + discussion_date = format_datetime(discussion_date_raw, with_ms=False, fallback_to_now=False) pdf_col = "Pdf" if normalized_type == "story" else "PDF Urls" original_pdf = get_url_field(get_csv_value(row_dict, expected_cols, pdf_col)) From 5424d47dfffd35afaca71ef4fa9c4fa27c3a9b2e Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 14:29:33 +0530 Subject: [PATCH 20/23] fix(uploads): do not fabricate placeholder fallback strings for missing tenant, leader, or program metadata --- app/api/services/uploads.py | 16 ++++++++-------- 1 file changed, 8 insertions(+), 8 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 1a4a23d..5d9e44f 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -549,7 +549,7 @@ async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) if not isinstance(record_meta, dict): record_meta = {} - tenant_code = record_meta.get("tenant_code") or "mitra" + tenant_code = record_meta.get("tenant_code") try: async with db.pool.acquire() as conn: @@ -600,18 +600,18 @@ async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) except Exception as db_exc: logger.warning("Failed to query program/leader category metadata from DB: %s", db_exc) - # Fallbacks if DB query returned nothing - if not leader_info: + # Use record's values if DB query didn't find matching rows (do not fabricate placeholders) + if not leader_info and record.get("leader_category"): leader_info = { "id": str(uuid.uuid4()), - "name": record.get("leader_category") or "District Leader", - "description": f"Leader category: {record.get('leader_category') or 'District Leader'}", + "name": record.get("leader_category"), + "description": f"Leader category: {record.get('leader_category')}", } - if not program_info: + if not program_info and record.get("program_name"): program_info = { "id": str(uuid.uuid4()), - "name": record.get("program_name") or "My Program", - "description": f"Program: {record.get('program_name') or 'My Program'}", + "name": record.get("program_name"), + "description": f"Program: {record.get('program_name')}", } metadata = { From 0d293c7c332b9e487d143a08d0b1e0b0f3f17855 Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 14:32:22 +0530 Subject: [PATCH 21/23] fix(uploads): exclude failed status rows from check_duplicate_file query --- app/database/operations.py | 1 + 1 file changed, 1 insertion(+) diff --git a/app/database/operations.py b/app/database/operations.py index df65ef9..b8ee038 100644 --- a/app/database/operations.py +++ b/app/database/operations.py @@ -660,6 +660,7 @@ async def check_duplicate_file( AND report_type = $3 AND file_name = $4 AND file_size = $5 + AND status != 'failed' LIMIT 1 """, program_name, From 10d59df81858dd097da901c5a0ef7b27d517a83b Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 14:40:26 +0530 Subject: [PATCH 22/23] fix(uploads): add periodic background sweep loop for stale in_progress and stuck pending CSV uploads --- app/database/operations.py | 13 +++++++++++ main.py | 45 ++++++++++++++++++++++++++++---------- 2 files changed, 47 insertions(+), 11 deletions(-) diff --git a/app/database/operations.py b/app/database/operations.py index b8ee038..ddd4fc0 100644 --- a/app/database/operations.py +++ b/app/database/operations.py @@ -760,6 +760,19 @@ async def update_status( record_id, ) +async def list_by_status(status: str) -> list: + """List all tracker records with a given status.""" + from app.database.db import db + if not db.pool: + await db.connect() + + async with db.pool.acquire() as conn: + rows = await conn.fetch( + "SELECT * FROM csv_uploads WHERE status = $1 ORDER BY created_at", + status, + ) + return [dict(r) for r in rows] + async def try_claim_for_processing(record_id: int) -> Optional[str]: """ diff --git a/main.py b/main.py index 6f57044..5f01fbc 100644 --- a/main.py +++ b/main.py @@ -21,23 +21,46 @@ worker_running = False +async def _periodic_reclaim_loop(interval_seconds: int = 300): + from app.database import operations as ops + from app.api.services.uploads import process_csv_inline + + while True: + try: + reclaimed = await ops.reclaim_stale_in_progress(stale_minutes=30) + if reclaimed: + logger.info("Reclaimed %d stale in_progress CSV upload(s)", reclaimed) + + pending_records = await ops.list_by_status("pending") + for record in pending_records: + record_id = record["id"] + claim_status = await ops.try_claim_for_processing(record_id) + if claim_status == "success": + logger.info("Periodic sweep picked up pending CSV upload ID %s", record_id) + asyncio.create_task(process_csv_inline(record_id)) + except asyncio.CancelledError: + logger.info("Periodic CSV reclaim loop cancelled.") + break + except Exception as exc: + logger.exception("Periodic CSV reclaim loop encountered an error: %s", exc) + + await asyncio.sleep(interval_seconds) + + @asynccontextmanager async def lifespan(app: FastAPI): await db.connect() - # Reclaim records that were left stuck at in_progress by a previous - # crash / pod restart / OOM-kill before process_csv_inline could finish. - # Any record that has been in_progress for more than 30 minutes is - # considered stale and reset to pending so it can be retried. - try: - from app.database import operations as ops - reclaimed = await ops.reclaim_stale_in_progress(stale_minutes=30) - if reclaimed: - logger.info("Reclaimed %d stale in_progress CSV upload(s) back to pending on startup", reclaimed) - except Exception as exc: - logger.warning("Startup reclaim failed (non-fatal): %s", exc) + reclaim_task = asyncio.create_task(_periodic_reclaim_loop(interval_seconds=300)) yield + + reclaim_task.cancel() + try: + await reclaim_task + except asyncio.CancelledError: + pass + await db.disconnect() From c71366c315b55fb51cf35b70d08af6d88d64110a Mon Sep 17 00:00:00 2001 From: rohxnn Date: Wed, 9 Sep 2026 14:51:17 +0530 Subject: [PATCH 23/23] perf(uploads): hoist per-row schema parsing, reuse parsed DataFrame, and clean up dead code - Hoist expected_cols = json.loads(...) out of row_to_json per-row loop into rows_to_json batch scope. - Pass pre-parsed DataFrame df from handle_upload directly to process_csv_inline to avoid duplicate CSV parsing. - Remove dead normalized_sub_type assignment from app/database/operations.py. --- app/api/services/uploads.py | 75 +++++++++++++++++++++---------------- app/database/operations.py | 3 -- 2 files changed, 42 insertions(+), 36 deletions(-) diff --git a/app/api/services/uploads.py b/app/api/services/uploads.py index 5d9e44f..37d304a 100644 --- a/app/api/services/uploads.py +++ b/app/api/services/uploads.py @@ -185,15 +185,17 @@ def row_to_json( report_type: str, event_type: str = "create", metadata: Optional[dict] = None, + expected_cols: Optional[List[str]] = None, ) -> str: row_dict = {k: (None if pd.isna(v) else v) for k, v in row.to_dict().items()} normalized_type = report_type.lower().strip() - raw_cols = settings.STORY_CSV_COLUMN if normalized_type == "story" else settings.DISCUSSION_CSV_COLUMN - try: - expected_cols = json.loads(raw_cols) - except Exception: - expected_cols = [] + if expected_cols is None: + raw_cols = settings.STORY_CSV_COLUMN if normalized_type == "story" else settings.DISCUSSION_CSV_COLUMN + try: + expected_cols = json.loads(raw_cols) + except Exception: + expected_cols = [] try: submission_id = int(get_csv_value(row_dict, expected_cols, "id")) @@ -367,6 +369,13 @@ def rows_to_json( event_type: str = "create", metadata: Optional[dict] = None, ): + normalized_type = report_type.lower().strip() + raw_cols = settings.STORY_CSV_COLUMN if normalized_type == "story" else settings.DISCUSSION_CSV_COLUMN + try: + expected_cols = json.loads(raw_cols) + except Exception: + expected_cols = [] + for _, row in df.iterrows(): row_dict = {k: (None if pd.isna(v) else v) for k, v in row.to_dict().items()} is_complete, missing_fields = _is_row_complete(row_dict, report_type) @@ -376,7 +385,7 @@ def rows_to_json( missing_fields, ) continue - yield row_to_json(row, report_type, event_type, metadata) + yield row_to_json(row, report_type, event_type, metadata, expected_cols=expected_cols) # --------------------------------------------------------------------------- @@ -462,7 +471,11 @@ def _on_delivery(err, _msg): # Inline CSV Processing (replaces Temporal activities) # --------------------------------------------------------------------------- -async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) -> None: +async def process_csv_inline( + record_id: int, + file_bytes: Optional[bytes] = None, + df: Optional[pd.DataFrame] = None, +) -> None: """ Processes a single csv_upload record end-to-end: 1. Use in-memory CSV bytes (or fetch from cloud storage if file_bytes is None) @@ -483,33 +496,29 @@ async def process_csv_inline(record_id: int, file_bytes: Optional[bytes] = None) cloud_storage_path = record["cloud_storage_path"] report_type = record["report_type"] - # --- 1. Fetch/Parse CSV (use in-memory file_bytes if available, else fetch from storage) --- - # Retry GCS fetch up to 3 times with exponential backoff (1s, 2s, 4s) so that - # transient network blips self-heal instead of immediately going on_hold. - # In-memory file_bytes (fresh upload path) skips the retry since no network is - # involved. The old Temporal activity had retry_policy(maximum_attempts=3, - # backoff_coefficient=2.0) — this is the equivalent without Temporal. + # --- 1. Fetch/Parse CSV (use pre-parsed df/file_bytes if available, else fetch from storage) --- try: - if file_bytes is None: - last_exc: Exception = RuntimeError("unreachable") - for attempt in range(3): - try: - csv_file = await asyncio.to_thread(fetch_csv, cloud_storage_path) - break - except Exception as exc: - last_exc = exc - if attempt < 2: - wait = 2 ** attempt # 1s, 2s — then give up - logger.warning( - "GCS fetch attempt %d/3 failed for record %s (%s); retrying in %ds", - attempt + 1, record_id, exc, wait, - ) - await asyncio.sleep(wait) + if df is None: + if file_bytes is None: + last_exc: Exception = RuntimeError("unreachable") + for attempt in range(3): + try: + csv_file = await asyncio.to_thread(fetch_csv, cloud_storage_path) + break + except Exception as exc: + last_exc = exc + if attempt < 2: + wait = 2 ** attempt # 1s, 2s — then give up + logger.warning( + "GCS fetch attempt %d/3 failed for record %s (%s); retrying in %ds", + attempt + 1, record_id, exc, wait, + ) + await asyncio.sleep(wait) + else: + raise last_exc else: - raise last_exc - else: - csv_file = file_bytes - df = await asyncio.to_thread(load_csv, csv_file) + csv_file = file_bytes + df = await asyncio.to_thread(load_csv, csv_file) except Exception as exc: logger.exception("Failed to fetch/load CSV for record %s", record_id) error_meta = { @@ -789,7 +798,7 @@ async def handle_upload( claim_status = await operations.try_claim_for_processing(record_id) if claim_status == "success": - background_tasks.add_task(process_csv_inline, record_id, file_bytes) + background_tasks.add_task(process_csv_inline, record_id, file_bytes, df) logger.info("Scheduled inline CSV processing for upload ID %s", record_id) return { diff --git a/app/database/operations.py b/app/database/operations.py index ddd4fc0..e139091 100644 --- a/app/database/operations.py +++ b/app/database/operations.py @@ -230,9 +230,6 @@ async def insert_or_update_submission( # created at. This field is only present on create events (absent from # partial update payloads), so COALESCE in the SQL below preserves the # original DB value on updates — no special-casing needed. - # discussion_date (when the discussion took place) is a separate field - # (data.discussionDate) stored in discussion_submissions, handled below. - normalized_sub_type = submission_type.lower().strip() sub_date_str = data.get("submissionDate") if sub_date_str: