goworker / internal /processor /processor.go
ChartDB Admin
feat: optimize MongoDB storage with 10-day TTL and minimal schema
ad3c586
Raw
History Blame Contribute Delete
20.2 kB
// Package processor provides the main processing logic for GDELT timestamps.
// It uses parallel processing with bounded concurrency, structural validation,
// and existential validation via HEAD probing.
package processor
import (
"context"
"crypto/sha256"
"encoding/hex"
"fmt"
"io"
"net/url"
"strings"
"sync"
"time"
"gdelt-engine/internal/constants"
"gdelt-engine/internal/downloader"
"gdelt-engine/internal/parser"
"gdelt-engine/internal/schema"
"go.uber.org/zap"
)
// Storage defines the interface for database operations.
type Storage interface {
IsTimestampCompleted(ctx context.Context, timestamp string) (bool, error)
GetTimestampStatus(ctx context.Context, timestamp string) (*schema.CompletedTimestamp, error)
MarkTimestampStarted(ctx context.Context, timestamp string) error
MarkTimestampCompleted(ctx context.Context, ts *schema.CompletedTimestamp) error
GetAllTimestamps(ctx context.Context) (completed []schema.CompletedTimestamp, processing []string, err error)
MarkURLProcessed(ctx context.Context, url *schema.ProcessedURL) error
GetProcessedURLs(ctx context.Context, timestamp string) ([]schema.ProcessedURL, error)
BulkInsertArticles(ctx context.Context, articles []schema.Article) (int, error)
}
// ValidationResult represents the result of timestamp validation
type ValidationResult struct {
Timestamp string
Valid bool
Reason string
Status schema.ProcessingStatus
}
// Task represents a single file processing task
type Task struct {
Timestamp string
URL string
FileType string
}
// Result represents the result of processing a single file
type Result struct {
Task Task
RowCount int
ArticleCount int
Duration time.Duration
Err error
}
// ProcessorOption is a functional option for Processor
type ProcessorOption func(*Processor)
func WithWorkers(n int) ProcessorOption {
return func(p *Processor) { p.workers = n }
}
func WithTimeout(t time.Duration) ProcessorOption {
return func(p *Processor) { p.timeout = t }
}
func WithBatchSize(n int) ProcessorOption {
return func(p *Processor) { p.batchSize = n }
}
func WithMaxParallelTimestamps(n int) ProcessorOption {
return func(p *Processor) { p.maxParallelTimestamps = n }
}
func WithFilterThemes(themes []string) ProcessorOption {
return func(p *Processor) { p.filterThemes = themes }
}
func WithFilterLocations(locations []string) ProcessorOption {
return func(p *Processor) { p.filterLocations = locations }
}
// Processor handles concurrent timestamp processing
type Processor struct {
workers int
timeout time.Duration
batchSize int
maxParallelTimestamps int
downloader *downloader.StreamingDownloader
parser parser.StreamParser
storage Storage
logger *zap.Logger
// Filters - if set, only articles matching these will be saved
filterThemes []string // e.g., ["ECON_", "TAX_"]
filterLocations []string // e.g., ["US", "IN"]
// Semaphore for limiting parallel timestamps
semaphore chan struct{}
// Track processing state
mu sync.RWMutex
processing map[string]bool
}
// NewProcessor creates a new Processor with functional options
func NewProcessor(storage Storage, logger *zap.Logger, opts ...ProcessorOption) *Processor {
p := &Processor{
workers: constants.DefaultWorkers,
timeout: constants.DefaultTimeout,
batchSize: constants.DefaultBatchSize,
maxParallelTimestamps: constants.MaxParallelTimestamps,
storage: storage,
logger: logger,
processing: make(map[string]bool),
}
for _, opt := range opts {
opt(p)
}
// Initialize semaphore for parallel timestamp limit
p.semaphore = make(chan struct{}, p.maxParallelTimestamps)
// Initialize downloader and parser
p.downloader = downloader.NewStreamingDownloader(
downloader.WithLogger(logger),
)
p.parser = parser.NewCSVStreamParser(
parser.WithParserLogger(logger),
)
return p
}
// ValidateTimestampStructural performs structural validation (cheap, deterministic).
// Checks: length, parsable, minute โˆˆ {00,15,30,45}, not in future.
func ValidateTimestampStructural(ts string) *ValidationResult {
result := &ValidationResult{Timestamp: ts}
// Check length
if len(ts) != constants.TimestampLength {
result.Valid = false
result.Reason = fmt.Sprintf("invalid length: expected %d, got %d", constants.TimestampLength, len(ts))
result.Status = schema.StatusFailed
return result
}
// Parse as time
t, err := time.Parse(constants.TimestampFormat, ts)
if err != nil {
result.Valid = false
result.Reason = fmt.Sprintf("invalid format: %v", err)
result.Status = schema.StatusFailed
return result
}
// Check minute is valid GDELT interval (00, 15, 30, 45)
m := t.Minute()
if m%15 != 0 {
result.Valid = false
result.Reason = fmt.Sprintf("invalid minute %d: must be 00, 15, 30, or 45", m)
result.Status = schema.StatusFailed
return result
}
// Check not in future
if t.After(time.Now().UTC()) {
result.Valid = false
result.Reason = "timestamp is in the future"
result.Status = schema.StatusFailed
return result
}
result.Valid = true
return result
}
// isWithinGraceWindow checks if timestamp is within the grace window for recent data
func isWithinGraceWindow(ts string) bool {
t, err := time.Parse(constants.TimestampFormat, ts)
if err != nil {
return false
}
graceTime := time.Now().UTC().Add(-time.Duration(constants.GraceWindowMinutes) * time.Minute)
return t.After(graceTime)
}
// QueueTimestamps validates and queues timestamps for processing.
// Returns per-timestamp validation results.
func (p *Processor) QueueTimestamps(ctx context.Context, timestamps []string) (queued []string, rejected []ValidationResult) {
p.mu.Lock()
defer p.mu.Unlock()
for _, ts := range timestamps {
// Layer 1: Structural validation
validation := ValidateTimestampStructural(ts)
if !validation.Valid {
rejected = append(rejected, *validation)
p.logger.Debug("Rejected invalid timestamp",
zap.String("timestamp", ts),
zap.String("reason", validation.Reason),
)
continue
}
// Check if already processing
if p.processing[ts] {
rejected = append(rejected, ValidationResult{
Timestamp: ts,
Valid: false,
Reason: "already processing",
Status: schema.StatusProcessing,
})
continue
}
// Check if already completed
completed, err := p.storage.IsTimestampCompleted(ctx, ts)
if err != nil {
p.logger.Error("Failed to check timestamp status", zap.String("timestamp", ts), zap.Error(err))
continue
}
if completed {
rejected = append(rejected, ValidationResult{
Timestamp: ts,
Valid: false,
Reason: "already completed",
Status: schema.StatusCompleted,
})
continue
}
// Mark as processing
p.processing[ts] = true
queued = append(queued, ts)
}
// Start background processing for queued timestamps
if len(queued) > 0 {
go p.processTimestampsParallel(context.Background(), queued)
}
return queued, rejected
}
// processTimestampsParallel processes multiple timestamps in parallel with bounded concurrency.
func (p *Processor) processTimestampsParallel(ctx context.Context, timestamps []string) {
var wg sync.WaitGroup
for _, ts := range timestamps {
wg.Add(1)
go func(timestamp string) {
defer wg.Done()
// Acquire semaphore (limit parallel timestamps)
p.semaphore <- struct{}{}
defer func() { <-p.semaphore }()
p.processTimestamp(ctx, timestamp)
}(ts)
}
wg.Wait()
}
// processTimestamp processes a single timestamp with existential validation
func (p *Processor) processTimestamp(ctx context.Context, timestamp string) {
start := time.Now()
p.logger.Info("Processing timestamp", zap.String("timestamp", timestamp))
// Mark as started in database
if err := p.storage.MarkTimestampStarted(ctx, timestamp); err != nil {
p.logger.Error("Failed to mark timestamp started", zap.Error(err))
}
// Create timeout context
ctx, cancel := context.WithTimeout(ctx, p.timeout)
defer cancel()
// Generate tasks for all 3 files
tasks := p.generateTasks(timestamp)
// Layer 2: Existential validation via HEAD probing
// Check if at least the export file exists
exportURL := fmt.Sprintf(constants.ExportURLTemplate, timestamp)
headResult, err := p.downloader.HeadCheck(ctx, exportURL)
if headResult == downloader.HeadNotFound {
// Check if within grace window
if isWithinGraceWindow(timestamp) {
p.logger.Info("Timestamp not yet published, within grace window",
zap.String("timestamp", timestamp),
)
p.markTimestampResult(ctx, timestamp, start, schema.StatusPending,
"GDELT file not yet published (within grace window)", 0, 0, 0)
} else {
p.logger.Warn("Timestamp does not exist on GDELT",
zap.String("timestamp", timestamp),
)
p.markTimestampResult(ctx, timestamp, start, schema.StatusFailed,
"GDELT file does not exist", 0, 0, 0)
}
p.removeFromProcessing(timestamp)
return
}
if headResult == downloader.HeadError && err != nil {
p.logger.Error("HEAD check failed", zap.String("timestamp", timestamp), zap.Error(err))
p.markTimestampResult(ctx, timestamp, start, schema.StatusFailed,
fmt.Sprintf("HEAD check failed: %v", err), 0, 0, 0)
p.removeFromProcessing(timestamp)
return
}
// Process files in parallel with batching
results := p.processFilesParallel(ctx, tasks)
// Aggregate results
var totalRows, totalArticles, filesProcessed int
var lastError error
for result := range results {
if result.Err != nil {
p.logger.Error("Task failed",
zap.String("file_type", result.Task.FileType),
zap.Error(result.Err),
)
lastError = result.Err
continue
}
filesProcessed++
totalRows += result.RowCount
totalArticles += result.ArticleCount
p.logger.Info("File processed",
zap.String("file_type", result.Task.FileType),
zap.Int("rows", result.RowCount),
zap.Int("articles", result.ArticleCount),
zap.Duration("duration", result.Duration),
)
}
// Mark completion
status := schema.StatusCompleted
errMsg := ""
if lastError != nil {
status = schema.StatusFailed
errMsg = lastError.Error()
}
p.markTimestampResult(ctx, timestamp, start, status, errMsg, filesProcessed, totalRows, totalArticles)
p.removeFromProcessing(timestamp)
p.logger.Info("Timestamp completed",
zap.String("timestamp", timestamp),
zap.Int("files", filesProcessed),
zap.Int("rows", totalRows),
zap.Int("articles", totalArticles),
zap.Duration("duration", time.Since(start)),
)
}
// processFilesParallel processes multiple files in parallel with batching
func (p *Processor) processFilesParallel(ctx context.Context, tasks []Task) <-chan Result {
results := make(chan Result, len(tasks))
var wg sync.WaitGroup
// Process all files in parallel (limited by constants.MaxParallelFiles)
fileSemaphore := make(chan struct{}, constants.MaxParallelFiles)
for _, task := range tasks {
wg.Add(1)
go func(t Task) {
defer wg.Done()
// Acquire file semaphore
fileSemaphore <- struct{}{}
defer func() { <-fileSemaphore }()
result := p.processTask(ctx, t)
results <- result
}(task)
}
go func() {
wg.Wait()
close(results)
}()
return results
}
// markTimestampResult saves the timestamp processing result to database
func (p *Processor) markTimestampResult(ctx context.Context, timestamp string, start time.Time,
status schema.ProcessingStatus, errMsg string, filesProcessed, totalRows, totalArticles int) {
completedAt := time.Now()
ts := &schema.CompletedTimestamp{
Timestamp: timestamp,
Status: status,
StartedAt: start,
CompletedAt: &completedAt,
FilesTotal: constants.FilesPerTimestamp,
FilesProcessed: filesProcessed,
ArticlesCount: totalArticles,
TotalRows: totalRows,
DurationMs: time.Since(start).Milliseconds(),
Error: errMsg,
}
if err := p.storage.MarkTimestampCompleted(ctx, ts); err != nil {
p.logger.Error("Failed to mark timestamp result", zap.Error(err))
}
}
// removeFromProcessing removes a timestamp from the processing map
func (p *Processor) removeFromProcessing(timestamp string) {
p.mu.Lock()
delete(p.processing, timestamp)
p.mu.Unlock()
}
// processTask handles a single file download and processing
func (p *Processor) processTask(ctx context.Context, task Task) Result {
start := time.Now()
// Download stream
reader, _, err := p.downloader.StreamDownload(ctx, task.URL)
if err != nil {
return Result{Task: task, Err: fmt.Errorf("download: %w", err)}
}
defer reader.Close()
// Read ZIP content (required for zip.NewReader)
zipData, err := io.ReadAll(reader)
if err != nil {
return Result{Task: task, Err: fmt.Errorf("read zip: %w", err)}
}
// Parse ZIP stream
records := p.parser.ParseZipStream(ctx, zipData, task.FileType, task.Timestamp)
// Collect and extract articles
articles := make(map[string]*schema.Article)
var rowCount int
for result := range records {
if result.Err != nil {
if result.Err == context.Canceled || result.Err == context.DeadlineExceeded {
return Result{Task: task, Err: result.Err}
}
continue
}
rowCount++
// Extract URL based on record type
url := p.extractURL(result.Record, task.FileType)
if url == "" {
continue
}
// Apply filters when processing GKG files (only GKG has theme/location data)
// When filters are set and we're processing GKG, check if record matches
if p.hasFilters() && task.FileType == constants.FileTypeGKG {
if !p.matchesFilters(result.Record) {
continue // Skip articles that don't match filters
}
} else if p.hasFilters() && task.FileType != constants.FileTypeGKG {
// When filters are set, only process GKG (has theme/location data)
continue
}
// Create or update article
normalized := normalizeURL(url)
if normalized == "" {
continue
}
id := hashURL(normalized)
if _, ok := articles[id]; ok {
// Article already seen, refresh expiration time
articles[id].ExpiresAt = time.Now().Add(time.Duration(schema.ArticleTTLDays) * 24 * time.Hour)
} else {
article := &schema.Article{
ID: id,
URL: url,
ExpiresAt: time.Now().Add(time.Duration(schema.ArticleTTLDays) * 24 * time.Hour),
}
articles[id] = article
}
}
// Save articles in batches
articleSlice := make([]schema.Article, 0, len(articles))
for _, a := range articles {
articleSlice = append(articleSlice, *a)
}
insertedCount := 0
if len(articleSlice) > 0 {
inserted, err := p.storage.BulkInsertArticles(ctx, articleSlice)
if err != nil {
p.logger.Error("Failed to insert articles", zap.Error(err))
}
insertedCount = inserted
}
// Mark URL as processed
processedURL := &schema.ProcessedURL{
URL: task.URL,
Timestamp: task.Timestamp,
FileType: task.FileType,
Status: schema.StatusCompleted,
ProcessedAt: time.Now(),
RowCount: rowCount,
ArticlesExtracted: len(articles),
DurationMs: time.Since(start).Milliseconds(),
}
if err := p.storage.MarkURLProcessed(ctx, processedURL); err != nil {
p.logger.Error("Failed to mark URL processed", zap.Error(err))
}
return Result{
Task: task,
RowCount: rowCount,
ArticleCount: insertedCount,
Duration: time.Since(start),
}
}
// generateTasks creates Task objects for all 3 GDELT files for a timestamp
func (p *Processor) generateTasks(timestamp string) []Task {
return []Task{
{Timestamp: timestamp, URL: fmt.Sprintf(constants.ExportURLTemplate, timestamp), FileType: constants.FileTypeExport},
{Timestamp: timestamp, URL: fmt.Sprintf(constants.MentionsURLTemplate, timestamp), FileType: constants.FileTypeMentions},
{Timestamp: timestamp, URL: fmt.Sprintf(constants.GKGURLTemplate, timestamp), FileType: constants.FileTypeGKG},
}
}
// extractURL extracts the article URL from a parsed record
func (p *Processor) extractURL(record interface{}, fileType string) string {
switch fileType {
case constants.FileTypeExport:
if event, ok := record.(*schema.Event); ok {
return event.SourceURL
}
case constants.FileTypeMentions:
if mention, ok := record.(*schema.Mention); ok {
return mention.MentionIdentifier
}
case constants.FileTypeGKG:
if gkg, ok := record.(*schema.GKG); ok {
return gkg.DocumentIdentifier
}
}
return ""
}
// hasFilters returns true if any filters are configured
func (p *Processor) hasFilters() bool {
return len(p.filterThemes) > 0 || len(p.filterLocations) > 0
}
// matchesFilters checks if a GKG record matches the configured filters
// Returns true if record matches any theme OR location filter (OR logic)
func (p *Processor) matchesFilters(record interface{}) bool {
gkg, ok := record.(*schema.GKG)
if !ok {
return false
}
// Check theme filters (prefix match)
if len(p.filterThemes) > 0 {
allThemes := append(gkg.Themes, gkg.V2Themes...)
for _, filterTheme := range p.filterThemes {
for _, theme := range allThemes {
if strings.HasPrefix(strings.ToUpper(theme), strings.ToUpper(filterTheme)) {
return true
}
}
}
}
// Check location filters (prefix match on location string)
if len(p.filterLocations) > 0 {
allLocations := append(gkg.Locations, gkg.V2Locations...)
for _, filterLoc := range p.filterLocations {
for _, loc := range allLocations {
if strings.HasPrefix(strings.ToUpper(loc), strings.ToUpper(filterLoc)) {
return true
}
}
}
}
// If themes filter is set but didn't match, and no location filter or didn't match
return false
}
// GetTimestampStatus returns the processing status for a timestamp
func (p *Processor) GetTimestampStatus(ctx context.Context, timestamp string) (interface{}, error) {
p.mu.RLock()
isProcessing := p.processing[timestamp]
p.mu.RUnlock()
if isProcessing {
return map[string]interface{}{
"timestamp": timestamp,
"status": "processing",
}, nil
}
ts, err := p.storage.GetTimestampStatus(ctx, timestamp)
if err != nil {
return nil, err
}
if ts == nil {
return map[string]interface{}{
"timestamp": timestamp,
"status": "not_found",
"message": "Timestamp not queued for processing",
}, nil
}
return ts, nil
}
// GetAllTimestamps returns all completed and processing timestamps
func (p *Processor) GetAllTimestamps(ctx context.Context) (interface{}, error) {
completed, processing, err := p.storage.GetAllTimestamps(ctx)
if err != nil {
return nil, err
}
p.mu.RLock()
for ts := range p.processing {
if !containsString(processing, ts) {
processing = append(processing, ts)
}
}
p.mu.RUnlock()
return map[string]interface{}{
"completed": completed,
"processing": processing,
"total_completed": len(completed),
"total_processing": len(processing),
}, nil
}
// Helper functions
func normalizeURL(rawURL string) string {
if rawURL == "" {
return ""
}
u, err := url.Parse(rawURL)
if err != nil {
return ""
}
if u.Scheme == "" {
u.Scheme = "http"
}
u.Host = strings.ToLower(u.Host)
u.Host = strings.TrimPrefix(u.Host, "www.")
q := u.Query()
for _, param := range constants.TrackingParams {
q.Del(param)
}
u.RawQuery = q.Encode()
u.Fragment = ""
u.Path = strings.TrimSuffix(u.Path, "/")
return u.String()
}
func hashURL(normalizedURL string) string {
hash := sha256.Sum256([]byte(normalizedURL))
return hex.EncodeToString(hash[:])
}
func contains(slice []string, item string) bool {
for _, s := range slice {
if s == item {
return true
}
}
return false
}
func containsString(slice []string, item string) bool {
return contains(slice, item)
}
func mergeUnique(a, b []string) []string {
seen := make(map[string]bool)
for _, s := range a {
seen[s] = true
}
for _, s := range b {
if !seen[s] {
a = append(a, s)
seen[s] = true
}
}
return a
}
// ValidateTimestamp is exported for use by API handlers
func ValidateTimestamp(timestamp string) error {
result := ValidateTimestampStructural(timestamp)
if !result.Valid {
return fmt.Errorf(result.Reason)
}
return nil
}