Merged in feature/remove-query (pull request #201)
remove query from codebase part 1 * remove query * fix localstack run
This commit is contained in:
+11
-69
@@ -5,7 +5,7 @@
|
||||
This document describes the SQS queue-based processing pipeline that orchestrates document processing through various stages. The system uses AWS SQS queues to connect microservices (runners) in a sequential processing pipeline.
|
||||
|
||||
## Visual Pipeline Flow
|
||||
[mermaid link](https://mermaid.live/edit#pako:eNqdlWtv2zYUhv8KwaDABiiZJOqOoUDjdFg_dGvktAVmFwVDHV0QmdQoqY1n57-PF8nxXBcYKoCASL08z7mR2mEmCsAZLlvxldVUDujues2Rel68QEuCXn8BPqClGCUDu74kq5_Uh-uRPcDw67385eWNYONGy953raDFz5_Q5eXL_WHzH2JoyobRoRF8j5a33pofCG9pw9G7poO24YBuRxihR5QXKB85B9lPyFtv1Q9CwmfQFg3UaA0J5d7OfDU4u9FIssRNvCdrIveMUwsJdFCIQjAkgQlZaI_8VTGF8LnhzbH92VWlsSh_p6RvlOiE488c33JqYA-oGLtWBw690d1v0es7WqGa9rXGkmdsv-XsLJZYLNHYpRKdYMmMNbr9B9o2heIh1jZzohjlxjwqW1ppbPCMZS1QfpYbWG6guQutOgEHMziw8WoJ6iYBE20LTJUE9YMqJpVFr8HhM3iAx_NpDi031Nw7JTrBhjPW6PavPi6RVknKrLk_Fzma5lO3Rau_R5Db72c4ssRoZ3RnchzPzMj2NauhGFtAWt9MpS1VrHNwmhpb6llgbIGxBZ7AkhlmVPvfVfomlDVWQAe8AM4U2YCODtPCFB3pENBv6jzbD6_evfFWt9qAfj2A3E9T4cyesdNtow0mK9s73ybMep9Y75OdlZ3JVzqHYKT7azqwGnVSMOhtsmjbTg2qc2bCIEdhWF8_qNOvSnguGmPWqp79Tqc6f7H7vud-at1PbfInyLdBpO4chNmwv5NNVYFEtCxVY0Nh29zEYPxPjvxfDlt1m1V2zlra9zdQ6hKOgMqmbbOLsixi13X6QYoHyC4IIdP75demGOrM7x4ddYaEzC5c1z0xJI2jk6UkZgD3P2hJX5q0mp1Kk_I-TX7QFO2aQ2zUjen_N3NkSN_zjrpp1SBqBGqEakRqxGokji6fyePxptxzct_JiZMHTh46eeTksZMnjqqcTdV_CGSO-njVNJUKwa5hB29AbmhTqP_iTq-t8VDDBtY4U6_qOntY4zV_Ujo6DkI3D84GOYKDpRirGmclbXs1s81509BK0s0s6Sj_S4jDFIpG-fPW_oTNv9jBldToyaI-7XIhRj7gzPdiz1jA2Q4_4ozEwVVCUt91Yz8KQ8938BZnl0mcXBEv8RMSxF7ok-DJwf8YpneVem4QpUkUpcSPI-_pXxutlHc)
|
||||
|
||||
```mermaid
|
||||
flowchart TB
|
||||
%% S3 Event Source
|
||||
@@ -25,46 +25,34 @@ flowchart TB
|
||||
R4 -->|Clean per<br/>collector standards| SQ5[document_text<br/>Queue]
|
||||
|
||||
SQ5 --> R5{docTextRunner<br/>:8085}
|
||||
R5 -->|AWS Textract<br/>OCR extraction| SQ6[query_sync<br/>Queue]
|
||||
|
||||
SQ6 --> R6{querySyncRunner<br/>:8087}
|
||||
R6 -->|Schedule queries<br/>for document| SQ7[query<br/>Queue]
|
||||
|
||||
SQ7 --> R7{queryRunner<br/>:8088}
|
||||
R7 -->|Handle query<br/>dependencies| SQ7
|
||||
R5 -->|AWS Textract<br/>OCR extraction| DB[(PostgreSQL)]
|
||||
|
||||
%% Client Sync Flow
|
||||
API1[Query API<br/>:8080] -->|Client update| SQ8[client_sync<br/>Queue]
|
||||
SQ8 --> R8{clientSyncRunner<br/>:8089}
|
||||
R8 -->|Batch process<br/>all client docs| SQ3
|
||||
|
||||
%% Query Version Sync Flow
|
||||
API1 -->|Query update| SQ9[query_version_sync<br/>Queue]
|
||||
SQ9 --> R9{queryVersionSyncRunner<br/>:8090}
|
||||
R9 -->|Trigger affected<br/>clients| SQ8
|
||||
|
||||
%% Styling
|
||||
classDef queue fill:#ffd700,stroke:#333,stroke-width:2px,color:#000
|
||||
classDef runner fill:#87ceeb,stroke:#333,stroke-width:2px,color:#000
|
||||
classDef storage fill:#98fb98,stroke:#333,stroke-width:2px,color:#000
|
||||
classDef api fill:#ffa07a,stroke:#333,stroke-width:2px,color:#000
|
||||
|
||||
class SQ1,SQ2,SQ3,SQ4,SQ5,SQ6,SQ7,SQ8,SQ9 queue
|
||||
class R1,R2,R3,R4,R5,R6,R7,R8,R9 runner
|
||||
class S3 storage
|
||||
class SQ1,SQ2,SQ3,SQ4,SQ5,SQ8 queue
|
||||
class R1,R2,R3,R4,R5,R8 runner
|
||||
class S3,DB storage
|
||||
class API1 api
|
||||
```
|
||||
|
||||
### Flow Legend
|
||||
- 🟨 **Yellow boxes**: SQS Queues
|
||||
- 🟦 **Blue diamonds**: Runner Services (with port numbers)
|
||||
- 🟩 **Green cylinder**: S3 Storage
|
||||
- 🟩 **Green cylinder**: S3 Storage / PostgreSQL Database
|
||||
- 🟧 **Orange box**: REST API endpoint
|
||||
- **Solid arrows**: Message flow direction
|
||||
- **Self-loop on queryRunner**: Handles query dependencies by re-queueing
|
||||
|
||||
## Detailed Process Flow
|
||||
[mermaid link](https://mermaid.live/edit#pako:eNp9lW9v2jAQxr-K5VfblFaEtCvkRaU1MAlp0xoSadKEVBnnAI_EoY4zlSK--84OgQYSkJAc--fnnjv_21GeJ0B9OpMFvJYgOYwEWyqWzSTB34YpLbjYMKlJ5LX0jQkrSKFzBS_wD7AnLKGENnB6JMcGnJZSgroERxPDJTkvM6RehBSdmqPJ9MBOkOoUjBqCxVbyTsGoFoyQ6hQMGoI8BSY7FYNaMTBYp2TckNTw1p10XCvGSHUJhjZpXFC1vZpxWGVswWs5h0eqU-kk1Jnlk0Ge80IvFUThj5msmMi7eXyMxj42iN0b5NOv-V_gOlDANCT-l88HcGzBqY8iaUoyKAq2hHpsioOjJ58EK-BrIhZEZJtcabSgVw1mgpFAJqaGRMhFTnbzkq9BO2QNW4esWLHa195GEzujNSR2N0Im5SYVHB2T-daqXGA2n-M6EwU8V0mDig7edjU0GZ28RJZo9xI1vfBUmAicSbv-DShoiWHqBalYinkKx3iBpdvjBRflZimml2yJPRDn1DOoRa6yalDIZWM8vpZ0bIl2E_GFCXt08K8Yx51zwmqRgKHIt98RiQ_QuVJkbqiTghX8yITXViiMLNFqNjyu0HeB09nGbhYstz00AoqPWHgZxKkOl1krrCQBxlfHsGZGe9CzHQobVMVb_hSuBsZvwEtdmdmejVU1UVCUqT7Nq2xO4ebV3Aim-B_liWniQlOHZqAyJhJ8ZXZm8ozqFWQwoz42E6bWM3x99sixUufmEqK-ViU4VOXlckX9BUsL_Co3CR6ew-t07MWL5U-eZ_UUSAR6_Vm9afZpc-hSmdgHSeNPBXkpNfX77p0VoP6OvlHf7bm3D3f9njcY9j13eN936Nb09m693sD1BgPPux_2eoO9Q99tyN7tYOC6_Ydhf-h-7WPjYf8f7gRf3A)
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
participant S3
|
||||
@@ -78,10 +66,6 @@ sequenceDiagram
|
||||
participant DCR as docCleanRunner
|
||||
participant DT as document_text Queue
|
||||
participant DTR as docTextRunner
|
||||
participant QS as query_sync Queue
|
||||
participant QSR as querySyncRunner
|
||||
participant Q as query Queue
|
||||
participant QR as queryRunner
|
||||
participant DB as PostgreSQL
|
||||
|
||||
S3->>SE: S3 Event (ObjectCreated:*)
|
||||
@@ -107,17 +91,6 @@ sequenceDiagram
|
||||
DTR->>DB: Check if text extracted
|
||||
DTR-->>DTR: Call AWS Textract
|
||||
DTR->>DB: Store extracted text
|
||||
DTR->>QS: Send {documentID}
|
||||
|
||||
QS->>QSR: Poll message
|
||||
QSR->>DB: Find applicable queries
|
||||
QSR->>Q: Send {documentID, queryID} for each
|
||||
|
||||
Q->>QR: Poll message
|
||||
QR->>DB: Check dependencies
|
||||
QR->>DB: Execute query
|
||||
QR->>DB: Store results
|
||||
QR-->>Q: Re-queue if dependencies pending
|
||||
```
|
||||
|
||||
## S3 Event Trigger Configuration
|
||||
@@ -171,9 +144,7 @@ aws s3api put-bucket-notification-configuration \
|
||||
| **2** | `docInitRunner` | `document_init` | `document_sync` | Creates document records, checks for duplicates based on ETag hash |
|
||||
| **3** | `docSyncRunner` | `document_sync` | `document_clean` | Validates client eligibility for document processing (can_sync flag) |
|
||||
| **4** | `docCleanRunner` | `document_clean` | `document_text` | Cleans document content according to collector standards |
|
||||
| **5** | `docTextRunner` | `document_text` | `query_sync` | Extracts text using AWS Textract OCR service |
|
||||
| **6** | `querySyncRunner` | `query_sync` | `query` | Syncs and schedules queries that need to run for documents |
|
||||
| **7** | `queryRunner` | `query` | `query` (self-referential) | Executes queries with dependency chains, can re-queue for dependencies |
|
||||
| **5** | `docTextRunner` | `document_text` | - | Extracts text using AWS Textract OCR service, stores in database |
|
||||
|
||||
### Additional Async Runners
|
||||
|
||||
@@ -182,7 +153,6 @@ These runners are triggered by API calls or other events outside the main pipeli
|
||||
| Runner Service | Reads From Queue | Writes To Queue | Trigger |
|
||||
|---------------|------------------|-----------------|---------|
|
||||
| `clientSyncRunner` | `client_sync` | `document_sync` | Triggered by API calls or when client configuration changes |
|
||||
| `queryVersionSyncRunner` | `query_version_sync` | `client_sync` | Triggered when query versions are updated via API |
|
||||
|
||||
## Pipeline Flow Diagrams
|
||||
|
||||
@@ -210,13 +180,7 @@ docCleanRunner
|
||||
↓
|
||||
docTextRunner
|
||||
↓
|
||||
[query_sync queue]
|
||||
↓
|
||||
querySyncRunner
|
||||
↓
|
||||
[query queue]
|
||||
↓
|
||||
queryRunner
|
||||
PostgreSQL (text stored)
|
||||
```
|
||||
|
||||
### 2. Client Sync Flow (Batch Processing)
|
||||
@@ -230,21 +194,6 @@ clientSyncRunner
|
||||
[document_sync queue] → (feeds into main pipeline)
|
||||
```
|
||||
|
||||
### 3. Query Version Update Flow
|
||||
```
|
||||
Query Update (API)
|
||||
↓
|
||||
[query_version_sync queue]
|
||||
↓
|
||||
queryVersionSyncRunner
|
||||
↓
|
||||
[client_sync queue]
|
||||
↓
|
||||
clientSyncRunner
|
||||
↓
|
||||
[document_sync queue] → (feeds into main pipeline)
|
||||
```
|
||||
|
||||
## Queue Names and Environment Variables
|
||||
|
||||
| Queue Purpose | Queue Name Variable | Queue URL Variable |
|
||||
@@ -254,19 +203,15 @@ clientSyncRunner
|
||||
| Document Sync | `QNAME_DOCUMENT_SYNC` | `DOCUMENT_SYNC_URL` |
|
||||
| Document Clean | `QNAME_DOCUMENT_CLEAN` | `DOCUMENT_CLEAN_URL` |
|
||||
| Document Text | `QNAME_DOCUMENT_TEXT` | `DOCUMENT_TEXT_URL` |
|
||||
| Query Sync | `QNAME_QUERY_SYNC` | `QUERY_SYNC_URL` |
|
||||
| Query Runner | `QNAME_QUERY_RUNNER` | `QUERY_URL` |
|
||||
| Client Sync | `QNAME_CLIENT_SYNC` | `CLIENT_SYNC_URL` |
|
||||
| Query Version Sync | `QNAME_QUERY_VERSION_SYNC` | `QUERY_VERSION_SYNC_URL` |
|
||||
|
||||
## Key Design Principles
|
||||
|
||||
1. **Single Responsibility**: Each runner has one specific task in the pipeline
|
||||
2. **Sequential Processing**: Documents flow through defined stages in order
|
||||
3. **Idempotent Operations**: Duplicate detection prevents reprocessing
|
||||
4. **Async Reprocessing**: Client/query updates trigger document reprocessing
|
||||
5. **Self-Referential Queuing**: Query runner can re-queue items for dependency handling
|
||||
6. **Scalability**: Each runner can be scaled independently based on workload
|
||||
4. **Async Reprocessing**: Client updates trigger document reprocessing
|
||||
5. **Scalability**: Each runner can be scaled independently based on workload
|
||||
|
||||
## IAM Permissions Required
|
||||
|
||||
@@ -293,7 +238,4 @@ Runner ports for local development:
|
||||
- `docSyncRunner`: 8083
|
||||
- `docCleanRunner`: 8084
|
||||
- `docTextRunner`: 8085
|
||||
- `querySyncRunner`: 8087
|
||||
- `queryRunner`: 8088
|
||||
- `clientSyncRunner`: 8089
|
||||
- `queryVersionSyncRunner`: 8090
|
||||
Reference in New Issue
Block a user