Update complete pdf code for entire system * docs and code for generation of updated docs
13 KiB
Service Processing Details
This document provides detailed information about what each service in the document processing pipeline actually does.
Document Processing Pipeline Services
1. storeEventRunner
Purpose: Entry point for document processing, triggered by S3 upload events
Code Locations:
cmd/storeEventRunner/- Service entry point and main.goapi/storeEventRunner/- Runner implementation and message handlinginternal/document/store/- Core document storage logicinternal/serviceconfig/queue/documentinit/- Queue configuration for next stage
Processing Logic:
- Receives S3 event notifications (ObjectCreated:Put, ObjectCreated:CompleteMultipartUpload)
- Validates that the event type is supported
- Parses the S3 object key to extract structured information (client ID, location, timestamps)
- Filters for objects in the "import" location only
- Extracts document metadata (bucket, key, hash/ETag)
- Sends message to DOCINIT queue with document location information
Key Validations:
- Only processes supported S3 events (Put and CompleteMultipartUpload)
- Only processes objects with "import" location in their key path
- Validates key can be parsed into expected structure
Output: Sends to DOCINIT queue with {bucket, key, hash}
2. docInitRunner
Purpose: Document registration and deduplication
Code Locations:
cmd/docInitRunner/- Service entry point and main.goapi/docInitRunner/- Runner implementation and message handlinginternal/document/init/- Core document initialization logicinternal/serviceconfig/queue/documentsync/- Queue configuration for next stage
Processing Logic:
- Parses the S3 key to extract client ID, batch ID, and other metadata
- Checks if document already exists by looking up hash in database for the client
- If document exists (duplicate):
- Adds new document entry to existing document record
- Does NOT trigger downstream processing (deduplication)
- If new document:
- Creates new document record in database with client ID and hash
- Associates batch ID if present
- Adds document entry with S3 location details
- Sends message to DOCSYNC queue to continue processing
Key Features:
- Hash-based deduplication per client
- Maintains multiple entries for same document (different uploads)
- Transaction-based database operations for consistency
Output: Sends to DOCSYNC queue with {document_id} (only for new documents)
3. docSyncRunner
Purpose: Client synchronization validation
Code Locations:
cmd/docSyncRunner/- Service entry point and main.goapi/docSyncRunner/- Runner implementation and message handlinginternal/document/sync/- Core document synchronization logicinternal/client/- Client management utilitiesinternal/document/- Shared document utilitiesinternal/serviceconfig/queue/documentclean/- Queue configuration for next stage
Processing Logic:
- Retrieves document summary from database
- Fetches client configuration
- Checks if client has "CanSync" flag enabled
- If CanSync is true:
- Sends document to DOCCLEAN queue for processing
- If CanSync is false:
- Logs that document won't be synced
- Processing stops here for this document
Key Validations:
- Client must exist in database
- Client must have CanSync=true to proceed
Output: Sends to DOCCLEAN queue with {document_id} (only if client can sync)
4. docCleanRunner
Purpose: Document validation and preparation
Code Locations:
cmd/docCleanRunner/- Service entry point and main.goapi/docCleanRunner/- Runner implementation and message handlinginternal/document/clean/- Core document cleaning and validation logicinternal/document/types/- Document type definitions and MIME type handlinginternal/serviceconfig/queue/documenttext/- Queue configuration for next stage
Processing Logic:
- Checks if document has already been cleaned (idempotency)
- If not cleaned:
- Retrieves document entry from database
- Makes HEAD request to S3 to get content metadata
- Validates content type/MIME type is acceptable (PDF)
- Downloads document from S3
- Validates document integrity (not corrupt)
- Stores clean record with:
- Success: bucket, key, hash, mimetype
- Failure: reason for failure (invalid MIME, corrupt file)
- Creates versioned clean entry for tracking
- Sends to DOCTEXT queue regardless of success/failure
Key Validations:
- MIME type validation (must be acceptable format)
- Document corruption check
- Hash verification against S3 ETag
Output: Sends to DOCTEXT queue with {document_id}
5. docTextRunner
Purpose: Text extraction from documents using AWS Textract
Code Locations:
cmd/docTextRunner/- Service entry point and main.goapi/docTextRunner/- Runner implementation and message handlinginternal/document/text/- Core text extraction logic using AWS Textractinternal/document/types/- Document type utilities for text processinginternal/serviceconfig/textract/- AWS Textract client configurationinternal/serviceconfig/queue/querysync/- Queue configuration for next stage
Tests for this use mocks only. No tests to the actual aws textract service.
Processing Logic:
- Checks if text has already been extracted (idempotency)
- If not extracted:
- Retrieves clean document record
- Skips if document failed cleaning
- Downloads document from S3
- Determines page count
- For each page (in parallel using thread pool):
- Sends page to AWS Textract for OCR
- Processes base text blocks
- Identifies required features (tables, forms)
- Performs additional extraction if needed
- Merges text with page markers
- Stores extracted text in S3
- Records text extraction in database with versioning
- Sends to QUERYSYNC queue for query processing
Key Features:
- Parallel page processing for performance
- Supports tables and forms extraction
- Stores full text with page boundaries preserved
- Version tracking for text extractions
Output: Sends to QUERYSYNC queue with {document_id}
6. querySyncRunner
Purpose: Query dependency resolution and scheduling
Code Locations:
cmd/querySyncRunner/- Service entry point and main.goapi/querySyncRunner/- Runner implementation and message handlinginternal/query/sync/- Core query synchronization logicinternal/query/result/sync/- Result synchronization utilitiesinternal/serviceconfig/queue/query/- Queue configuration for next stage
Processing Logic:
- Verifies document has extracted text
- Queries database for unsynced queries that:
- Are configured for this client (via collector)
- Have no dependencies OR
- Have all dependencies already satisfied
- For each eligible query:
- Sends message to QUERY queue for execution
Key Features:
- Dependency graph resolution
- Only schedules queries with satisfied dependencies
- Client-specific query configuration via collectors
Output: Sends to QUERY queue with {document_id, query_id} for each eligible query
7. queryRunner
Purpose: Query execution and result storage
Code Locations:
cmd/queryRunner/- Service entry point and main.goapi/queryRunner/- Runner implementation and message handlinginternal/query/result/set/- Core query result processinginternal/query/result/- Query execution logic and processorsinternal/query/result/processor/- Base processor interface and utilitiesinternal/query/types/contextFull/- CONTEXT_FULL query type implementationinternal/query/types/jsonExtractor/- JSON_EXTRACTOR query type implementationinternal/query/- Query management utilitiesinternal/query/result/sync/- Result synchronization for triggering dependents
Processing Logic:
- Retrieves document text version information
- Fetches query definition and active version
- If query has dependencies:
- Retrieves results from required queries
- Validates all dependencies are satisfied
- Executes query based on type:
- CONTEXT_FULL: Returns complete document text
- JSON_EXTRACTOR: Extracts specific JSON paths from dependent results
- Stores result in database with:
- Query result value
- Query version used
- Text version processed
- Dependencies tracked
- Identifies queries that depend on this result
- Sends dependent queries to QUERY queue (recursive processing)
Key Features:
- Supports multiple query types with different processors
- Dependency tracking and validation
- Recursive processing for dependent queries
- Version tracking for reproducibility
Output: Sends to QUERY queue with {document_id, query_id} for each dependent query
Additional Synchronization Services
8. clientSyncRunner
Purpose: Client-level document synchronization
Code Locations:
cmd/clientSyncRunner/- Service entry point and main.goapi/clientSyncRunner/- Runner implementation and message handlinginternal/client/sync/- Core client synchronization logicinternal/serviceconfig/queue/documentsync/- Queue configuration for document sync
Processing Logic:
- Processes CLIENTSYNC queue messages
- Retrieves all documents for a specific client
- Sends each document to DOCSYNC queue for processing
- Enables bulk reprocessing of client documents
Output: Sends to DOCSYNC queue with {document_id} for each client document
9. queryVersionSyncRunner
Purpose: Query version change propagation
Code Locations:
cmd/queryVersionSyncRunner/- Service entry point and main.goapi/queryVersionSyncRunner/- Runner implementation and message handlinginternal/query/versionsync/- Core query version synchronization logicinternal/serviceconfig/queue/clientsync/- Queue configuration for client sync
Processing Logic:
- Processes QUERYVERSIONSYNC queue messages triggered by query updates
- Identifies all clients affected by query version changes
- Sends affected clients to CLIENTSYNC queue for reprocessing
- Enables cascading updates when query logic changes
Output: Sends to CLIENTSYNC queue with {client_id} for each affected client
Data Flow Summary
Primary Document Processing Pipeline
S3 Upload Event
↓
storeEventRunner (validates event, extracts metadata)
↓
docInitRunner (deduplication, document registration)
↓
docSyncRunner (client sync validation)
↓
docCleanRunner (document validation, MIME check)
↓
docTextRunner (OCR via AWS Textract)
↓
querySyncRunner (find executable queries)
↓
queryRunner (execute queries, trigger dependents)
↓
queryRunner (recursive for dependent queries)
Synchronization Flows
Query Version Update
↓
queryVersionSyncRunner (find affected clients)
↓
clientSyncRunner (reprocess all client documents)
↓
docSyncRunner → ... (normal processing pipeline)
Key Design Patterns
Idempotency
- Each service checks if its work has already been done
- Prevents duplicate processing on retries
- Version tracking ensures consistency
Dependency Management
- Queries can depend on other query results
- Dependency graph is resolved before execution
- Results are chained through recursive processing
Version Tracking
- Document clean versions
- Text extraction versions
- Query versions
- Enables reprocessing when configurations change
Error Handling
- Failed documents are marked but continue through pipeline
- Allows partial processing and debugging
- Failure reasons are stored for analysis
Client Isolation
- All processing is scoped to client ID
- Deduplication is per-client
- Query configurations are per-client via collectors
Shared Code Directories
These directories contain utilities and configurations used across multiple services:
Core Infrastructure
internal/database/- Database queries, migrations, and repository patternsinternal/serviceconfig/- Configuration management for all servicesinternal/serviceconfig/queue/- SQS queue client configurationsinternal/serviceconfig/objectstore/- S3 client configurationsinternal/serviceconfig/aws/- AWS service configurationsinternal/server/runner/- Base runner framework for queue-based services
Business Logic Libraries
internal/document/- General document utilities (shared across document services)internal/client/- Client management utilitiesinternal/query/- Query management and utilities (shared across query services)
Testing & Utilities
internal/test/- Test utilities and helpersinternal/validation/- Input validation utilitiesmocks/- Auto-generated mock interfaces for testing
API Layer
api/queryAPI/- REST API handlers and controllers (not a queue-based runner)pkg/queryAPI/- Generated API client libraries