Files
query-orchestration/internal/query
Jay Brown 730f3ba11a Merged in feature/preserve-folder-paths (pull request #179)
support folder paths

* support folder paths

* s3 storage docs
2025-09-05 18:25:31 +00:00
..
2025-08-04 10:54:22 -07:00
2025-07-23 09:47:01 -07:00

Package query

Overview

Package query implements the core query processing engine for the document orchestration platform. It provides a sophisticated dependency-driven execution model where queries can extract data from documents and depend on results from other queries. The package supports multiple query types, versioning, and asynchronous processing through AWS SQS queues.

Architecture

Query Processing Model

The system implements a directed acyclic graph (DAG) of query dependencies:

  1. Independent Queries: Process document text directly
  2. Dependent Queries: Use results from other queries as input
  3. Version Management: Each query maintains version history
  4. Result Caching: Processed results stored for reuse

Query Types

JsonExtractor

Extracts values from JSON documents using JSONPath expressions:

  • Supports nested path extraction
  • Handles arrays and complex objects
  • Returns extracted values as strings
  • Configurable default values for missing paths

ContextFull

Processes complete document context:

  • Access to full document text
  • Can reference other query results
  • Supports complex transformations
  • Used for advanced extraction logic

Core Components

Service Layer (service.go)

Central service orchestrating all query operations:

  • Query creation and updates
  • Version management
  • Result processing coordination
  • Dependency resolution

Query Management

  • create.go: Query creation with validation
  • update.go: Query updates with version tracking
  • get.go: Query retrieval operations
  • list.go: Query listing with filtering
  • normalize.go: Configuration normalization
  • parse.go: Query parsing and validation

Result Processing (result/)

  • processor/: Query type-specific processors
  • get.go: Result retrieval
  • process.go: Result computation logic
  • set/: Result storage operations
  • sync/: Result synchronization

Synchronization

  • sync/: Query dependency synchronization
  • versionsync/: Version update propagation
  • test/: Query testing utilities

Query Lifecycle

1. Query Creation

// Create a new JSON extractor query
query, err := queryService.Create(ctx, &CreateRequest{
    ClientID:    clientID,
    Name:        "extract_invoice_number",
    Type:        "JsonExtractor",
    Description: "Extracts invoice number from document",
    Config: map[string]interface{}{
        "path": "$.invoice.number",
        "default": "N/A",
    },
})

2. Query Dependencies

// Create dependent query
dependentQuery, err := queryService.Create(ctx, &CreateRequest{
    ClientID:     clientID,
    Name:         "calculate_total",
    Type:         "ContextFull",
    DependsOn:    []string{"extract_line_items"},
    Description:  "Calculates total from line items",
    Config:       config,
})

3. Version Management

  • Each update creates a new version
  • Active version tracked separately
  • Previous versions retained for history
  • Dependency chains updated on version changes

4. Result Processing

// Process query for a document
result, err := resultProcessor.Process(ctx, &ProcessRequest{
    DocumentID: documentID,
    QueryID:    queryID,
})

Configuration

Query Configuration Structure

JsonExtractor Config

{
    "path": "$.data.field",
    "default": "default_value",
    "type": "string"
}

ContextFull Config

{
    "template": "Process {{.Document}} with {{.Dependencies.other_query}}",
    "format": "json",
    "validations": []
}

Service Configuration

  • Database connection for query storage
  • Queue configuration for async processing
  • Logging configuration
  • Thread pool for parallel processing

Query Processing Flow

Synchronous Processing

  1. Validation: Check query configuration
  2. Dependency Resolution: Ensure dependencies are satisfied
  3. Execution: Run query processor
  4. Storage: Save results to database
  5. Notification: Trigger dependent queries

Asynchronous Processing

  1. Queue Message: Receive processing request
  2. Lock Acquisition: Prevent concurrent processing
  3. Processing: Execute query logic
  4. Result Storage: Persist results
  5. Dependency Trigger: Queue dependent queries

Dependency Management

Dependency Resolution

  • Recursive dependency checking
  • Circular dependency prevention
  • Missing dependency handling
  • Version-aware dependency tracking

Dependency Types

  • Hard Dependencies: Must complete before processing
  • Soft Dependencies: Optional, use if available
  • Version Dependencies: Specific version requirements

Error Handling

Error Types

  • Validation Errors: Invalid configuration or input
  • Dependency Errors: Missing or failed dependencies
  • Processing Errors: Query execution failures
  • Storage Errors: Database operation failures

Error Recovery

  • Automatic retry for transient failures
  • Dead letter queue for persistent failures
  • Error context preservation
  • Detailed error logging

Testing

Unit Testing

  • Mock-based testing for isolated components
  • Configuration validation tests
  • Dependency resolution tests
  • Error scenario coverage

Integration Testing

  • Full query lifecycle tests
  • Dependency chain testing
  • Concurrent processing validation
  • Performance benchmarks

Test Utilities (test/)

  • Query creation helpers
  • Result validation utilities
  • Mock data generators
  • Test scenario builders

Performance Optimization

Caching Strategy

  • Result caching for repeated queries
  • Dependency result reuse
  • Configuration caching
  • Connection pooling

Parallel Processing

  • Concurrent query execution where possible
  • Thread pool management
  • Resource throttling
  • Queue-based load distribution

Database Optimization

  • Indexed query lookups
  • Batch result storage
  • Efficient dependency queries
  • View-based aggregations