diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 49c8047..98c5e1a 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -39,7 +39,9 @@ jobs: uses: actions/checkout@v4 - name: Validate Compose configuration - run: docker compose config + run: | + docker compose config + docker compose -f deploy/realtime/compose.yaml config --quiet - name: Build NewsLens image run: docker build --tag newslens-api:ci . @@ -89,4 +91,33 @@ jobs: - name: Remove container if: always() - run: docker rm --force newslens-api-ci || true \ No newline at end of file + run: docker rm --force newslens-api-ci || true + + ingestion: + runs-on: ubuntu-latest + + defaults: + run: + working-directory: services/ingestion + + steps: + - name: Check out repository + uses: actions/checkout@v4 + + - name: Set up Go + uses: actions/setup-go@v5 + with: + go-version-file: services/ingestion/go.mod + cache-dependency-path: services/ingestion/go.sum + + - name: Check formatting + run: test -z "$(gofmt -l .)" + + - name: Vet + run: go vet ./... + + - name: Test + run: go test ./... + + - name: Build ingestion binary + run: go build ./cmd/newslens-ingestion diff --git a/README.md b/README.md index 81ee59a..05698d1 100644 --- a/README.md +++ b/README.md @@ -4,11 +4,17 @@ [![Publish container](https://github.com/triasha72/NewsLens/actions/workflows/publish-container.yml/badge.svg)](https://github.com/triasha72/NewsLens/actions/workflows/publish-container.yml) [![Release](https://img.shields.io/badge/release-v0.3.0-blue)](https://github.com/triasha72/NewsLens/releases/tag/v0.3.0) [![Python](https://img.shields.io/badge/python-3.11%20%7C%203.12-blue)](https://www.python.org/) +[![Go](https://img.shields.io/badge/go-1.23-blue)](https://go.dev/) [![License](https://img.shields.io/badge/license-MIT-green)](LICENSE) NewsLens began with a simple question: how much of a news recommender's apparent quality survives once recommendations are evaluated in the order they could actually have been made? -It grew into a leakage-aware news search and recommendation system built on the Microsoft MIND news-recommendation dataset. The project follows several connected questions—temporal leakage, sparse histories, cold-start routing, incompatible score scales, and reproducible serving—from raw records to a tested API and published container. +It grew into a leakage-aware news recommendation system and a real-time search +platform. The offline path follows temporal leakage, sparse histories, cold-start +routing, incompatible score scales, and reproducible serving from Microsoft MIND +records to a tested API. The real-time path follows a new article through Go, +Kafka, PostgreSQL, and freshness-aware search, including the failures that can +happen between acceptance and indexing. ## What was built and why @@ -36,6 +42,17 @@ NewsLens combines: - automated Python and container validation in GitHub Actions; and - multi-platform container publication to GitHub Container Registry. +The real-time search path adds: + +- a Go ingestion API with strict event validation and bounded publish timeouts; +- Kafka with three ordered partitions and a two-member consumer group; +- transactional PostgreSQL idempotency and stale-update protection; +- bounded retries, dead-letter records, fetch backoff, and graceful shutdown; +- category, entity, and freshness query understanding; +- inspectable relevance, freshness, and popularity score components; +- Prometheus counters and lag/freshness gauges; and +- repeatable load, duplicate, DLQ, failover, and backlog-recovery exercises. + The repository contains more than 400 automated tests. ### Experiments and engineering decisions @@ -69,6 +86,16 @@ Docker Desktop Kubernetes cluster. The manifest fixes, bounded service result, and limits of that single-node evidence are recorded in [`docs/KUBERNETES_LOCAL_EVIDENCE.md`](docs/KUBERNETES_LOCAL_EVIDENCE.md). +The separate real-time stack was also built and exercised end to end on Docker +Desktop. In a 500-event, concurrency-20 local run, all events were accepted at +1,225 events/s; publish p99 was 44.12 ms and sampled produced-to-indexed p95 was +78.76 ms. All 25 duplicate replays were recognized. A deliberately stopped +consumer's assigned partition recovered in 5.67 seconds, a retained backlog event +became searchable 2.43 seconds after a consumer restarted, and a malformed Kafka +record reached the DLQ with its original bytes intact. These are bounded +single-machine observations, documented with limitations in +[`docs/REALTIME_LOCAL_EVIDENCE.md`](docs/REALTIME_LOCAL_EVIDENCE.md). + ## Questions that shaped NewsLens The system was not designed from a predetermined architecture checklist. Its components were added as earlier experiments exposed new questions: @@ -100,6 +127,12 @@ flowchart TD H --> I["Frozen evaluation, diagnostics, and selection reports"] I --> J["Versioned v0.3.0 model artifact"] J --> K["FastAPI, Docker, and GitHub Container Registry"] + + L["Article publisher"] --> M["Go ingestion API"] + M --> N["Kafka: 3 partitions"] + N --> O["Go consumer group"] + O --> P["PostgreSQL idempotent article store"] + P --> Q["FastAPI query understanding and freshness ranking"] ``` The selected model uses TF-IDF content recommendations when the user history produces a positive similarity signal. Cold-start and zero-signal requests are routed to a popularity model trained only on the appropriate training partition. @@ -332,12 +365,27 @@ Only load artifacts produced by a trusted NewsLens training workflow. The artifa - `GET /ready` for model-serving readiness; - `GET /model-info` for model metadata; - `POST /recommend` for candidate ranking; +- `GET /realtime/ready` for PostgreSQL-backed search readiness; +- `GET /search` for query understanding and freshness-aware ranking; - fail-fast startup for corrupt configured artifacts; - HTTP `503` when inference is unavailable; - request IDs; - HTTP and inference latency reporting; and - structured request and recommendation logs. +### Real-time article path + +- a statically linked Go producer and consumer binary; +- keyed Kafka delivery across three partitions; +- two consumers in one group with bounded failure detection; +- transactional event-ledger idempotency in PostgreSQL; +- stale article update protection; +- bounded database and broker retry delays; +- a dead-letter topic that preserves malformed source bytes; +- graceful process shutdown; and +- Prometheus metrics for accepted, processed, duplicate, failed, dead-lettered, + lag, and freshness signals. + ### Container and CI/CD - non-root Docker runtime; @@ -345,6 +393,7 @@ Only load artifacts produced by a trusted NewsLens training workflow. The artifa - artifact-free image construction; - read-only artifact mounting through Docker Compose; - Python linting and test execution in CI; +- Go formatting, vet, test, and build checks in CI; - container build validation in CI; - artifact-free liveness and readiness-contract checks; - multi-platform `linux/amd64` and `linux/arm64` images; @@ -589,6 +638,38 @@ A successful response includes: See [`docs/API.md`](docs/API.md) for the complete API contract. +## Run real-time search locally + +This path does not need the MIND files or a recommendation artifact. It needs +Docker Desktop with Compose v2: + +```bash +./scripts/realtime/start_stack.sh +``` + +Publish and search one article: + +```bash +curl -X POST http://127.0.0.1:8080/events \ + -H 'Content-Type: application/json' \ + -d '{ + "event_id": "readme-event-1", + "article_id": "readme-article-1", + "title": "Apple launches an AI chip today", + "category": "technology", + "published_at": "2026-08-23T12:00:00Z", + "produced_at": "2026-08-23T12:00:01Z", + "body": "A local end-to-end example." + }' + +curl --get \ + --data-urlencode 'q=latest Apple AI chip' \ + http://127.0.0.1:8000/search +``` + +The operations guide covers load, duplicate, DLQ, ordering, and recovery checks: +[`docs/REALTIME_OPERATIONS.md`](docs/REALTIME_OPERATIONS.md). + ## Docker deployment Before starting Docker Compose, generate the local artifact: @@ -680,16 +761,26 @@ NewsLens/ │ ├── DECISIONS.md │ ├── DEPLOYMENT.md │ ├── EVALUATION.md +│ ├── REALTIME_ARCHITECTURE.md +│ ├── REALTIME_LOCAL_EVIDENCE.md +│ ├── REALTIME_OPERATIONS.md │ └── RESEARCH_QUESTIONS.md +├── deploy/ +│ └── realtime/ ├── reports/ │ ├── content_metrics.json │ ├── fallback_metrics.json │ ├── mindsmall_dev_audit.json │ ├── mindsmall_train_audit.json +│ ├── realtime_load_v0_1.json +│ ├── realtime_recovery_v0_1.json │ └── popularity_metrics.json ├── scripts/ +│ ├── realtime/ │ ├── setup.ps1 │ └── setup.sh +├── services/ +│ └── ingestion/ ├── src/ │ └── newslens/ │ ├── api/ @@ -726,6 +817,7 @@ NewsLens/ │ │ ├── fallback.py │ │ ├── popularity.py │ │ └── tfidf.py +│ ├── realtime/ │ ├── __init__.py │ ├── __main__.py │ └── cli.py @@ -786,12 +878,16 @@ The automated suite covers: - request and response validation; - readiness and failure behavior; - request observability; and -- API integration. +- API integration; +- Go event validation and HTTP behavior; and +- idempotent worker outcomes, retries, and dead-letter routing. -GitHub Actions runs two main CI jobs: +GitHub Actions runs three CI jobs: 1. **Quality:** installs the project, runs Ruff, and executes the full Python test suite. 2. **Container:** validates Compose, builds the Docker image, starts an artifact-free container, checks liveness, and verifies that readiness correctly returns HTTP `503` without a model. +3. **Ingestion:** checks Go formatting, runs `go vet` and unit tests, and builds + the ingestion binary. Tagged releases also publish multi-platform images to GitHub Container Registry. @@ -814,9 +910,13 @@ Tagged releases also publish multi-platform images to GitHub Container Registry. - Category and exposure cohorts can overlap. - Subgroup-specific confidence intervals are not currently reported. - The published container does not include the licensed dataset or generated model artifact. -- The DuckDB layer is a local analytical batch database, not a hosted PostgreSQL - service, streaming ingestion system, or online feature store. -- Request logs and latency headers provide service-level observability, but there is no external metrics store or alerting system. +- The DuckDB layer remains a local analytical batch database. The separate + PostgreSQL store contains streamed articles and is not an online feature store + for the recommendation model. +- The local real-time stack has Prometheus scraping, but no durable external + metrics store, alerting policy, or operated DLQ replay process. +- The real-time evidence uses one Kafka broker and one PostgreSQL instance on one + Docker Desktop host; it does not establish replicated or multi-zone operation. - A versioned container is published, but NewsLens is not operated as a public, always-on hosted service. - No online experiment has been conducted. - Offline metric improvements do not establish production or business impact. diff --git a/ROADMAP.md b/ROADMAP.md index 4340ef9..f2a32f8 100644 --- a/ROADMAP.md +++ b/ROADMAP.md @@ -44,6 +44,20 @@ Paired bootstrap intervals favor the fallback system over content-only ranking f The selected model can be exported as a versioned, checksummed artifact, loaded by a typed API, mounted read-only into a non-root container, validated in CI, and published for both `linux/amd64` and `linux/arm64`. The same artifact was also served through two replicas on Docker Desktop Kubernetes with probes, service routing, restricted egress checks, and a bounded load run. +### Real-time search and ingestion + +A separate Go service now validates and publishes article events to a +three-partition Kafka topic. Two group consumers apply bounded retries and write +to PostgreSQL through a transactional event ledger, so redelivery is idempotent +and an older update cannot overwrite a newer article. Invalid or exhausted +records move to a DLQ, and Prometheus scrapes producer and consumer signals. + +FastAPI reads a bounded article set and returns category, entity, and freshness +intent with inspectable ranking components. The local end-to-end evidence covers +500-event load, 25 duplicate replays, a targeted partition reassignment, retained +backlog recovery, and a malformed-event DLQ envelope. Exact results and non-claims +are in [`docs/REALTIME_LOCAL_EVIDENCE.md`](docs/REALTIME_LOCAL_EVIDENCE.md). + ### Reproducible analytical data Validated MIND records can be materialized into normalized DuckDB tables for @@ -111,6 +125,14 @@ health and service-routing checks, and completed the documented bounded load test. This is local deployment evidence, not cloud or multi-node production evidence. +### Can a new article move from acceptance to searchable state through a recoverable event path? + +Yes, within the tested single-host topology. Go publishes keyed events to Kafka, +group consumers persist them idempotently to PostgreSQL, and FastAPI ranks the +stored article. Local evidence also covers duplicate delivery, malformed-event +dead lettering, consumer reassignment, and backlog recovery. Broker replication, +database high availability, and live-traffic relevance remain outside that claim. + ## Next investigations ### 1. History recency and weighting @@ -170,14 +192,17 @@ Content and popularity scores must remain separate unless a defensible calibrati ### 5. Performance and reliability -**Question:** Where are the serving limits of the current artifact and API? +**Question:** Where are the sustained and multi-host limits beyond the current +bounded local evidence? **Plan:** -- benchmark startup time, memory, throughput, and p50/p95/p99 latency; +- extend the current p50/p95/p99, freshness, and recovery measurements into a + sustained soak workload; - test multiple history and candidate-set sizes; - exercise corrupt, missing, and incompatible artifacts; -- add graceful shutdown and concurrency tests; +- add broker and PostgreSQL failure exercises beyond the existing consumer + process recovery test; - publish the hardware and workload used for every benchmark; and - define thresholds before optimizing. diff --git a/deploy/realtime/compose.yaml b/deploy/realtime/compose.yaml new file mode 100644 index 0000000..023c34f --- /dev/null +++ b/deploy/realtime/compose.yaml @@ -0,0 +1,132 @@ +name: newslens-realtime + +x-consumer: &consumer + build: + context: ../.. + dockerfile: services/ingestion/Dockerfile + environment: &consumer-environment + NEWSLENS_INGESTION_MODE: consumer + NEWSLENS_KAFKA_BROKERS: kafka:9092 + NEWSLENS_KAFKA_TOPIC: news-events + NEWSLENS_KAFKA_DLQ_TOPIC: news-events-dlq + NEWSLENS_KAFKA_GROUP_ID: newslens-indexers + NEWSLENS_DATABASE_URL: postgresql://newslens:newslens@postgres:5432/newslens + depends_on: + kafka-init: + condition: service_completed_successfully + postgres: + condition: service_healthy + restart: unless-stopped + init: true + +services: + postgres: + build: + context: ../.. + dockerfile: deploy/realtime/postgres/Dockerfile + environment: + POSTGRES_DB: newslens + POSTGRES_USER: newslens + POSTGRES_PASSWORD: newslens + volumes: + - realtime-postgres:/var/lib/postgresql/data + healthcheck: + test: ["CMD-SHELL", "pg_isready -U newslens -d newslens"] + interval: 3s + timeout: 3s + retries: 20 + + kafka: + image: apache/kafka:3.9.1 + hostname: kafka + environment: + KAFKA_NODE_ID: 1 + KAFKA_PROCESS_ROLES: broker,controller + KAFKA_LISTENERS: PLAINTEXT://:9092,CONTROLLER://:9093 + KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092 + KAFKA_CONTROLLER_LISTENER_NAMES: CONTROLLER + KAFKA_LISTENER_SECURITY_PROTOCOL_MAP: CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT + KAFKA_CONTROLLER_QUORUM_VOTERS: 1@kafka:9093 + KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR: 1 + KAFKA_TRANSACTION_STATE_LOG_MIN_ISR: 1 + KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS: 0 + KAFKA_AUTO_CREATE_TOPICS_ENABLE: "false" + healthcheck: + test: ["CMD-SHELL", "/opt/kafka/bin/kafka-topics.sh --bootstrap-server localhost:9092 --list >/dev/null"] + interval: 5s + timeout: 5s + retries: 30 + + kafka-init: + image: apache/kafka:3.9.1 + depends_on: + kafka: + condition: service_healthy + command: + - /bin/bash + - -ec + - | + /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic news-events --partitions 3 --replication-factor 1 + /opt/kafka/bin/kafka-topics.sh --bootstrap-server kafka:9092 --create --if-not-exists --topic news-events-dlq --partitions 1 --replication-factor 1 + + ingestion-api: + build: + context: ../.. + dockerfile: services/ingestion/Dockerfile + environment: + NEWSLENS_INGESTION_MODE: api + NEWSLENS_KAFKA_BROKERS: kafka:9092 + NEWSLENS_KAFKA_TOPIC: news-events + ports: + - "8080:8080" + depends_on: + kafka-init: + condition: service_completed_successfully + restart: unless-stopped + init: true + + consumer-1: + <<: *consumer + environment: + <<: *consumer-environment + NEWSLENS_KAFKA_CLIENT_ID: newslens-consumer-1 + ports: + - "8081:8080" + + consumer-2: + <<: *consumer + environment: + <<: *consumer-environment + NEWSLENS_KAFKA_CLIENT_ID: newslens-consumer-2 + ports: + - "8082:8080" + + search-api: + build: + context: ../.. + dockerfile: Dockerfile + environment: + NEWSLENS_REALTIME_DATABASE_URL: postgresql://newslens:newslens@postgres:5432/newslens + ports: + - "8000:8000" + depends_on: + postgres: + condition: service_healthy + restart: unless-stopped + init: true + + prometheus: + build: + context: ../.. + dockerfile: deploy/realtime/prometheus.Dockerfile + command: ["--config.file=/etc/prometheus/prometheus.yml"] + ports: + - "9090:9090" + depends_on: + - ingestion-api + - consumer-1 + - consumer-2 + +volumes: + realtime-postgres: diff --git a/deploy/realtime/postgres/Dockerfile b/deploy/realtime/postgres/Dockerfile new file mode 100644 index 0000000..9b845ae --- /dev/null +++ b/deploy/realtime/postgres/Dockerfile @@ -0,0 +1,2 @@ +FROM postgres:17-alpine +COPY deploy/realtime/postgres/init.sql /docker-entrypoint-initdb.d/001-realtime.sql diff --git a/deploy/realtime/postgres/init.sql b/deploy/realtime/postgres/init.sql new file mode 100644 index 0000000..bf1ed42 --- /dev/null +++ b/deploy/realtime/postgres/init.sql @@ -0,0 +1,23 @@ +CREATE TABLE IF NOT EXISTS realtime_ingestion_events ( + event_id TEXT PRIMARY KEY, + article_id TEXT NOT NULL, + produced_at TIMESTAMPTZ NOT NULL, + indexed_at TIMESTAMPTZ NOT NULL +); + +CREATE TABLE IF NOT EXISTS realtime_articles ( + article_id TEXT PRIMARY KEY, + title TEXT NOT NULL, + body TEXT NOT NULL DEFAULT '', + category TEXT NOT NULL, + published_at TIMESTAMPTZ NOT NULL, + produced_at TIMESTAMPTZ NOT NULL, + indexed_at TIMESTAMPTZ NOT NULL, + popularity INTEGER NOT NULL DEFAULT 0 CHECK (popularity >= 0) +); + +CREATE INDEX IF NOT EXISTS realtime_articles_category_published_idx + ON realtime_articles (category, published_at DESC); + +CREATE INDEX IF NOT EXISTS realtime_articles_published_idx + ON realtime_articles (published_at DESC); diff --git a/deploy/realtime/prometheus.Dockerfile b/deploy/realtime/prometheus.Dockerfile new file mode 100644 index 0000000..e88d3fb --- /dev/null +++ b/deploy/realtime/prometheus.Dockerfile @@ -0,0 +1,2 @@ +FROM prom/prometheus:v3.5.0 +COPY deploy/realtime/prometheus.yml /etc/prometheus/prometheus.yml diff --git a/deploy/realtime/prometheus.yml b/deploy/realtime/prometheus.yml new file mode 100644 index 0000000..967ba05 --- /dev/null +++ b/deploy/realtime/prometheus.yml @@ -0,0 +1,10 @@ +global: + scrape_interval: 5s + +scrape_configs: + - job_name: newslens-ingestion-api + static_configs: + - targets: ["ingestion-api:8080"] + - job_name: newslens-consumers + static_configs: + - targets: ["consumer-1:8080", "consumer-2:8080"] diff --git a/docs/API.md b/docs/API.md index 5228a9b..dfa4be3 100644 --- a/docs/API.md +++ b/docs/API.md @@ -1,7 +1,7 @@ # NewsLens API -NewsLens provides an artifact-backed FastAPI service for candidate-based news -recommendation. +NewsLens provides an artifact-backed recommendation API and a PostgreSQL-backed +search API for articles delivered by the real-time ingestion path. ## Current capabilities @@ -17,6 +17,10 @@ The API provides: - automatic OpenAPI documentation; and - fail-fast startup for missing or corrupt configured artifacts. +The search route also provides deterministic category, entity, and freshness +intent, bounded candidate retrieval, freshness-aware ranking, and score-component +diagnostics. Recommendation readiness and search-store readiness are independent. + ## Artifact configuration The API reads the model location from the @@ -35,6 +39,13 @@ Only load artifacts generated by a trusted NewsLens training workflow. Joblib uses pickle-compatible deserialization. Checksums protect against accidental corruption but do not make an untrusted pickle safe. +The streamed-article database is configured separately: + +```bash +export NEWSLENS_REALTIME_DATABASE_URL=postgresql://newslens:newslens@localhost:5432/newslens +export NEWSLENS_REALTIME_CANDIDATE_LIMIT=200 +``` + ## Running locally Install NewsLens with development dependencies: @@ -124,6 +135,66 @@ The recommendation source is: - `content` when the user's history produces a usable TF-IDF profile; or - `popularity` when the user is cold-start or has no usable content signal. +## `GET /realtime/ready` + +Returns `200` only when the configured PostgreSQL article store answers a probe: + +```json +{ + "status": "ready", + "realtime_store_ready": true +} +``` + +It returns `503` when the database is not configured or reachable. This does not +depend on the recommendation model artifact. + +## `GET /search` + +Searches the bounded streamed-article candidate set: + +```bash +curl --get \ + --data-urlencode 'q=latest Apple AI chip' \ + --data-urlencode 'top_k=5' \ + http://127.0.0.1:8000/search +``` + +Example response: + +```json +{ + "request_id": "request-identifier", + "intent": { + "normalized_query": "latest Apple AI chip", + "category": "technology", + "entity": "Apple", + "prefers_freshness": true + }, + "candidate_count": 2, + "returned_count": 2, + "search_ms": 1.8, + "results": [ + { + "article_id": "article-123", + "title": "Apple launches an AI chip today", + "category": "technology", + "published_at": "2026-08-23T12:00:00+00:00", + "score": 0.91, + "relevance_score": 0.90, + "freshness_score": 0.98, + "popularity_score": 0.60, + "index_freshness_ms": 62.5 + } + ] +} +``` + +`index_freshness_ms` is the event's stored `produced_at` to PostgreSQL +`indexed_at` duration. `search_ms` measures candidate retrieval and ranking for +this request. Query length is limited to 500 characters, `top_k` to 1–100, and +database candidates to `NEWSLENS_REALTIME_CANDIDATE_LIMIT`. + ## Validation behavior The API rejects: @@ -137,7 +208,9 @@ The API rejects: Invalid requests return HTTP `422`. -Inference without a loaded model returns HTTP `503`. +Inference without a loaded model returns HTTP `503`. Search without a ready +real-time store also returns `503`; either route can remain ready while the other +is unavailable. ## Testing @@ -277,4 +350,4 @@ Recommendation logs include: - routing source; and - inference latency. -These measurements provide local service observability. They do not yet represent a complete production monitoring, metrics-storage, or alerting system. \ No newline at end of file +These measurements provide local service observability. They do not yet represent a complete production monitoring, metrics-storage, or alerting system. diff --git a/docs/DECISIONS.md b/docs/DECISIONS.md index bc1eb97..b0f8d16 100644 --- a/docs/DECISIONS.md +++ b/docs/DECISIONS.md @@ -71,3 +71,41 @@ replacement improve reproducibility and failure safety, but DuckDB is not being used as an online feature store. A PostgreSQL service should be considered only if future work establishes a need for shared access, concurrent writes, or operational event ingestion. + +--- + +## D003 - Separate real-time transport from offline recommendation research + +### Decision + +Use a small Go service and Kafka consumer group for streamed article ingestion, +PostgreSQL for the searchable write model, and the existing Python API for query +understanding and ranking. + +### Context + +The MIND pipeline is a batch research workflow whose DuckDB database and model +artifacts are deliberately local and immutable. New articles introduce concurrent +writes, bursts, redelivery, failure recovery, and freshness expectations. Putting +those concerns into the recommendation process would couple model readiness to +event transport and create an unbounded in-memory failure boundary. + +### Rationale + +Go provides a small statically linked server for the I/O path. Kafka retains +bursts, preserves per-article key order, and distributes three partitions across +a consumer group. PostgreSQL supplies concurrent transactions and makes the event +ledger plus article update one atomic idempotency boundary. Python remains the +right place for transparent search features and same-repository evaluation. + +### Consequences + +The system now has independent recommendation and real-time readiness contracts. +Its at-least-once Kafka delivery can redeliver records during the one-second +offset-commit window, so stable event IDs are part of the API contract. A late +event cannot overwrite a newer article version. Invalid or exhausted records +require an operated DLQ review and replay process. + +The local Compose system has one broker and one PostgreSQL instance. The split +creates clear production scaling and recovery boundaries, but it does not claim +broker replication, database high availability, or multi-zone operation. diff --git a/docs/DEPLOYMENT.md b/docs/DEPLOYMENT.md index 5cfc218..300df7e 100644 --- a/docs/DEPLOYMENT.md +++ b/docs/DEPLOYMENT.md @@ -2,6 +2,10 @@ NewsLens provides a Dockerized FastAPI recommendation service. +It also includes a separate local Compose topology for real-time ingestion and +search. The two deployments share the Python package but do not require one +another to be ready. + ## Prerequisites - Docker Desktop @@ -150,6 +154,44 @@ curl --fail http://localhost/ready The HPA requires the Kubernetes Metrics Server. The NetworkPolicy requires a network plugin that enforces `networking.k8s.io/v1` policies. +## Real-time search stack + +Start Kafka, PostgreSQL, the Go producer, two Go consumers, the Python search +API, and Prometheus: + +```bash +./scripts/realtime/start_stack.sh +``` + +The script uses `deploy/realtime/compose.yaml`, builds immutable PostgreSQL and +Prometheus configuration images, waits for container health, and verifies all +four application readiness endpoints. Ports are: + +| Port | Service | +|---:|---| +| 8000 | search and recommendation API | +| 8080 | Go event producer | +| 8081 | consumer 1 health and metrics | +| 8082 | consumer 2 health and metrics | +| 9090 | Prometheus | + +Run the evidence suite only after readiness succeeds: + +```bash +PYTHONPATH=src python scripts/evaluate_realtime_search.py +python scripts/realtime/benchmark_ingestion.py \ + --events 500 --concurrency 20 --freshness-samples 25 --duplicate-probes 25 +python scripts/realtime/failure_recovery.py +python scripts/realtime/verify_dlq.py +python scripts/realtime/verify_ordering.py +``` + +The failure script stops and starts consumers in this named local Compose +project. See [`REALTIME_OPERATIONS.md`](REALTIME_OPERATIONS.md) for the event +contract, DLQ inspection, recovery procedure, and data-reset command. Measured +local results and their limits are in +[`REALTIME_LOCAL_EVIDENCE.md`](REALTIME_LOCAL_EVIDENCE.md). + ## Repeatable load test Create a request using article IDs present in the mounted artifact: diff --git a/docs/REALTIME_ARCHITECTURE.md b/docs/REALTIME_ARCHITECTURE.md new file mode 100644 index 0000000..b4035a2 --- /dev/null +++ b/docs/REALTIME_ARCHITECTURE.md @@ -0,0 +1,191 @@ +# Real-time search architecture + +NewsLens originally answered an offline recommendation question: which ranking +decisions still hold under a chronological evaluation? The real-time path asks a +different question: how does a newly published article become searchable without +turning the research code into an unbounded in-process queue? + +The implementation separates event transport from search serving. A small Go +service accepts and validates article events, Kafka absorbs bursts and assigns +ordered partitions, Go consumers write idempotently to PostgreSQL, and FastAPI +retrieves and ranks the searchable records. This keeps the existing MIND +recommendation route intact while adding an independently ready search route. + +## Request and event flow + +```mermaid +sequenceDiagram + participant P as Producer + participant G as Go ingestion API + participant K as Kafka (3 partitions) + participant C as Go consumer group + participant D as PostgreSQL + participant S as FastAPI search + + P->>G: POST /events + G->>G: bound, decode, normalize, validate + G->>K: keyed event, required acknowledgements + K-->>G: accepted + G-->>P: 202 + event_id + K->>C: one partition record + C->>D: transaction: event ledger + article upsert + D-->>C: committed or duplicate + C->>K: commit offset + S->>D: bounded candidate query + S->>S: query intent + relevance/freshness rank + S-->>P: ranked results + diagnostics +``` + +The producer response means Kafka accepted the event. It does not mean the +article is already searchable. Produced-to-indexed time is stored with each +article and returned as `index_freshness_ms`; the benchmark polls the search API +and records the distribution separately from publish latency. + +## Event contract + +`POST /events` accepts one JSON object no larger than 1 MiB: + +```json +{ + "event_id": "publisher-unique-id", + "article_id": "article-123", + "title": "Apple launches an AI chip", + "category": "technology", + "published_at": "2026-08-23T12:00:00Z", + "body": "Article text", + "produced_at": "2026-08-23T12:00:01Z" +} +``` + +`article_id`, `title`, `category`, and `published_at` are required. The API can +generate `event_id` and `produced_at`, though stable publisher-supplied event IDs +are preferred because they provide retry idempotency. Unknown fields, multiple +JSON objects, missing fields, excessive strings, and oversized payloads are +rejected before Kafka. + +Events are keyed by `article_id`. Kafka therefore preserves order for changes to +one article within a partition, while unrelated articles can be handled in +parallel across three partitions. + +## Delivery and consistency semantics + +Kafka consumption is at least once. A consumer queues an offset for commit only +after the PostgreSQL transaction succeeds or the event is placed on the +dead-letter topic. Kafka flushes those commits every second, so a crash can +redeliver the most recent successful records; the database idempotency boundary +handles that window. + +PostgreSQL turns redelivery into an idempotent operation: + +1. `realtime_ingestion_events.event_id` is the idempotency key. +2. The event ledger insert and article upsert share one transaction. +3. An existing event ID is counted as a duplicate and leaves the article alone. +4. A different event ID may update the article only when its `produced_at` is not + older than the stored version. + +The result is effectively-once database mutation for a stable event ID, built on +at-least-once delivery. It is not global exactly-once delivery: an external side +effect added outside this transaction would need its own idempotency boundary. + +Search is read-after-index, not read-after-publish. PostgreSQL is the current +searchable store, so a successful consumer transaction becomes visible to new +queries without a separate index refresh job. + +## Consumer groups, retries, and dead letters + +The local stack starts two consumers with one group ID. Kafka assigns each +partition to at most one active group member and rebalances when a member enters +or leaves. A two-second heartbeat and six-second session timeout bound local +process-failure detection. Processing stays deliberately synchronous inside each consumer. That +bounds memory and lets Kafka retain the backlog when PostgreSQL is slower than +the producer. + +Database writes receive three bounded attempts by default with increasing +100 ms delays. Invalid events skip database retries. Invalid events and exhausted +writes are published to `news-events-dlq` with their source location, failure +reason, timestamp, and base64-encoded original bytes. The source offset is +committed only after the dead-letter write succeeds, so a DLQ outage does not +silently discard the record. + +Kafka fetch failures also back off up to one second instead of forming a hot +retry loop. + +## Backpressure and overload behavior + +The HTTP producer performs a synchronous Kafka write with required broker +acknowledgements and a five-second request deadline. Slow or unavailable Kafka +therefore slows or rejects producers instead of filling an in-memory queue. + +Consumers fetch and process one record at a time. When PostgreSQL slows, consumer +lag grows in Kafka; records are not accumulated in application memory. The API +limits event bodies to 1 MiB and search retrieval to a configurable bounded +candidate count (`NEWSLENS_REALTIME_CANDIDATE_LIMIT`, default 200). + +These choices favor an inspectable failure mode over maximum throughput. Batch +writes and asynchronous producer buffering should only be introduced with +measured latency, loss, and memory budgets. + +## Query understanding and ranking + +The search endpoint normalizes whitespace and extracts three small, inspectable +signals: + +- a category from deterministic keyword groups; +- a capitalized or uppercase entity phrase; and +- freshness intent from words such as `latest`, `today`, and `breaking`. + +PostgreSQL applies category and entity filters when they have matches, then +returns a recent bounded candidate set. Ranking combines lexical overlap, +exponential freshness decay, and log-normalized popularity. Freshness receives a +20% weight only when the query asks for recent information; otherwise that weight +returns to relevance. Responses expose every component so a surprising order can +be inspected rather than attributed to an opaque score. + +The deterministic contract fixture in +[`reports/realtime_search_evaluation_v0_1.json`](../reports/realtime_search_evaluation_v0_1.json) +contains six hand-labeled queries and ten articles. It is useful evidence that +the intended rules work, but it is too small and synthetic for a general search +quality claim. + +## Observability + +Every Go process exposes: + +- `GET /health` for process liveness; +- `GET /ready` for its configured Kafka or PostgreSQL dependencies; and +- `GET /metrics` in Prometheus text format. + +Metrics cover accepted, processed, duplicate, invalid, dead-lettered, and failed +events, consumer lag when reported by the Kafka client, and the last +produced-to-indexed duration. Prometheus scrapes the producer and both consumers. +The Python API retains request IDs, process time, structured search completion +logs, and per-result index freshness. + +The local metrics are signals, not an alerting policy. Production operation would +still need durable metric storage, dashboards, paging thresholds, and an owned +DLQ replay process. + +## Failure behavior + +| Failure | Visible behavior | Recovery path | Data-loss boundary | +|---|---|---|---| +| One consumer stops | Kafka rebalances its partitions to the other consumer | Restart the process; idempotency absorbs redelivery | No loss expected while Kafka retains the event | +| All consumers stop | Publish can continue and lag grows; new articles are not searchable | Restore a consumer and drain the backlog | Kafka retention policy | +| PostgreSQL is unavailable | Writes retry, then move to DLQ if Kafka remains available | Restore PostgreSQL and replay reviewed DLQ events | DLQ retention and replay discipline | +| Kafka is unavailable | Producer readiness fails and `POST /events` returns 503/504 | Restore Kafka and let producers retry with the same event ID | Producer retry policy | +| Consumer crashes after DB commit | Kafka may redeliver the same event | Ledger recognizes the event ID as a duplicate | Stable event ID required | +| Older article update arrives late | Event is recorded but cannot overwrite a newer article version | No operator action | Producer timestamps must be trustworthy | +| Invalid payload reaches Kafka | Consumer writes the bytes and reason to DLQ | Correct the producer or replay a repaired event | DLQ retention | + +## Scaling boundaries + +Consumer parallelism cannot usefully exceed the number of Kafka partitions. Add +partitions before adding more active consumers, while accounting for the fact +that partition increases change key distribution. PostgreSQL is the shared write +and search bottleneck; connection-pool, index, query, and write amplification +must be measured before horizontal API scaling is treated as capacity evidence. + +The included Compose topology is a single-host development system with one Kafka +broker and one PostgreSQL instance. It demonstrates component boundaries and +failure handling, not broker replication, database high availability, multi-zone +recovery, or production capacity. diff --git a/docs/REALTIME_LOCAL_EVIDENCE.md b/docs/REALTIME_LOCAL_EVIDENCE.md new file mode 100644 index 0000000..f4dbf67 --- /dev/null +++ b/docs/REALTIME_LOCAL_EVIDENCE.md @@ -0,0 +1,104 @@ +# Real-time search local evidence + +This record separates behavior observed in a live stack from behavior covered +only by unit tests or design analysis. It was produced on August 23, 2026 after +building the repository's own images and starting the complete Compose topology. + +## Environment + +| Property | Value | +|---|---| +| Runtime | Docker Desktop 29.7.2, Linux containers on Apple silicon | +| Docker allocation | 10 CPUs, 8.32 GB memory | +| Kafka | Apache Kafka 3.9.1, one broker, three `news-events` partitions | +| Consumers | Two Go 1.23 services in `newslens-indexers` | +| Search store | PostgreSQL 17 | +| Search API | Python 3.12 and FastAPI | +| Metrics | Prometheus 3.5.0 | + +The commands and report generators are documented in +[`REALTIME_OPERATIONS.md`](REALTIME_OPERATIONS.md). Raw machine-readable results +are stored in `reports/`. + +## End-to-end load and freshness + +The workload published 500 events with concurrency 20, sampled 25 articles for +produced-to-indexed freshness, and replayed 25 stable event IDs. + +| Observation | Result | +|---|---:| +| Accepted publish requests | 500 / 500 | +| Publish failure rate | 0% | +| Publish throughput | 1,225 events/s | +| Publish latency p50 / p95 / p99 | 15.15 / 20.02 / 44.12 ms | +| Produced-to-indexed p50 / p95 / p99 | 62.51 / 78.76 / 82.07 ms | +| Duplicate probes recognized | 25 / 25 | + +The first run exposed a one-second Kafka producer batch delay. A second run +showed that synchronous per-event offset commits moved a 500-event burst's p95 +freshness above 18 seconds. The final configuration uses a 10 ms producer batch +timeout and one-second periodic consumer commits; the table reports only the +post-fix rerun. + +This is one single-machine development workload, not a production capacity or +service-level claim. The freshness sample records each article's own +`produced_at` to PostgreSQL `indexed_at` duration; search polling confirms the +article is visible but is not added to that stored duration. + +Evidence: [`realtime_load_v0_1.json`](../reports/realtime_load_v0_1.json). + +## Consumer failure and backlog recovery + +The failure probe first read the active group assignment, chose an article key +owned by consumer 1, and stopped that consumer. Kafka reassigned its partition +and the event became searchable in 5.67 seconds. The probe then stopped both +consumers, published an event into the backlog, and measured 2.43 seconds from +starting one consumer until the event was searchable. + +This proves local process reassignment and retained-backlog recovery. It does not +exercise broker, PostgreSQL, host, or multi-zone failure. + +Evidence: [`realtime_recovery_v0_1.json`](../reports/realtime_recovery_v0_1.json). + +## Dead-letter behavior + +A malformed non-JSON record was inserted directly into `news-events`. A consumer +counted one invalid record, published one dead-letter record, and preserved the +original bytes in a valid base64 envelope with the source position and reason. +Database retry exhaustion is covered by the Go worker test; it was not induced in +the live PostgreSQL service. + +Evidence: [`realtime_dlq_v0_1.json`](../reports/realtime_dlq_v0_1.json). + +## Late-event ordering + +Two distinct events updated one article, with the older `produced_at` event sent +second. Both entered the event ledger, while the searchable title remained the +newer version. This verifies the PostgreSQL conditional upsert in the live stack; +its correctness still depends on trustworthy producer timestamps. + +Evidence: [`realtime_ordering_v0_1.json`](../reports/realtime_ordering_v0_1.json). + +## Query and ranking contract + +A deterministic fixture checks 18 category, entity, and freshness fields across +six labeled queries. All 18 match. On the same ten-article fixture, +freshness-aware ranking raises NDCG@10 from 0.869 to 1.000 and +freshness-weighted NDCG@10 from 0.634 to 0.878 relative to the relevance-only +configuration. + +This small hand-labeled fixture shows that the implemented rules behave as +specified. It is not an independently judged corpus and cannot support a broad +search-quality claim. + +Evidence: +[`realtime_search_evaluation_v0_1.json`](../reports/realtime_search_evaluation_v0_1.json). + +## What this evidence does not establish + +- Kafka broker replication or loss behavior; +- PostgreSQL high availability or point-in-time recovery; +- multi-host, multi-zone, or managed-cloud operation; +- sustained capacity beyond this short workload; +- a production alerting and DLQ replay process; or +- relevance quality on live or independently judged search traffic. diff --git a/docs/REALTIME_OPERATIONS.md b/docs/REALTIME_OPERATIONS.md new file mode 100644 index 0000000..2614d2e --- /dev/null +++ b/docs/REALTIME_OPERATIONS.md @@ -0,0 +1,142 @@ +# Real-time search operations + +This runbook exercises the NewsLens event path on one machine. It uses the same +commands for success, load, and failure evidence so a report can be reproduced +instead of copied from an unrelated environment. + +## Start and verify + +Requirements are Docker Engine with Compose v2, `curl`, Bash, and Python 3.11 or +3.12. The stack uses ports 8000, 8080, 8081, 8082, and 9090. + +```bash +./scripts/realtime/start_stack.sh +docker compose -f deploy/realtime/compose.yaml ps +``` + +The start script builds both applications, waits for Compose health checks, and +then verifies the producer, both consumers, and the search store. Prometheus is +available at . + +Publish one event: + +```bash +curl --fail-with-body \ + -X POST http://127.0.0.1:8080/events \ + -H 'Content-Type: application/json' \ + -d '{ + "event_id": "manual-event-1", + "article_id": "manual-article-1", + "title": "Apple launches an AI chip today", + "category": "technology", + "published_at": "2026-08-23T12:00:00Z", + "produced_at": "2026-08-23T12:00:01Z", + "body": "A manual end-to-end probe." + }' + +curl --get --fail-with-body \ + --data-urlencode 'q=latest Apple AI chip' \ + --data-urlencode 'top_k=5' \ + http://127.0.0.1:8000/search +``` + +Publishing the same request again should increment the duplicate metric without +creating a second article mutation: + +```bash +curl http://127.0.0.1:8081/metrics +curl http://127.0.0.1:8082/metrics +``` + +## Quality contract + +The query and ranking fixture does not need Docker: + +```bash +PYTHONPATH=src python scripts/evaluate_realtime_search.py +``` + +It writes `reports/realtime_search_evaluation_v0_1.json`, including each expected +and observed intent field plus relevance-only and freshness-aware ranking metrics. + +## Load and freshness evidence + +After the stack is ready: + +```bash +python scripts/realtime/benchmark_ingestion.py \ + --events 500 \ + --concurrency 20 \ + --freshness-samples 25 +``` + +The report records publish throughput, failure rate, p50/p95/p99 acknowledgement +latency, and sampled produced-to-indexed p50/p95/p99. Quote a result only with its +event count, concurrency, machine, Compose topology, and report limitations. + +## Failure and recovery evidence + +The recovery exercise stops services and must run only against the local Compose +project: + +```bash +python scripts/realtime/failure_recovery.py +python scripts/realtime/verify_ordering.py +``` + +It first stops one consumer and proves the other can index a probe. It then stops +both consumers, publishes a backlog event, restores one consumer, and measures +the time until that event is searchable. A `finally` block starts both consumers +again even when an assertion fails. + +The ordering probe publishes two distinct events for one article, with the older +event deliberately arriving second. It waits for both events to be processed and +confirms the newer title remains searchable. + +Useful inspection commands: + +```bash +docker compose -f deploy/realtime/compose.yaml logs --tail=200 consumer-1 consumer-2 +docker compose -f deploy/realtime/compose.yaml exec kafka \ + /opt/kafka/bin/kafka-consumer-groups.sh \ + --bootstrap-server kafka:9092 \ + --group newslens-indexers \ + --describe +``` + +## DLQ inspection and replay + +Inspect dead-letter records before replaying them: + +```bash +docker compose -f deploy/realtime/compose.yaml exec kafka \ + /opt/kafka/bin/kafka-console-consumer.sh \ + --bootstrap-server kafka:9092 \ + --topic news-events-dlq \ + --from-beginning \ + --max-messages 20 +``` + +The original bytes are base64 encoded because malformed JSON must itself remain +representable in the DLQ envelope. There is intentionally no automatic replay: +an operator must identify the cause, decode and repair the event when appropriate, +assign a new stable event ID, and publish it through `POST /events`. + +## Shutdown and data reset + +Stop processes while retaining the PostgreSQL volume: + +```bash +./scripts/realtime/stop_stack.sh +``` + +To delete the local real-time database as well, explicitly target this Compose +project and its volumes: + +```bash +docker compose -f deploy/realtime/compose.yaml down --volumes +``` + +That last command is destructive and removes only the named Compose stack's +local PostgreSQL volume. It does not affect the licensed MIND archives or the +separate recommendation artifact. diff --git a/pyproject.toml b/pyproject.toml index 495860b..5a6d01c 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -20,6 +20,7 @@ dependencies = [ "scipy>=1.13", "uvicorn>=0.30,<1.0", "joblib>=1.4,<2.0", + "psycopg[binary]>=3.2,<4.0", ] [project.optional-dependencies] diff --git a/reports/realtime_dlq_v0_1.json b/reports/realtime_dlq_v0_1.json new file mode 100644 index 0000000..dde0eea --- /dev/null +++ b/reports/realtime_dlq_v0_1.json @@ -0,0 +1,11 @@ +{ + "schema_version": "newslens.realtime-dlq.v1", + "generated_at": "2026-08-23T20:24:14.971361+00:00", + "environment": "local Docker Compose", + "invalid_events_observed": 1, + "dead_letter_events_observed": 1, + "dead_letter_envelope_valid": true, + "limitations": [ + "This probes malformed JSON; database-exhaustion DLQ behavior is unit tested." + ] +} diff --git a/reports/realtime_load_v0_1.json b/reports/realtime_load_v0_1.json new file mode 100644 index 0000000..15c144d --- /dev/null +++ b/reports/realtime_load_v0_1.json @@ -0,0 +1,38 @@ +{ + "schema_version": "newslens.realtime-load.v1", + "generated_at": "2026-08-23T20:22:13.981615+00:00", + "environment": "local Docker Compose", + "configuration": { + "events": 500, + "concurrency": 20, + "freshness_samples": 25, + "duplicate_probes": 25 + }, + "publish": { + "accepted": 500, + "failed": 0, + "failure_rate": 0.0, + "throughput_events_per_second": 1225.46579188993, + "latency_ms": { + "mean": 16.248251977958716, + "p50": 15.150916500715539, + "p95": 20.018879147392, + "p99": 44.12482450177776 + } + }, + "produced_to_searchable_ms": { + "successful_samples": 25, + "failed_samples": 0, + "p50": 62.514, + "p95": 78.7624, + "p99": 82.07351999999999 + }, + "idempotency": { + "duplicate_probes_accepted": 25, + "duplicates_observed_by_consumers": 25 + }, + "limitations": [ + "This is a single-machine Docker Compose result, not a production capacity claim.", + "Freshness is sampled and includes client polling time only through the stored produced-to-indexed timestamp." + ] +} diff --git a/reports/realtime_ordering_v0_1.json b/reports/realtime_ordering_v0_1.json new file mode 100644 index 0000000..0793715 --- /dev/null +++ b/reports/realtime_ordering_v0_1.json @@ -0,0 +1,11 @@ +{ + "schema_version": "newslens.realtime-ordering.v1", + "generated_at": "2026-08-23T20:27:13.738678+00:00", + "environment": "local Docker Compose", + "distinct_events_processed": 2, + "newer_version_preserved": true, + "observed_title": "Ordering70a9d1542b9b4b2a9e3a05ce68b65828 Current", + "limitations": [ + "Producer clocks must be trustworthy because ordering uses produced_at." + ] +} diff --git a/reports/realtime_recovery_v0_1.json b/reports/realtime_recovery_v0_1.json new file mode 100644 index 0000000..93d8436 --- /dev/null +++ b/reports/realtime_recovery_v0_1.json @@ -0,0 +1,11 @@ +{ + "schema_version": "newslens.realtime-recovery.v1", + "stopped_consumer_partition": 1, + "one_consumer_stopped_searchable_ms": 5668.895958006033, + "full_consumer_outage_recovery_ms": 2429.43025000568, + "generated_at": "2026-08-23T20:23:33.563047+00:00", + "environment": "local Docker Compose", + "limitations": [ + "Single-host process failure exercise; broker and host failures are out of scope." + ] +} diff --git a/reports/realtime_search_evaluation_v0_1.json b/reports/realtime_search_evaluation_v0_1.json new file mode 100644 index 0000000..b6017ac --- /dev/null +++ b/reports/realtime_search_evaluation_v0_1.json @@ -0,0 +1,187 @@ +{ + "schema_version": "newslens.realtime-search-evaluation.v1", + "generated_at": "2026-08-23T20:04:08.761624+00:00", + "fixture": { + "queries": 6, + "articles": 10, + "reference_time": "2026-08-23T12:00:00+00:00", + "source": "small, hand-labeled deterministic contract fixture" + }, + "query_understanding": { + "field_accuracy": 1.0, + "correct_fields": 18, + "total_fields": 18, + "cases": [ + { + "query": "latest Apple AI chip", + "expected": { + "category": "technology", + "entity": "Apple", + "prefers_freshness": true + }, + "observed": { + "category": "technology", + "entity": "Apple", + "prefers_freshness": true + } + }, + { + "query": "Apple AI chip analysis", + "expected": { + "category": "technology", + "entity": "Apple", + "prefers_freshness": false + }, + "observed": { + "category": "technology", + "entity": "Apple", + "prefers_freshness": false + } + }, + { + "query": "NASA recent Mars research", + "expected": { + "category": "science", + "entity": "NASA Mars", + "prefers_freshness": true + }, + "observed": { + "category": "science", + "entity": "NASA Mars", + "prefers_freshness": true + } + }, + { + "query": "Falcons football score tonight", + "expected": { + "category": "sports", + "entity": "Falcons", + "prefers_freshness": true + }, + "observed": { + "category": "sports", + "entity": "Falcons", + "prefers_freshness": true + } + }, + { + "query": "latest Microsoft market earnings", + "expected": { + "category": "business", + "entity": "Microsoft", + "prefers_freshness": true + }, + "observed": { + "category": "business", + "entity": "Microsoft", + "prefers_freshness": true + } + }, + { + "query": "Pfizer vaccine research", + "expected": { + "category": "health", + "entity": "Pfizer", + "prefers_freshness": false + }, + "observed": { + "category": "health", + "entity": "Pfizer", + "prefers_freshness": false + } + } + ] + }, + "ranking": { + "relevance_only": { + "ndcg_at_10": 0.8686305730985519, + "mrr_at_10": 1.0, + "recall_at_10": 1.0, + "freshness_weighted_ndcg_at_10": 0.6335919732394809 + }, + "freshness_aware": { + "ndcg_at_10": 1.0, + "mrr_at_10": 1.0, + "recall_at_10": 1.0, + "freshness_weighted_ndcg_at_10": 0.8775572100937762 + } + }, + "ranking_configuration": { + "relevance_only": { + "relevance": 0.95, + "freshness": 0.0, + "popularity": 0.05 + }, + "freshness_aware": { + "relevance": 0.75, + "freshness": 0.2, + "popularity": 0.05 + } + }, + "cases": [ + { + "query": "latest Apple AI chip", + "category": "technology", + "entity": "Apple", + "prefers_freshness": true, + "relevance": { + "tech-new": 3, + "tech-old": 2 + } + }, + { + "query": "Apple AI chip analysis", + "category": "technology", + "entity": "Apple", + "prefers_freshness": false, + "relevance": { + "tech-old": 3, + "tech-new": 2 + } + }, + { + "query": "NASA recent Mars research", + "category": "science", + "entity": "NASA Mars", + "prefers_freshness": true, + "relevance": { + "space-new": 3, + "space-old": 2 + } + }, + { + "query": "Falcons football score tonight", + "category": "sports", + "entity": "Falcons", + "prefers_freshness": true, + "relevance": { + "sport-new": 3, + "sport-old": 1 + } + }, + { + "query": "latest Microsoft market earnings", + "category": "business", + "entity": "Microsoft", + "prefers_freshness": true, + "relevance": { + "market-new": 3, + "market-old": 2 + } + }, + { + "query": "Pfizer vaccine research", + "category": "health", + "entity": "Pfizer", + "prefers_freshness": false, + "relevance": { + "health-old": 3, + "health-new": 2 + } + } + ], + "limitations": [ + "This fixture verifies ranking behavior; it is not evidence of generalization to live traffic.", + "A larger independently judged query set is still needed before product-quality claims." + ] +} diff --git a/scripts/evaluate_realtime_search.py b/scripts/evaluate_realtime_search.py new file mode 100755 index 0000000..ae0f7bd --- /dev/null +++ b/scripts/evaluate_realtime_search.py @@ -0,0 +1,197 @@ +#!/usr/bin/env python3 +"""Evaluate deterministic query understanding and freshness-aware ranking.""" + +from __future__ import annotations + +import argparse +import json +import math +from dataclasses import asdict, dataclass +from datetime import UTC, datetime, timedelta +from pathlib import Path + +from newslens.realtime import FreshnessRanker, SearchDocument, understand_query + + +@dataclass(frozen=True) +class QueryCase: + query: str + category: str | None + entity: str | None + prefers_freshness: bool + relevance: dict[str, int] + + +def documents(now: datetime) -> tuple[SearchDocument, ...]: + def article( + article_id: str, + title: str, + body: str, + category: str, + age_hours: float, + popularity: int, + ) -> SearchDocument: + published_at = now - timedelta(hours=age_hours) + return SearchDocument( + article_id=article_id, + title=title, + body=body, + category=category, + published_at=published_at, + produced_at=published_at + timedelta(minutes=1), + indexed_at=published_at + timedelta(minutes=1, seconds=2), + popularity=popularity, + ) + + return ( + article("tech-new", "Apple launches AI chip today", "New performance details.", "technology", 1, 15), + article("tech-old", "Apple AI chip performance analysis", "Detailed chip review.", "technology", 168, 900), + article("space-new", "NASA shares latest Mars mission update", "Live mission status.", "science", 2, 20), + article("space-old", "NASA Mars mission research archive", "Long-form mission research.", "science", 240, 700), + article("sport-new", "Falcons football score tonight", "The current game score.", "sports", 0.5, 30), + article("sport-old", "Falcons football season analysis", "Full season analysis.", "sports", 120, 600), + article("market-new", "Microsoft shares rise in latest market", "Fresh earnings reaction.", "business", 3, 25), + article("market-old", "Microsoft market earnings analysis", "Detailed historical analysis.", "business", 96, 800), + article("health-new", "Pfizer vaccine update today", "New hospital guidance.", "health", 4, 10), + article("health-old", "Pfizer vaccine research review", "A detailed medicine review.", "health", 336, 500), + ) + + +CASES = ( + QueryCase("latest Apple AI chip", "technology", "Apple", True, {"tech-new": 3, "tech-old": 2}), + QueryCase("Apple AI chip analysis", "technology", "Apple", False, {"tech-old": 3, "tech-new": 2}), + QueryCase("NASA recent Mars research", "science", "NASA Mars", True, {"space-new": 3, "space-old": 2}), + QueryCase("Falcons football score tonight", "sports", "Falcons", True, {"sport-new": 3, "sport-old": 1}), + QueryCase("latest Microsoft market earnings", "business", "Microsoft", True, {"market-new": 3, "market-old": 2}), + QueryCase("Pfizer vaccine research", "health", "Pfizer", False, {"health-old": 3, "health-new": 2}), +) + + +def dcg(grades: list[float]) -> float: + return sum((2**grade - 1) / math.log2(index + 2) for index, grade in enumerate(grades)) + + +def ndcg(order: list[str], relevance: dict[str, float], cutoff: int = 10) -> float: + observed = [relevance.get(article_id, 0.0) for article_id in order[:cutoff]] + ideal = sorted(relevance.values(), reverse=True)[:cutoff] + denominator = dcg(ideal) + return dcg(observed) / denominator if denominator else 0.0 + + +def reciprocal_rank(order: list[str], relevance: dict[str, int], cutoff: int = 10) -> float: + for rank, article_id in enumerate(order[:cutoff], start=1): + if relevance.get(article_id, 0) > 0: + return 1.0 / rank + return 0.0 + + +def recall(order: list[str], relevance: dict[str, int], cutoff: int = 10) -> float: + relevant = {article_id for article_id, grade in relevance.items() if grade > 0} + if not relevant: + return 0.0 + return len(relevant & set(order[:cutoff])) / len(relevant) + + +def evaluate_ranker( + ranker: FreshnessRanker, + corpus: tuple[SearchDocument, ...], + now: datetime, +) -> dict[str, float]: + ndcgs: list[float] = [] + mrrs: list[float] = [] + recalls: list[float] = [] + freshness_ndcgs: list[float] = [] + by_id = {document.article_id: document for document in corpus} + for case in CASES: + intent = understand_query(case.query) + order = [item.document.article_id for item in ranker.rank(intent, corpus, top_k=10, now=now)] + ndcgs.append(ndcg(order, case.relevance)) + mrrs.append(reciprocal_rank(order, case.relevance)) + recalls.append(recall(order, case.relevance)) + freshness_relevance = { + article_id: grade + * math.exp( + -math.log(2) + * max(0.0, (now - by_id[article_id].published_at).total_seconds() / 3600) + / 24.0 + ) + for article_id, grade in case.relevance.items() + } + freshness_ndcgs.append(ndcg(order, freshness_relevance)) + return { + "ndcg_at_10": sum(ndcgs) / len(ndcgs), + "mrr_at_10": sum(mrrs) / len(mrrs), + "recall_at_10": sum(recalls) / len(recalls), + "freshness_weighted_ndcg_at_10": sum(freshness_ndcgs) / len(freshness_ndcgs), + } + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--output", type=Path, default=Path("reports/realtime_search_evaluation_v0_1.json")) + args = parser.parse_args() + now = datetime(2026, 8, 23, 12, 0, tzinfo=UTC) + corpus = documents(now) + + intent_results = [] + correct_fields = 0 + total_fields = 0 + for case in CASES: + actual = understand_query(case.query) + expected = { + "category": case.category, + "entity": case.entity, + "prefers_freshness": case.prefers_freshness, + } + observed = { + "category": actual.category, + "entity": actual.entity, + "prefers_freshness": actual.prefers_freshness, + } + correct_fields += sum(expected[key] == observed[key] for key in expected) + total_fields += len(expected) + intent_results.append({"query": case.query, "expected": expected, "observed": observed}) + + relevance_only = FreshnessRanker( + relevance_weight=0.95, + freshness_weight=0.0, + popularity_weight=0.05, + ) + freshness_aware = FreshnessRanker() + report = { + "schema_version": "newslens.realtime-search-evaluation.v1", + "generated_at": datetime.now(UTC).isoformat(), + "fixture": { + "queries": len(CASES), + "articles": len(corpus), + "reference_time": now.isoformat(), + "source": "small, hand-labeled deterministic contract fixture", + }, + "query_understanding": { + "field_accuracy": correct_fields / total_fields, + "correct_fields": correct_fields, + "total_fields": total_fields, + "cases": intent_results, + }, + "ranking": { + "relevance_only": evaluate_ranker(relevance_only, corpus, now), + "freshness_aware": evaluate_ranker(freshness_aware, corpus, now), + }, + "ranking_configuration": { + "relevance_only": {"relevance": 0.95, "freshness": 0.0, "popularity": 0.05}, + "freshness_aware": {"relevance": 0.75, "freshness": 0.20, "popularity": 0.05}, + }, + "cases": [asdict(case) for case in CASES], + "limitations": [ + "This fixture verifies ranking behavior; it is not evidence of generalization to live traffic.", + "A larger independently judged query set is still needed before product-quality claims.", + ], + } + args.output.parent.mkdir(parents=True, exist_ok=True) + args.output.write_text(json.dumps(report, indent=2) + "\n", encoding="utf-8") + print(json.dumps(report, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/realtime/benchmark_ingestion.py b/scripts/realtime/benchmark_ingestion.py new file mode 100755 index 0000000..8f8c163 --- /dev/null +++ b/scripts/realtime/benchmark_ingestion.py @@ -0,0 +1,210 @@ +#!/usr/bin/env python3 +"""Measure publish throughput and produced-to-searchable freshness.""" + +from __future__ import annotations + +import argparse +import json +import math +import statistics +import time +import urllib.error +import urllib.parse +import urllib.request +import uuid +from concurrent.futures import ThreadPoolExecutor, as_completed +from datetime import UTC, datetime +from pathlib import Path +from typing import Any + + +def percentile(values: list[float], percentage: float) -> float: + if not values: + return 0.0 + ordered = sorted(values) + position = (len(ordered) - 1) * percentage + lower = math.floor(position) + upper = math.ceil(position) + if lower == upper: + return ordered[lower] + return ordered[lower] + (ordered[upper] - ordered[lower]) * (position - lower) + + +def request_json(url: str, *, payload: dict[str, Any] | None = None) -> tuple[int, Any]: + data = json.dumps(payload).encode() if payload is not None else None + request = urllib.request.Request( + url, + data=data, + headers={"Content-Type": "application/json"} if data else {}, + method="POST" if data else "GET", + ) + try: + with urllib.request.urlopen(request, timeout=10) as response: + return response.status, json.load(response) + except urllib.error.HTTPError as error: + response = error.read().decode(errors="replace") + try: + return error.code, json.loads(response) + except json.JSONDecodeError: + return error.code, {"error": response} + + +def publish(url: str, payload: dict[str, Any]) -> tuple[int, float]: + started = time.perf_counter() + status, _ = request_json(f"{url}/events", payload=payload) + return status, (time.perf_counter() - started) * 1_000 + + +def wait_until_searchable(search_url: str, token: str, timeout: float) -> dict[str, Any]: + deadline = time.monotonic() + timeout + encoded = urllib.parse.urlencode({"q": token, "top_k": 1}) + while time.monotonic() < deadline: + status, body = request_json(f"{search_url}/search?{encoded}") + if status == 200 and body["returned_count"]: + return body["results"][0] + time.sleep(0.05) + raise TimeoutError(f"article with token {token!r} did not become searchable") + + +def metric_value(url: str, name: str) -> int: + with urllib.request.urlopen(f"{url}/metrics", timeout=10) as response: + for line in response.read().decode().splitlines(): + if line.startswith(f"{name} "): + return int(float(line.split()[1])) + raise ValueError(f"metric {name!r} was not found at {url}") + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--events", type=int, default=500) + parser.add_argument("--concurrency", type=int, default=20) + parser.add_argument("--freshness-samples", type=int, default=25) + parser.add_argument("--duplicate-probes", type=int, default=25) + parser.add_argument("--ingestion-url", default="http://127.0.0.1:8080") + parser.add_argument("--search-url", default="http://127.0.0.1:8000") + parser.add_argument( + "--consumer-url", + action="append", + default=["http://127.0.0.1:8081", "http://127.0.0.1:8082"], + ) + parser.add_argument("--output", type=Path, default=Path("reports/realtime_load_v0_1.json")) + args = parser.parse_args() + if ( + args.events < 1 + or args.concurrency < 1 + or args.freshness_samples < 1 + or args.duplicate_probes < 0 + ): + parser.error("event, concurrency, and sample counts must be positive") + + run_id = uuid.uuid4().hex[:10] + published_at = datetime.now(UTC).isoformat() + payloads = [] + for index in range(args.events): + token = f"Loadtest{run_id}{index}" + payloads.append( + { + "event_id": f"load-{run_id}-{index}", + "article_id": f"load-{run_id}-{index}", + "title": f"{token} technology update", + "category": "technology", + "published_at": published_at, + "body": "A deterministic article used to measure ingestion behavior.", + "produced_at": datetime.now(UTC).isoformat(), + } + ) + + started = time.perf_counter() + latencies: list[float] = [] + statuses: list[int] = [] + with ThreadPoolExecutor(max_workers=args.concurrency) as executor: + futures = [executor.submit(publish, args.ingestion_url, payload) for payload in payloads] + for future in as_completed(futures): + status, latency = future.result() + statuses.append(status) + latencies.append(latency) + duration = time.perf_counter() - started + + sample_count = min(args.freshness_samples, len(payloads)) + freshness: list[float] = [] + search_failures = 0 + for payload in payloads[:sample_count]: + token = payload["title"].split()[0] + try: + result = wait_until_searchable(args.search_url, token, timeout=30.0) + freshness.append(float(result["index_freshness_ms"])) + except TimeoutError: + search_failures += 1 + + duplicate_count = min(args.duplicate_probes, len(payloads)) + duplicate_metric = "newslens_ingestion_duplicates_total" + duplicates_before = sum( + metric_value(url, duplicate_metric) for url in args.consumer_url + ) + duplicate_statuses = [ + publish(args.ingestion_url, payload)[0] for payload in payloads[:duplicate_count] + ] + duplicate_deadline = time.monotonic() + 30 + duplicates_observed = 0 + while time.monotonic() < duplicate_deadline: + duplicates_after = sum( + metric_value(url, duplicate_metric) for url in args.consumer_url + ) + duplicates_observed = duplicates_after - duplicates_before + if duplicates_observed >= duplicate_count: + break + time.sleep(0.1) + + accepted = sum(status == 202 for status in statuses) + report = { + "schema_version": "newslens.realtime-load.v1", + "generated_at": datetime.now(UTC).isoformat(), + "environment": "local Docker Compose", + "configuration": { + "events": args.events, + "concurrency": args.concurrency, + "freshness_samples": sample_count, + "duplicate_probes": duplicate_count, + }, + "publish": { + "accepted": accepted, + "failed": len(statuses) - accepted, + "failure_rate": (len(statuses) - accepted) / len(statuses), + "throughput_events_per_second": accepted / duration, + "latency_ms": { + "mean": statistics.fmean(latencies), + "p50": percentile(latencies, 0.50), + "p95": percentile(latencies, 0.95), + "p99": percentile(latencies, 0.99), + }, + }, + "produced_to_searchable_ms": { + "successful_samples": len(freshness), + "failed_samples": search_failures, + "p50": percentile(freshness, 0.50), + "p95": percentile(freshness, 0.95), + "p99": percentile(freshness, 0.99), + }, + "idempotency": { + "duplicate_probes_accepted": sum(status == 202 for status in duplicate_statuses), + "duplicates_observed_by_consumers": duplicates_observed, + }, + "limitations": [ + "This is a single-machine Docker Compose result, not a production capacity claim.", + "Freshness is sampled and includes client polling time only through the stored produced-to-indexed timestamp.", + ], + } + args.output.parent.mkdir(parents=True, exist_ok=True) + args.output.write_text(json.dumps(report, indent=2) + "\n", encoding="utf-8") + print(json.dumps(report, indent=2)) + successful = ( + accepted == args.events + and search_failures == 0 + and all(status == 202 for status in duplicate_statuses) + and duplicates_observed >= duplicate_count + ) + return 0 if successful else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/realtime/failure_recovery.py b/scripts/realtime/failure_recovery.py new file mode 100755 index 0000000..bd0b9fb --- /dev/null +++ b/scripts/realtime/failure_recovery.py @@ -0,0 +1,147 @@ +#!/usr/bin/env python3 +"""Measure consumer failover and backlog recovery in the local Compose stack.""" + +from __future__ import annotations + +import argparse +import json +import subprocess +import time +import urllib.parse +import urllib.request +import uuid +from datetime import UTC, datetime +from pathlib import Path +from typing import Any + + +def compose(file: Path, *arguments: str) -> None: + subprocess.run(["docker", "compose", "-f", str(file), *arguments], check=True) + + +def consumer_partitions(file: Path, client_id: str) -> set[int]: + command = [ + "docker", + "compose", + "-f", + str(file), + "exec", + "-T", + "kafka", + "/opt/kafka/bin/kafka-consumer-groups.sh", + "--bootstrap-server", + "kafka:9092", + "--group", + "newslens-indexers", + "--describe", + ] + output = subprocess.run(command, check=True, capture_output=True, text=True).stdout + partitions = set() + for line in output.splitlines(): + fields = line.split() + if len(fields) >= 8 and fields[1] == "news-events" and fields[-1] == client_id: + partitions.add(int(fields[2])) + if not partitions: + raise RuntimeError(f"no partition assignment found for {client_id}") + return partitions + + +def fnv_partition(key: str, count: int = 3) -> int: + value = 2_166_136_261 + for byte in key.encode(): + value ^= byte + value = (value * 16_777_619) & 0xFFFFFFFF + signed = value if value < 0x80000000 else value - 0x100000000 + return abs(signed) % count + + +def request_json(url: str, payload: dict[str, Any] | None = None) -> Any: + data = json.dumps(payload).encode() if payload else None + request = urllib.request.Request( + url, + data=data, + headers={"Content-Type": "application/json"} if data else {}, + method="POST" if data else "GET", + ) + with urllib.request.urlopen(request, timeout=10) as response: + return json.load(response) + + +def event(label: str, target_partition: int | None = None) -> tuple[dict[str, Any], str]: + identifier = uuid.uuid4().hex + while target_partition is not None and fnv_partition(f"{label}-{identifier}") != target_partition: + identifier = uuid.uuid4().hex + token = f"Recovery{identifier}" + now = datetime.now(UTC).isoformat() + return ( + { + "event_id": f"{label}-{identifier}", + "article_id": f"{label}-{identifier}", + "title": f"{token} technology update", + "category": "technology", + "published_at": now, + "produced_at": now, + "body": "Failure-recovery probe.", + }, + token, + ) + + +def wait_search(search_url: str, token: str, timeout: float = 60.0) -> float: + started = time.perf_counter() + query = urllib.parse.urlencode({"q": token, "top_k": 1}) + while time.perf_counter() - started < timeout: + body = request_json(f"{search_url}/search?{query}") + if body["returned_count"]: + return (time.perf_counter() - started) * 1_000 + time.sleep(0.1) + raise TimeoutError(f"{token} did not become searchable") + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--compose-file", type=Path, default=Path("deploy/realtime/compose.yaml")) + parser.add_argument("--ingestion-url", default="http://127.0.0.1:8080") + parser.add_argument("--search-url", default="http://127.0.0.1:8000") + parser.add_argument("--output", type=Path, default=Path("reports/realtime_recovery_v0_1.json")) + args = parser.parse_args() + + report: dict[str, Any] = {"schema_version": "newslens.realtime-recovery.v1"} + try: + target_partition = min( + consumer_partitions(args.compose_file, "newslens-consumer-1") + ) + compose(args.compose_file, "stop", "consumer-1") + failover_event, token = event("failover", target_partition) + request_json(f"{args.ingestion_url}/events", failover_event) + report["stopped_consumer_partition"] = target_partition + report["one_consumer_stopped_searchable_ms"] = wait_search(args.search_url, token) + + compose(args.compose_file, "stop", "consumer-2") + backlog_event, token = event("backlog") + request_json(f"{args.ingestion_url}/events", backlog_event) + time.sleep(2) + recovery_started = time.perf_counter() + compose(args.compose_file, "start", "consumer-1") + wait_search(args.search_url, token) + report["full_consumer_outage_recovery_ms"] = ( + time.perf_counter() - recovery_started + ) * 1_000 + finally: + compose(args.compose_file, "start", "consumer-1", "consumer-2") + + report.update( + { + "generated_at": datetime.now(UTC).isoformat(), + "environment": "local Docker Compose", + "limitations": ["Single-host process failure exercise; broker and host failures are out of scope."], + } + ) + args.output.parent.mkdir(parents=True, exist_ok=True) + args.output.write_text(json.dumps(report, indent=2) + "\n", encoding="utf-8") + print(json.dumps(report, indent=2)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/realtime/health_check.sh b/scripts/realtime/health_check.sh new file mode 100755 index 0000000..716e922 --- /dev/null +++ b/scripts/realtime/health_check.sh @@ -0,0 +1,24 @@ +#!/usr/bin/env bash +set -euo pipefail + +check_json_endpoint() { + local label="$1" + local url="$2" + local attempts=30 + + for ((attempt = 1; attempt <= attempts; attempt++)); do + if curl --fail --silent --show-error "$url" >/dev/null; then + echo "$label is ready: $url" + return 0 + fi + sleep 2 + done + + echo "$label did not become ready: $url" >&2 + return 1 +} + +check_json_endpoint "ingestion API" "${NEWSLENS_INGESTION_URL:-http://127.0.0.1:8080}/ready" +check_json_endpoint "consumer 1" "${NEWSLENS_CONSUMER_1_URL:-http://127.0.0.1:8081}/ready" +check_json_endpoint "consumer 2" "${NEWSLENS_CONSUMER_2_URL:-http://127.0.0.1:8082}/ready" +check_json_endpoint "search API" "${NEWSLENS_SEARCH_URL:-http://127.0.0.1:8000}/realtime/ready" diff --git a/scripts/realtime/start_stack.sh b/scripts/realtime/start_stack.sh new file mode 100755 index 0000000..ac07d54 --- /dev/null +++ b/scripts/realtime/start_stack.sh @@ -0,0 +1,8 @@ +#!/usr/bin/env bash +set -euo pipefail + +repository_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)" +compose_file="$repository_root/deploy/realtime/compose.yaml" + +docker compose -f "$compose_file" up --detach --build --wait +"$repository_root/scripts/realtime/health_check.sh" diff --git a/scripts/realtime/stop_stack.sh b/scripts/realtime/stop_stack.sh new file mode 100755 index 0000000..95ec8d9 --- /dev/null +++ b/scripts/realtime/stop_stack.sh @@ -0,0 +1,7 @@ +#!/usr/bin/env bash +set -euo pipefail + +repository_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)" +compose_file="$repository_root/deploy/realtime/compose.yaml" + +docker compose -f "$compose_file" down diff --git a/scripts/realtime/verify_dlq.py b/scripts/realtime/verify_dlq.py new file mode 100755 index 0000000..5afb621 --- /dev/null +++ b/scripts/realtime/verify_dlq.py @@ -0,0 +1,118 @@ +#!/usr/bin/env python3 +"""Inject malformed Kafka bytes and verify the dead-letter contract.""" + +from __future__ import annotations + +import argparse +import base64 +import json +import subprocess +import time +import urllib.request +from datetime import UTC, datetime +from pathlib import Path + + +def metric_value(url: str, name: str) -> int: + with urllib.request.urlopen(f"{url}/metrics", timeout=10) as response: + for line in response.read().decode().splitlines(): + if line.startswith(f"{name} "): + return int(float(line.split()[1])) + raise ValueError(f"metric {name!r} was not found at {url}") + + +def compose_command(file: Path, *arguments: str) -> list[str]: + return ["docker", "compose", "-f", str(file), *arguments] + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--compose-file", type=Path, default=Path("deploy/realtime/compose.yaml")) + parser.add_argument( + "--consumer-url", + action="append", + default=["http://127.0.0.1:8081", "http://127.0.0.1:8082"], + ) + parser.add_argument("--output", type=Path, default=Path("reports/realtime_dlq_v0_1.json")) + args = parser.parse_args() + + invalid_metric = "newslens_ingestion_invalid_total" + dlq_metric = "newslens_ingestion_dead_letter_total" + invalid_before = sum(metric_value(url, invalid_metric) for url in args.consumer_url) + dlq_before = sum(metric_value(url, dlq_metric) for url in args.consumer_url) + + malformed = b"not-json" + subprocess.run( + compose_command( + args.compose_file, + "exec", + "-T", + "kafka", + "/opt/kafka/bin/kafka-console-producer.sh", + "--bootstrap-server", + "kafka:9092", + "--topic", + "news-events", + ), + input=malformed + b"\n", + check=True, + ) + + deadline = time.monotonic() + 30 + invalid_delta = 0 + dlq_delta = 0 + while time.monotonic() < deadline: + invalid_delta = ( + sum(metric_value(url, invalid_metric) for url in args.consumer_url) + - invalid_before + ) + dlq_delta = sum(metric_value(url, dlq_metric) for url in args.consumer_url) - dlq_before + if invalid_delta >= 1 and dlq_delta >= 1: + break + time.sleep(0.1) + + consumed = subprocess.run( + compose_command( + args.compose_file, + "exec", + "-T", + "kafka", + "/opt/kafka/bin/kafka-console-consumer.sh", + "--bootstrap-server", + "kafka:9092", + "--topic", + "news-events-dlq", + "--from-beginning", + "--max-messages", + "1", + "--timeout-ms", + "10000", + ), + capture_output=True, + check=True, + text=True, + ) + envelope = json.loads(consumed.stdout.strip().splitlines()[0]) + expected_payload = base64.b64encode(malformed).decode() + envelope_valid = ( + envelope.get("source") == "news-events" + and envelope.get("original_payload_base64") == expected_payload + and bool(envelope.get("reason")) + ) + report = { + "schema_version": "newslens.realtime-dlq.v1", + "generated_at": datetime.now(UTC).isoformat(), + "environment": "local Docker Compose", + "invalid_events_observed": invalid_delta, + "dead_letter_events_observed": dlq_delta, + "dead_letter_envelope_valid": envelope_valid, + "limitations": ["This probes malformed JSON; database-exhaustion DLQ behavior is unit tested."], + } + args.output.parent.mkdir(parents=True, exist_ok=True) + args.output.write_text(json.dumps(report, indent=2) + "\n", encoding="utf-8") + print(json.dumps(report, indent=2)) + return 0 if invalid_delta >= 1 and dlq_delta >= 1 and envelope_valid else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/scripts/realtime/verify_ordering.py b/scripts/realtime/verify_ordering.py new file mode 100755 index 0000000..d196575 --- /dev/null +++ b/scripts/realtime/verify_ordering.py @@ -0,0 +1,109 @@ +#!/usr/bin/env python3 +"""Verify that an older article event cannot overwrite a newer version.""" + +from __future__ import annotations + +import argparse +import json +import time +import urllib.parse +import urllib.request +import uuid +from datetime import UTC, datetime, timedelta +from pathlib import Path +from typing import Any + + +def request_json(url: str, payload: dict[str, Any] | None = None) -> Any: + data = json.dumps(payload).encode() if payload else None + request = urllib.request.Request( + url, + data=data, + headers={"Content-Type": "application/json"} if data else {}, + method="POST" if data else "GET", + ) + with urllib.request.urlopen(request, timeout=10) as response: + return json.load(response) + + +def metric_value(url: str, name: str) -> int: + with urllib.request.urlopen(f"{url}/metrics", timeout=10) as response: + for line in response.read().decode().splitlines(): + if line.startswith(f"{name} "): + return int(float(line.split()[1])) + raise ValueError(f"metric {name!r} was not found at {url}") + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--ingestion-url", default="http://127.0.0.1:8080") + parser.add_argument("--search-url", default="http://127.0.0.1:8000") + parser.add_argument( + "--consumer-url", + action="append", + default=["http://127.0.0.1:8081", "http://127.0.0.1:8082"], + ) + parser.add_argument("--output", type=Path, default=Path("reports/realtime_ordering_v0_1.json")) + args = parser.parse_args() + + identifier = uuid.uuid4().hex + token = f"Ordering{identifier}" + article_id = f"ordering-{identifier}" + now = datetime.now(UTC) + base = { + "article_id": article_id, + "category": "technology", + "published_at": now.isoformat(), + "body": "Event ordering probe.", + } + current = { + **base, + "event_id": f"ordering-current-{identifier}", + "title": f"{token} Current", + "produced_at": now.isoformat(), + } + stale = { + **base, + "event_id": f"ordering-stale-{identifier}", + "title": f"{token} Stale", + "produced_at": (now - timedelta(days=1)).isoformat(), + } + processed_metric = "newslens_ingestion_processed_total" + processed_before = sum( + metric_value(url, processed_metric) for url in args.consumer_url + ) + request_json(f"{args.ingestion_url}/events", current) + request_json(f"{args.ingestion_url}/events", stale) + + deadline = time.monotonic() + 30 + processed_delta = 0 + while time.monotonic() < deadline: + processed_delta = ( + sum(metric_value(url, processed_metric) for url in args.consumer_url) + - processed_before + ) + if processed_delta >= 2: + break + time.sleep(0.1) + + query = urllib.parse.urlencode({"q": token, "top_k": 1}) + search = request_json(f"{args.search_url}/search?{query}") + title = search["results"][0]["title"] if search["results"] else None + current_preserved = title == current["title"] + report = { + "schema_version": "newslens.realtime-ordering.v1", + "generated_at": datetime.now(UTC).isoformat(), + "environment": "local Docker Compose", + "distinct_events_processed": processed_delta, + "newer_version_preserved": current_preserved, + "observed_title": title, + "limitations": ["Producer clocks must be trustworthy because ordering uses produced_at."], + } + args.output.parent.mkdir(parents=True, exist_ok=True) + args.output.write_text(json.dumps(report, indent=2) + "\n", encoding="utf-8") + print(json.dumps(report, indent=2)) + return 0 if processed_delta >= 2 and current_preserved else 1 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/services/ingestion/Dockerfile b/services/ingestion/Dockerfile new file mode 100644 index 0000000..3c8c7bf --- /dev/null +++ b/services/ingestion/Dockerfile @@ -0,0 +1,12 @@ +FROM golang:1.23-bookworm AS build + +WORKDIR /src +COPY services/ingestion/go.mod services/ingestion/go.sum* ./ +RUN go mod download +COPY services/ingestion/ ./ +RUN CGO_ENABLED=0 GOOS=linux go build -trimpath -ldflags="-s -w" -o /out/newslens-ingestion ./cmd/newslens-ingestion + +FROM gcr.io/distroless/static-debian12:nonroot +COPY --from=build /out/newslens-ingestion /newslens-ingestion +EXPOSE 8080 +ENTRYPOINT ["/newslens-ingestion"] diff --git a/services/ingestion/cmd/newslens-ingestion/main.go b/services/ingestion/cmd/newslens-ingestion/main.go new file mode 100644 index 0000000..56facd6 --- /dev/null +++ b/services/ingestion/cmd/newslens-ingestion/main.go @@ -0,0 +1,148 @@ +package main + +import ( + "context" + "errors" + "fmt" + "log/slog" + "net/http" + "os" + "os/signal" + "strconv" + "strings" + "syscall" + "time" + + "github.com/triasha72/NewsLens/services/ingestion/internal/adapters" + "github.com/triasha72/NewsLens/services/ingestion/internal/ingestion" +) + +type config struct { + mode string + address string + brokers []string + topic string + dlqTopic string + groupID string + clientID string + databaseURL string + retries int +} + +func environment(name, fallback string) string { + if value := strings.TrimSpace(os.Getenv(name)); value != "" { + return value + } + return fallback +} + +func loadConfig() (config, error) { + retries, err := strconv.Atoi(environment("NEWSLENS_PROCESSING_RETRIES", "3")) + if err != nil || retries < 1 || retries > 10 { + return config{}, errors.New("NEWSLENS_PROCESSING_RETRIES must be between 1 and 10") + } + configured := config{ + mode: environment("NEWSLENS_INGESTION_MODE", "all"), + address: environment("NEWSLENS_HTTP_ADDRESS", ":8080"), + brokers: strings.Split(environment("NEWSLENS_KAFKA_BROKERS", "localhost:9092"), ","), + topic: environment("NEWSLENS_KAFKA_TOPIC", "news-events"), + dlqTopic: environment("NEWSLENS_KAFKA_DLQ_TOPIC", "news-events-dlq"), + groupID: environment("NEWSLENS_KAFKA_GROUP_ID", "newslens-indexers"), + clientID: environment("NEWSLENS_KAFKA_CLIENT_ID", "newslens-ingestion"), + databaseURL: strings.TrimSpace(os.Getenv("NEWSLENS_DATABASE_URL")), + retries: retries, + } + if configured.mode != "api" && configured.mode != "consumer" && configured.mode != "all" { + return config{}, errors.New("NEWSLENS_INGESTION_MODE must be api, consumer, or all") + } + if configured.mode != "api" && configured.databaseURL == "" { + return config{}, errors.New("NEWSLENS_DATABASE_URL is required in consumer and all modes") + } + return configured, nil +} + +func run(ctx context.Context, configured config, logger *slog.Logger) error { + metrics := &ingestion.Metrics{} + var publisher ingestion.EventPublisher + var source ingestion.EventSource + var deadLetter ingestion.DeadLetterPublisher + var store ingestion.ArticleStore + + if configured.mode == "api" || configured.mode == "all" { + publisher = adapters.NewKafkaPublisher(configured.brokers, configured.topic) + defer publisher.Close() + } + if configured.mode == "consumer" || configured.mode == "all" { + storeContext, cancel := context.WithTimeout(ctx, 10*time.Second) + defer cancel() + postgresStore, err := adapters.NewPostgresStore(storeContext, configured.databaseURL) + if err != nil { + return err + } + store = postgresStore + defer store.Close() + + source = adapters.NewKafkaSource( + configured.brokers, + configured.topic, + configured.groupID, + configured.clientID, + ) + defer source.Close() + deadLetter = adapters.NewKafkaDeadLetterPublisher(configured.brokers, configured.dlqTopic) + defer deadLetter.Close() + } + + server := &http.Server{ + Addr: configured.address, + Handler: ingestion.NewHTTPHandler(publisher, store, metrics, logger), + ReadHeaderTimeout: 5 * time.Second, + ReadTimeout: 10 * time.Second, + WriteTimeout: 10 * time.Second, + IdleTimeout: 60 * time.Second, + } + result := make(chan error, 2) + go func() { + logger.Info("HTTP service started", "address", configured.address, "mode", configured.mode) + result <- server.ListenAndServe() + }() + + if source != nil && store != nil && deadLetter != nil { + worker := ingestion.NewWorker( + source, + store, + deadLetter, + metrics, + logger, + configured.retries, + ) + go func() { result <- worker.Run(ctx) }() + } + + select { + case <-ctx.Done(): + shutdownContext, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + return server.Shutdown(shutdownContext) + case err := <-result: + if errors.Is(err, http.ErrServerClosed) { + return nil + } + return err + } +} + +func main() { + logger := slog.New(slog.NewJSONHandler(os.Stdout, nil)) + configured, err := loadConfig() + if err != nil { + logger.Error("invalid configuration", "error", err) + os.Exit(2) + } + ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM) + defer stop() + if err := run(ctx, configured, logger); err != nil { + logger.Error("ingestion service stopped", "error", fmt.Sprintf("%v", err)) + os.Exit(1) + } +} diff --git a/services/ingestion/go.mod b/services/ingestion/go.mod new file mode 100644 index 0000000..76c1928 --- /dev/null +++ b/services/ingestion/go.mod @@ -0,0 +1,19 @@ +module github.com/triasha72/NewsLens/services/ingestion + +go 1.23.0 + +require ( + github.com/jackc/pgx/v5 v5.7.6 + github.com/segmentio/kafka-go v0.4.49 +) + +require ( + github.com/jackc/pgpassfile v1.0.0 // indirect + github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 // indirect + github.com/jackc/puddle/v2 v2.2.2 // indirect + github.com/klauspost/compress v1.15.9 // indirect + github.com/pierrec/lz4/v4 v4.1.15 // indirect + golang.org/x/crypto v0.37.0 // indirect + golang.org/x/sync v0.13.0 // indirect + golang.org/x/text v0.24.0 // indirect +) diff --git a/services/ingestion/go.sum b/services/ingestion/go.sum new file mode 100644 index 0000000..3edee80 --- /dev/null +++ b/services/ingestion/go.sum @@ -0,0 +1,42 @@ +github.com/davecgh/go-spew v1.1.0/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/davecgh/go-spew v1.1.1 h1:vj9j/u1bqnvCEfJOwUhtlOARqs3+rkHYY13jYWTU97c= +github.com/davecgh/go-spew v1.1.1/go.mod h1:J7Y8YcW2NihsgmVo/mv3lAwl/skON4iLHjSsI+c5H38= +github.com/jackc/pgpassfile v1.0.0 h1:/6Hmqy13Ss2zCq62VdNG8tM1wchn8zjSGOBJ6icpsIM= +github.com/jackc/pgpassfile v1.0.0/go.mod h1:CEx0iS5ambNFdcRtxPj5JhEz+xB6uRky5eyVu/W2HEg= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761 h1:iCEnooe7UlwOQYpKFhBabPMi4aNAfoODPEFNiAnClxo= +github.com/jackc/pgservicefile v0.0.0-20240606120523-5a60cdf6a761/go.mod h1:5TJZWKEWniPve33vlWYSoGYefn3gLQRzjfDlhSJ9ZKM= +github.com/jackc/pgx/v5 v5.7.6 h1:rWQc5FwZSPX58r1OQmkuaNicxdmExaEz5A2DO2hUuTk= +github.com/jackc/pgx/v5 v5.7.6/go.mod h1:aruU7o91Tc2q2cFp5h4uP3f6ztExVpyVv88Xl/8Vl8M= +github.com/jackc/puddle/v2 v2.2.2 h1:PR8nw+E/1w0GLuRFSmiioY6UooMp6KJv0/61nB7icHo= +github.com/jackc/puddle/v2 v2.2.2/go.mod h1:vriiEXHvEE654aYKXXjOvZM39qJ0q+azkZFrfEOc3H4= +github.com/klauspost/compress v1.15.9 h1:wKRjX6JRtDdrE9qwa4b/Cip7ACOshUI4smpCQanqjSY= +github.com/klauspost/compress v1.15.9/go.mod h1:PhcZ0MbTNciWF3rruxRgKxI5NkcHHrHUDtV4Yw2GlzU= +github.com/pierrec/lz4/v4 v4.1.15 h1:MO0/ucJhngq7299dKLwIMtgTfbkoSPF6AoMYDd8Q4q0= +github.com/pierrec/lz4/v4 v4.1.15/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4= +github.com/pmezard/go-difflib v1.0.0 h1:4DBwDE0NGyQoBHbLQYPwSUPoCMWR5BEzIk/f1lZbAQM= +github.com/pmezard/go-difflib v1.0.0/go.mod h1:iKH77koFhYxTK1pcRnkKkqfTogsbg7gZNVY4sRDYZ/4= +github.com/segmentio/kafka-go v0.4.49 h1:GJiNX1d/g+kG6ljyJEoi9++PUMdXGAxb7JGPiDCuNmk= +github.com/segmentio/kafka-go v0.4.49/go.mod h1:Y1gn60kzLEEaW28YshXyk2+VCUKbJ3Qr6DrnT3i4+9E= +github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME= +github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI= +github.com/stretchr/testify v1.7.0/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg= +github.com/stretchr/testify v1.8.1 h1:w7B6lhMri9wdJUVmEZPGGhZzrYTPvgJArz7wNPgYKsk= +github.com/stretchr/testify v1.8.1/go.mod h1:w2LPCIKwWwSfY2zedu0+kehJoqGctiVI29o6fzry7u4= +github.com/xdg-go/pbkdf2 v1.0.0 h1:Su7DPu48wXMwC3bs7MCNG+z4FhcyEuz5dlvchbq0B0c= +github.com/xdg-go/pbkdf2 v1.0.0/go.mod h1:jrpuAogTd400dnrH08LKmI/xc1MbPOebTwRqcT5RDeI= +github.com/xdg-go/scram v1.1.2 h1:FHX5I5B4i4hKRVRBCFRxq1iQRej7WO3hhBuJf+UUySY= +github.com/xdg-go/scram v1.1.2/go.mod h1:RT/sEzTbU5y00aCK8UOx6R7YryM0iF1N2MOmC3kKLN4= +github.com/xdg-go/stringprep v1.0.4 h1:XLI/Ng3O1Atzq0oBs3TWm+5ZVgkq2aqdlvP9JtoZ6c8= +github.com/xdg-go/stringprep v1.0.4/go.mod h1:mPGuuIYwz7CmR2bT9j4GbQqutWS1zV24gijq1dTyGkM= +golang.org/x/crypto v0.37.0 h1:kJNSjF/Xp7kU0iB2Z+9viTPMW4EqqsrywMXLJOOsXSE= +golang.org/x/crypto v0.37.0/go.mod h1:vg+k43peMZ0pUMhYmVAWysMK35e6ioLh3wB8ZCAfbVc= +golang.org/x/net v0.38.0 h1:vRMAPTMaeGqVhG5QyLJHqNDwecKTomGeqbnfZyKlBI8= +golang.org/x/net v0.38.0/go.mod h1:ivrbrMbzFq5J41QOQh0siUuly180yBYtLp+CKbEaFx8= +golang.org/x/sync v0.13.0 h1:AauUjRAJ9OSnvULf/ARrrVywoJDy0YS2AwQ98I37610= +golang.org/x/sync v0.13.0/go.mod h1:1dzgHSNfp02xaA81J2MS99Qcpr2w7fw1gpm99rleRqA= +golang.org/x/text v0.24.0 h1:dd5Bzh4yt5KYA8f9CJHCP4FB4D51c2c6JvN37xJJkJ0= +golang.org/x/text v0.24.0/go.mod h1:L8rBsPeo2pSS+xqN0d5u2ikmjtmoJbDBT1b7nHvFCdU= +gopkg.in/check.v1 v0.0.0-20161208181325-20d25e280405/go.mod h1:Co6ibVJAznAaIkqp8huTwlJQCZ016jof/cbN4VW5Yz0= +gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= +gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA= +gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM= diff --git a/services/ingestion/internal/adapters/kafka.go b/services/ingestion/internal/adapters/kafka.go new file mode 100644 index 0000000..6c798aa --- /dev/null +++ b/services/ingestion/internal/adapters/kafka.go @@ -0,0 +1,150 @@ +package adapters + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "strings" + "time" + + "github.com/segmentio/kafka-go" + "github.com/triasha72/NewsLens/services/ingestion/internal/ingestion" +) + +type KafkaPublisher struct { + brokers []string + writer *kafka.Writer +} + +func NewKafkaPublisher(brokers []string, topic string) *KafkaPublisher { + return &KafkaPublisher{ + brokers: brokers, + writer: &kafka.Writer{ + Addr: kafka.TCP(brokers...), + Topic: topic, + Balancer: &kafka.Hash{}, + RequiredAcks: kafka.RequireAll, + Async: false, + BatchSize: 100, + BatchTimeout: 10 * time.Millisecond, + }, + } +} + +func (p *KafkaPublisher) Publish(ctx context.Context, event ingestion.ArticleEvent) error { + payload, err := json.Marshal(event) + if err != nil { + return fmt.Errorf("encode event: %w", err) + } + return p.writer.WriteMessages(ctx, kafka.Message{ + Key: []byte(event.ArticleID), + Value: payload, + Time: event.ProducedAt, + }) +} + +func (p *KafkaPublisher) Ping(ctx context.Context) error { + if len(p.brokers) == 0 { + return errors.New("no Kafka brokers configured") + } + connection, err := kafka.DialContext(ctx, "tcp", p.brokers[0]) + if err != nil { + return fmt.Errorf("connect to Kafka: %w", err) + } + return connection.Close() +} + +func (p *KafkaPublisher) Close() error { return p.writer.Close() } + +type KafkaSource struct { + reader *kafka.Reader +} + +func NewKafkaSource(brokers []string, topic, groupID, clientID string) *KafkaSource { + return &KafkaSource{reader: kafka.NewReader(kafka.ReaderConfig{ + Brokers: brokers, + Topic: topic, + GroupID: groupID, + Dialer: &kafka.Dialer{ClientID: clientID, Timeout: 10 * time.Second}, + MinBytes: 1, + MaxBytes: ingestion.MaxEventBytes, + CommitInterval: 1 * time.Second, + HeartbeatInterval: 2 * time.Second, + SessionTimeout: 6 * time.Second, + RebalanceTimeout: 10 * time.Second, + StartOffset: kafka.FirstOffset, + })} +} + +func (s *KafkaSource) Fetch(ctx context.Context) (ingestion.Message, error) { + message, err := s.reader.FetchMessage(ctx) + if err != nil { + return ingestion.Message{}, err + } + return ingestion.Message{ + Topic: message.Topic, + Partition: message.Partition, + Offset: message.Offset, + Key: message.Key, + Value: message.Value, + Time: message.Time, + }, nil +} + +func (s *KafkaSource) Commit(ctx context.Context, message ingestion.Message) error { + return s.reader.CommitMessages(ctx, kafka.Message{ + Topic: message.Topic, + Partition: message.Partition, + Offset: message.Offset, + Key: message.Key, + Value: message.Value, + Time: message.Time, + }) +} + +func (s *KafkaSource) Lag() int64 { return s.reader.Stats().Lag } +func (s *KafkaSource) Close() error { return s.reader.Close() } + +type KafkaDeadLetterPublisher struct { + writer *kafka.Writer +} + +func NewKafkaDeadLetterPublisher(brokers []string, topic string) *KafkaDeadLetterPublisher { + return &KafkaDeadLetterPublisher{writer: &kafka.Writer{ + Addr: kafka.TCP(brokers...), + Topic: topic, + Balancer: &kafka.Hash{}, + RequiredAcks: kafka.RequireAll, + Async: false, + }} +} + +func (p *KafkaDeadLetterPublisher) PublishDeadLetter( + ctx context.Context, + message ingestion.Message, + reason string, +) error { + envelope := struct { + FailedAt time.Time `json:"failed_at"` + Reason string `json:"reason"` + Source string `json:"source"` + Partition int `json:"partition"` + Offset int64 `json:"offset"` + OriginalPayloadBase64 []byte `json:"original_payload_base64"` + }{ + FailedAt: time.Now().UTC(), + Reason: strings.TrimSpace(reason), + Source: message.Topic, + Partition: message.Partition, + Offset: message.Offset, + OriginalPayloadBase64: message.Value, + } + payload, err := json.Marshal(envelope) + if err != nil { + return fmt.Errorf("encode dead-letter event: %w", err) + } + return p.writer.WriteMessages(ctx, kafka.Message{Key: message.Key, Value: payload}) +} + +func (p *KafkaDeadLetterPublisher) Close() error { return p.writer.Close() } diff --git a/services/ingestion/internal/adapters/postgres.go b/services/ingestion/internal/adapters/postgres.go new file mode 100644 index 0000000..5435c20 --- /dev/null +++ b/services/ingestion/internal/adapters/postgres.go @@ -0,0 +1,99 @@ +package adapters + +import ( + "context" + "fmt" + "time" + + "github.com/jackc/pgx/v5/pgxpool" + "github.com/triasha72/NewsLens/services/ingestion/internal/ingestion" +) + +type PostgresStore struct { + pool *pgxpool.Pool + now func() time.Time +} + +func NewPostgresStore(ctx context.Context, databaseURL string) (*PostgresStore, error) { + pool, err := pgxpool.New(ctx, databaseURL) + if err != nil { + return nil, fmt.Errorf("configure PostgreSQL pool: %w", err) + } + store := &PostgresStore{pool: pool, now: time.Now} + if err := store.Ping(ctx); err != nil { + pool.Close() + return nil, err + } + return store, nil +} + +func (s *PostgresStore) Ping(ctx context.Context) error { + if err := s.pool.Ping(ctx); err != nil { + return fmt.Errorf("connect to PostgreSQL: %w", err) + } + return nil +} + +func (s *PostgresStore) Persist( + ctx context.Context, + event ingestion.ArticleEvent, +) (ingestion.PersistResult, error) { + transaction, err := s.pool.Begin(ctx) + if err != nil { + return ingestion.PersistResult{}, fmt.Errorf("begin transaction: %w", err) + } + defer func() { _ = transaction.Rollback(ctx) }() + + indexedAt := s.now().UTC() + command, err := transaction.Exec( + ctx, + `INSERT INTO realtime_ingestion_events (event_id, article_id, produced_at, indexed_at) + VALUES ($1, $2, $3, $4) + ON CONFLICT (event_id) DO NOTHING`, + event.EventID, + event.ArticleID, + event.ProducedAt, + indexedAt, + ) + if err != nil { + return ingestion.PersistResult{}, fmt.Errorf("record ingestion event: %w", err) + } + + if command.RowsAffected() == 0 { + if err := transaction.Commit(ctx); err != nil { + return ingestion.PersistResult{}, fmt.Errorf("commit duplicate event: %w", err) + } + return ingestion.PersistResult{Duplicate: true, IndexedAt: indexedAt}, nil + } + + _, err = transaction.Exec( + ctx, + `INSERT INTO realtime_articles ( + article_id, title, body, category, published_at, produced_at, indexed_at, popularity + ) VALUES ($1, $2, $3, $4, $5, $6, $7, 0) + ON CONFLICT (article_id) DO UPDATE SET + title = EXCLUDED.title, + body = EXCLUDED.body, + category = EXCLUDED.category, + published_at = EXCLUDED.published_at, + produced_at = EXCLUDED.produced_at, + indexed_at = EXCLUDED.indexed_at + WHERE realtime_articles.produced_at <= EXCLUDED.produced_at`, + event.ArticleID, + event.Title, + event.Body, + event.Category, + event.PublishedAt, + event.ProducedAt, + indexedAt, + ) + if err != nil { + return ingestion.PersistResult{}, fmt.Errorf("upsert article: %w", err) + } + if err := transaction.Commit(ctx); err != nil { + return ingestion.PersistResult{}, fmt.Errorf("commit article: %w", err) + } + return ingestion.PersistResult{IndexedAt: indexedAt}, nil +} + +func (s *PostgresStore) Close() { s.pool.Close() } diff --git a/services/ingestion/internal/ingestion/contracts.go b/services/ingestion/internal/ingestion/contracts.go new file mode 100644 index 0000000..4621180 --- /dev/null +++ b/services/ingestion/internal/ingestion/contracts.go @@ -0,0 +1,44 @@ +package ingestion + +import ( + "context" + "time" +) + +type EventPublisher interface { + Publish(context.Context, ArticleEvent) error + Ping(context.Context) error + Close() error +} + +type Message struct { + Topic string + Partition int + Offset int64 + Key []byte + Value []byte + Time time.Time +} + +type EventSource interface { + Fetch(context.Context) (Message, error) + Commit(context.Context, Message) error + Lag() int64 + Close() error +} + +type DeadLetterPublisher interface { + PublishDeadLetter(context.Context, Message, string) error + Close() error +} + +type PersistResult struct { + Duplicate bool + IndexedAt time.Time +} + +type ArticleStore interface { + Persist(context.Context, ArticleEvent) (PersistResult, error) + Ping(context.Context) error + Close() +} diff --git a/services/ingestion/internal/ingestion/event.go b/services/ingestion/internal/ingestion/event.go new file mode 100644 index 0000000..b86be22 --- /dev/null +++ b/services/ingestion/internal/ingestion/event.go @@ -0,0 +1,103 @@ +package ingestion + +import ( + "bytes" + "encoding/json" + "errors" + "fmt" + "io" + "strings" + "time" +) + +const MaxEventBytes = 1 << 20 + +type ArticleEvent struct { + EventID string `json:"event_id"` + ArticleID string `json:"article_id"` + Title string `json:"title"` + Category string `json:"category"` + PublishedAt time.Time `json:"published_at"` + Body string `json:"body"` + ProducedAt time.Time `json:"produced_at"` +} + +func DecodeEvent(payload []byte) (ArticleEvent, error) { + if len(payload) == 0 { + return ArticleEvent{}, errors.New("event payload cannot be empty") + } + if len(payload) > MaxEventBytes { + return ArticleEvent{}, fmt.Errorf("event payload exceeds %d bytes", MaxEventBytes) + } + + decoder := json.NewDecoder(bytes.NewReader(payload)) + decoder.DisallowUnknownFields() + + var event ArticleEvent + if err := decoder.Decode(&event); err != nil { + return ArticleEvent{}, fmt.Errorf("decode event: %w", err) + } + + var extra any + if err := decoder.Decode(&extra); err != io.EOF { + return ArticleEvent{}, errors.New("event payload must contain one JSON object") + } + + if err := event.Validate(); err != nil { + return ArticleEvent{}, err + } + return event, nil +} + +func (e *ArticleEvent) Normalize(now time.Time, eventID func() (string, error)) error { + e.EventID = strings.TrimSpace(e.EventID) + e.ArticleID = strings.TrimSpace(e.ArticleID) + e.Title = strings.TrimSpace(e.Title) + e.Category = strings.ToLower(strings.TrimSpace(e.Category)) + e.Body = strings.TrimSpace(e.Body) + + if e.EventID == "" { + generated, err := eventID() + if err != nil { + return fmt.Errorf("generate event id: %w", err) + } + e.EventID = generated + } + if e.ProducedAt.IsZero() { + e.ProducedAt = now.UTC() + } + return e.Validate() +} + +func (e ArticleEvent) Validate() error { + checks := []struct { + name string + value string + max int + }{ + {"event_id", e.EventID, 128}, + {"article_id", e.ArticleID, 128}, + {"title", e.Title, 500}, + {"category", e.Category, 100}, + } + + for _, check := range checks { + if strings.TrimSpace(check.value) == "" { + return fmt.Errorf("%s is required", check.name) + } + if len(check.value) > check.max { + return fmt.Errorf("%s exceeds %d characters", check.name, check.max) + } + } + + if len(e.Body) > 100_000 { + return errors.New("body exceeds 100000 characters") + } + if e.PublishedAt.IsZero() { + return errors.New("published_at is required") + } + if e.ProducedAt.IsZero() { + return errors.New("produced_at is required") + } + return nil +} diff --git a/services/ingestion/internal/ingestion/event_test.go b/services/ingestion/internal/ingestion/event_test.go new file mode 100644 index 0000000..2ed992c --- /dev/null +++ b/services/ingestion/internal/ingestion/event_test.go @@ -0,0 +1,39 @@ +package ingestion + +import ( + "testing" + "time" +) + +func TestNormalizeSuppliesEventMetadata(t *testing.T) { + now := time.Date(2026, time.August, 23, 12, 0, 0, 0, time.UTC) + event := ArticleEvent{ + ArticleID: " article-1 ", + Title: " A useful title ", + Category: " Technology ", + PublishedAt: now.Add(-time.Hour), + } + + err := event.Normalize(now, func() (string, error) { return "event-1", nil }) + if err != nil { + t.Fatalf("Normalize() error = %v", err) + } + if event.EventID != "event-1" || event.ArticleID != "article-1" { + t.Fatalf("unexpected identifiers: %#v", event) + } + if event.Category != "technology" || !event.ProducedAt.Equal(now) { + t.Fatalf("unexpected normalized event: %#v", event) + } +} + +func TestDecodeEventRejectsUnknownAndTrailingFields(t *testing.T) { + unknown := []byte(`{"event_id":"e","article_id":"a","title":"t","category":"c","published_at":"2026-08-23T12:00:00Z","produced_at":"2026-08-23T12:00:00Z","extra":true}`) + if _, err := DecodeEvent(unknown); err == nil { + t.Fatal("DecodeEvent() accepted an unknown field") + } + + trailing := []byte(`{"event_id":"e","article_id":"a","title":"t","category":"c","published_at":"2026-08-23T12:00:00Z","produced_at":"2026-08-23T12:00:00Z"} {}`) + if _, err := DecodeEvent(trailing); err == nil { + t.Fatal("DecodeEvent() accepted a second JSON object") + } +} diff --git a/services/ingestion/internal/ingestion/http.go b/services/ingestion/internal/ingestion/http.go new file mode 100644 index 0000000..8ecd04c --- /dev/null +++ b/services/ingestion/internal/ingestion/http.go @@ -0,0 +1,138 @@ +package ingestion + +import ( + "bytes" + "context" + "crypto/rand" + "encoding/hex" + "encoding/json" + "errors" + "io" + "log/slog" + "net/http" + "time" +) + +type HTTPService struct { + publisher EventPublisher + store ArticleStore + metrics *Metrics + logger *slog.Logger + now func() time.Time + newID func() (string, error) +} + +func NewHTTPHandler(publisher EventPublisher, store ArticleStore, metrics *Metrics, logger *slog.Logger) http.Handler { + service := &HTTPService{ + publisher: publisher, + store: store, + metrics: metrics, + logger: logger, + now: time.Now, + newID: randomID, + } + + mux := http.NewServeMux() + mux.HandleFunc("GET /health", service.health) + mux.HandleFunc("GET /ready", service.ready) + mux.HandleFunc("GET /metrics", service.prometheus) + mux.HandleFunc("POST /events", service.publishEvent) + return mux +} + +func randomID() (string, error) { + value := make([]byte, 16) + if _, err := rand.Read(value); err != nil { + return "", err + } + return hex.EncodeToString(value), nil +} + +func writeJSON(w http.ResponseWriter, status int, value any) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(status) + _ = json.NewEncoder(w).Encode(value) +} + +func (s *HTTPService) health(w http.ResponseWriter, _ *http.Request) { + writeJSON(w, http.StatusOK, map[string]string{"status": "ok", "service": "newslens-ingestion"}) +} + +func ping(ctx context.Context, publisher EventPublisher, store ArticleStore) error { + if publisher != nil { + if err := publisher.Ping(ctx); err != nil { + return err + } + } + if store != nil { + if err := store.Ping(ctx); err != nil { + return err + } + } + return nil +} + +func (s *HTTPService) ready(w http.ResponseWriter, r *http.Request) { + ctx, cancel := context.WithTimeout(r.Context(), 2*time.Second) + defer cancel() + if err := ping(ctx, s.publisher, s.store); err != nil { + writeJSON(w, http.StatusServiceUnavailable, map[string]string{"status": "not_ready", "error": err.Error()}) + return + } + writeJSON(w, http.StatusOK, map[string]string{"status": "ready"}) +} + +func (s *HTTPService) prometheus(w http.ResponseWriter, _ *http.Request) { + w.Header().Set("Content-Type", "text/plain; version=0.0.4") + s.metrics.WritePrometheus(w) +} + +func (s *HTTPService) publishEvent(w http.ResponseWriter, r *http.Request) { + if s.publisher == nil { + writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "publisher is disabled in consumer mode"}) + return + } + + reader := http.MaxBytesReader(w, r.Body, MaxEventBytes) + payload, err := io.ReadAll(reader) + if err != nil { + writeJSON(w, http.StatusRequestEntityTooLarge, map[string]string{"error": "event payload is too large"}) + return + } + + var event ArticleEvent + jsonDecoder := json.NewDecoder(bytes.NewReader(payload)) + jsonDecoder.DisallowUnknownFields() + if err := jsonDecoder.Decode(&event); err != nil { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": err.Error()}) + return + } + var extra any + if err := jsonDecoder.Decode(&extra); err != io.EOF { + writeJSON(w, http.StatusBadRequest, map[string]string{"error": "event payload must contain one JSON object"}) + return + } + if err := event.Normalize(s.now(), s.newID); err != nil { + writeJSON(w, http.StatusUnprocessableEntity, map[string]string{"error": err.Error()}) + return + } + + ctx, cancel := context.WithTimeout(r.Context(), 5*time.Second) + defer cancel() + if err := s.publisher.Publish(ctx, event); err != nil { + if errors.Is(err, context.DeadlineExceeded) { + writeJSON(w, http.StatusGatewayTimeout, map[string]string{"error": "publishing timed out"}) + return + } + s.logger.Error("publish failed", "error", err, "event_id", event.EventID) + writeJSON(w, http.StatusServiceUnavailable, map[string]string{"error": "event could not be published"}) + return + } + + s.metrics.IncPublished() + writeJSON(w, http.StatusAccepted, map[string]string{ + "status": "accepted", + "event_id": event.EventID, + "article_id": event.ArticleID, + }) +} diff --git a/services/ingestion/internal/ingestion/http_test.go b/services/ingestion/internal/ingestion/http_test.go new file mode 100644 index 0000000..490b535 --- /dev/null +++ b/services/ingestion/internal/ingestion/http_test.go @@ -0,0 +1,66 @@ +package ingestion + +import ( + "context" + "encoding/json" + "errors" + "log/slog" + "net/http" + "net/http/httptest" + "strings" + "testing" + "time" +) + +type fakePublisher struct { + event ArticleEvent + err error +} + +func (p *fakePublisher) Publish(_ context.Context, event ArticleEvent) error { + p.event = event + return p.err +} +func (p *fakePublisher) Ping(context.Context) error { return p.err } +func (p *fakePublisher) Close() error { return nil } + +func TestPublishEventNormalizesAndAcceptsValidPayload(t *testing.T) { + publisher := &fakePublisher{} + metrics := &Metrics{} + service := &HTTPService{ + publisher: publisher, + metrics: metrics, + logger: slog.Default(), + now: func() time.Time { return time.Date(2026, 8, 23, 12, 0, 0, 0, time.UTC) }, + newID: func() (string, error) { return "generated-id", nil }, + } + payload := `{"article_id":"article-1","title":"Fresh chip release","category":"Technology","published_at":"2026-08-23T11:00:00Z","body":"details"}` + request := httptest.NewRequest(http.MethodPost, "/events", strings.NewReader(payload)) + response := httptest.NewRecorder() + + service.publishEvent(response, request) + + if response.Code != http.StatusAccepted { + t.Fatalf("status = %d, body = %s", response.Code, response.Body.String()) + } + if publisher.event.EventID != "generated-id" || publisher.event.Category != "technology" { + t.Fatalf("published event = %#v", publisher.event) + } + var body map[string]string + if err := json.Unmarshal(response.Body.Bytes(), &body); err != nil { + t.Fatal(err) + } + if body["event_id"] != "generated-id" { + t.Fatalf("body = %#v", body) + } +} + +func TestReadinessReportsDependencyFailure(t *testing.T) { + publisher := &fakePublisher{err: errors.New("broker unavailable")} + handler := NewHTTPHandler(publisher, nil, &Metrics{}, slog.Default()) + response := httptest.NewRecorder() + handler.ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/ready", nil)) + if response.Code != http.StatusServiceUnavailable { + t.Fatalf("status = %d", response.Code) + } +} diff --git a/services/ingestion/internal/ingestion/metrics.go b/services/ingestion/internal/ingestion/metrics.go new file mode 100644 index 0000000..29b5d98 --- /dev/null +++ b/services/ingestion/internal/ingestion/metrics.go @@ -0,0 +1,50 @@ +package ingestion + +import ( + "fmt" + "io" + "sync/atomic" +) + +type Metrics struct { + published atomic.Int64 + processed atomic.Int64 + duplicates atomic.Int64 + invalid atomic.Int64 + deadLettered atomic.Int64 + processingFails atomic.Int64 + consumerLag atomic.Int64 + freshnessMicros atomic.Int64 +} + +func (m *Metrics) IncPublished() { m.published.Add(1) } +func (m *Metrics) IncProcessed() { m.processed.Add(1) } +func (m *Metrics) IncDuplicates() { m.duplicates.Add(1) } +func (m *Metrics) IncInvalid() { m.invalid.Add(1) } +func (m *Metrics) IncDeadLettered() { m.deadLettered.Add(1) } +func (m *Metrics) IncProcessingFails() { m.processingFails.Add(1) } +func (m *Metrics) SetConsumerLag(value int64) { + if value >= 0 { + m.consumerLag.Store(value) + } +} +func (m *Metrics) SetFreshnessMicros(value int64) { + if value >= 0 { + m.freshnessMicros.Store(value) + } +} + +func metric(w io.Writer, name, help, kind string, value int64) { + fmt.Fprintf(w, "# HELP %s %s\n# TYPE %s %s\n%s %d\n", name, help, name, kind, name, value) +} + +func (m *Metrics) WritePrometheus(w io.Writer) { + metric(w, "newslens_ingestion_published_total", "Accepted article events.", "counter", m.published.Load()) + metric(w, "newslens_ingestion_processed_total", "Persisted article events.", "counter", m.processed.Load()) + metric(w, "newslens_ingestion_duplicates_total", "Idempotently ignored events.", "counter", m.duplicates.Load()) + metric(w, "newslens_ingestion_invalid_total", "Rejected or dead-lettered invalid events.", "counter", m.invalid.Load()) + metric(w, "newslens_ingestion_dead_letter_total", "Events written to the dead-letter topic.", "counter", m.deadLettered.Load()) + metric(w, "newslens_ingestion_failures_total", "Processing attempts that exhausted retries.", "counter", m.processingFails.Load()) + metric(w, "newslens_ingestion_consumer_lag", "Current Kafka consumer lag when available.", "gauge", m.consumerLag.Load()) + metric(w, "newslens_ingestion_index_freshness_microseconds", "Produced-to-indexed latency of the last persisted event.", "gauge", m.freshnessMicros.Load()) +} diff --git a/services/ingestion/internal/ingestion/worker.go b/services/ingestion/internal/ingestion/worker.go new file mode 100644 index 0000000..e3e5d83 --- /dev/null +++ b/services/ingestion/internal/ingestion/worker.go @@ -0,0 +1,108 @@ +package ingestion + +import ( + "context" + "log/slog" + "time" +) + +type Worker struct { + source EventSource + store ArticleStore + deadLetter DeadLetterPublisher + metrics *Metrics + logger *slog.Logger + retries int +} + +func NewWorker(source EventSource, store ArticleStore, deadLetter DeadLetterPublisher, metrics *Metrics, logger *slog.Logger, retries int) *Worker { + if retries < 1 { + retries = 1 + } + return &Worker{source: source, store: store, deadLetter: deadLetter, metrics: metrics, logger: logger, retries: retries} +} + +func (w *Worker) Run(ctx context.Context) error { + fetchFailures := 0 + for { + message, err := w.source.Fetch(ctx) + if err != nil { + if ctx.Err() != nil { + return nil + } + w.logger.Error("fetch failed", "error", err) + fetchFailures++ + delay := time.Duration(min(fetchFailures, 10)) * 100 * time.Millisecond + select { + case <-ctx.Done(): + return nil + case <-time.After(delay): + } + continue + } + fetchFailures = 0 + + commit, err := w.process(ctx, message) + if err != nil { + w.logger.Error("processing failed", "error", err) + continue + } + if commit { + if err := w.source.Commit(ctx, message); err != nil { + w.logger.Error("commit failed", "error", err) + continue + } + } + w.metrics.SetConsumerLag(w.source.Lag()) + } +} + +func (w *Worker) process(ctx context.Context, message Message) (bool, error) { + event, err := DecodeEvent(message.Value) + if err != nil { + w.metrics.IncInvalid() + if dlqErr := w.deadLetter.PublishDeadLetter(ctx, message, err.Error()); dlqErr != nil { + return false, dlqErr + } + w.metrics.IncDeadLettered() + return true, nil + } + + var result PersistResult + for attempt := 1; attempt <= w.retries; attempt++ { + result, err = w.store.Persist(ctx, event) + if err == nil { + break + } + if attempt < w.retries { + delay := time.Duration(attempt*100) * time.Millisecond + select { + case <-ctx.Done(): + return false, ctx.Err() + case <-time.After(delay): + } + } + } + + if err != nil { + w.metrics.IncProcessingFails() + if dlqErr := w.deadLetter.PublishDeadLetter(ctx, message, err.Error()); dlqErr != nil { + return false, dlqErr + } + w.metrics.IncDeadLettered() + return true, nil + } + + if result.Duplicate { + w.metrics.IncDuplicates() + return true, nil + } + + w.metrics.IncProcessed() + freshness := result.IndexedAt.Sub(event.ProducedAt) + if freshness < 0 { + freshness = 0 + } + w.metrics.SetFreshnessMicros(freshness.Microseconds()) + return true, nil +} diff --git a/services/ingestion/internal/ingestion/worker_test.go b/services/ingestion/internal/ingestion/worker_test.go new file mode 100644 index 0000000..96e253a --- /dev/null +++ b/services/ingestion/internal/ingestion/worker_test.go @@ -0,0 +1,80 @@ +package ingestion + +import ( + "bytes" + "context" + "errors" + "log/slog" + "testing" + "time" +) + +type fakeSource struct{ lag int64 } + +func (s *fakeSource) Fetch(context.Context) (Message, error) { return Message{}, nil } +func (s *fakeSource) Commit(context.Context, Message) error { return nil } +func (s *fakeSource) Lag() int64 { return s.lag } +func (s *fakeSource) Close() error { return nil } + +type fakeStore struct { + result PersistResult + err error + calls int +} + +func (s *fakeStore) Persist(context.Context, ArticleEvent) (PersistResult, error) { + s.calls++ + return s.result, s.err +} +func (s *fakeStore) Ping(context.Context) error { return s.err } +func (s *fakeStore) Close() {} + +type fakeDLQ struct { + calls int + reason string +} + +func (d *fakeDLQ) PublishDeadLetter(_ context.Context, _ Message, reason string) error { + d.calls++ + d.reason = reason + return nil +} +func (d *fakeDLQ) Close() error { return nil } + +func testEventPayload() []byte { + return []byte(`{"event_id":"event-1","article_id":"article-1","title":"News","category":"technology","published_at":"2026-08-23T11:00:00Z","produced_at":"2026-08-23T12:00:00Z"}`) +} + +func TestWorkerProcessesAndCountsDuplicateEvents(t *testing.T) { + indexedAt := time.Date(2026, 8, 23, 12, 0, 1, 0, time.UTC) + store := &fakeStore{result: PersistResult{Duplicate: true, IndexedAt: indexedAt}} + metrics := &Metrics{} + worker := NewWorker(&fakeSource{}, store, &fakeDLQ{}, metrics, slog.Default(), 3) + + commit, err := worker.process(context.Background(), Message{Value: testEventPayload()}) + if err != nil || !commit { + t.Fatalf("process() = (%v, %v)", commit, err) + } + var output bytes.Buffer + metrics.WritePrometheus(&output) + if !bytes.Contains(output.Bytes(), []byte("newslens_ingestion_duplicates_total 1")) { + t.Fatalf("metrics = %s", output.String()) + } +} + +func TestWorkerDeadLettersInvalidAndExhaustedEvents(t *testing.T) { + dlq := &fakeDLQ{} + metrics := &Metrics{} + worker := NewWorker(&fakeSource{}, &fakeStore{}, dlq, metrics, slog.Default(), 2) + commit, err := worker.process(context.Background(), Message{Value: []byte("not-json")}) + if err != nil || !commit || dlq.calls != 1 { + t.Fatalf("invalid process = (%v, %v), dlq calls = %d", commit, err, dlq.calls) + } + + store := &fakeStore{err: errors.New("database unavailable")} + worker = NewWorker(&fakeSource{}, store, dlq, metrics, slog.Default(), 2) + commit, err = worker.process(context.Background(), Message{Value: testEventPayload()}) + if err != nil || !commit || store.calls != 2 || dlq.calls != 2 { + t.Fatalf("failed process = (%v, %v), store calls = %d, dlq calls = %d", commit, err, store.calls, dlq.calls) + } +} diff --git a/src/newslens/api/app.py b/src/newslens/api/app.py index 06c11eb..b8f322f 100644 --- a/src/newslens/api/app.py +++ b/src/newslens/api/app.py @@ -11,6 +11,7 @@ from fastapi import ( FastAPI, HTTPException, + Query, Request, status, ) @@ -20,6 +21,12 @@ from newslens.models import ( ContentPopularityFallbackRecommender, ) +from newslens.realtime import ( + ArticleRepository, + FreshnessRanker, + PostgresArticleRepository, + understand_query, +) from .observability import ( LOGGER, @@ -29,10 +36,14 @@ from .schemas import ( HealthResponse, ModelInfoResponse, + QueryIntentResponse, ReadinessResponse, RecommendationItem, RecommendationRequest, RecommendationResponse, + SearchReadinessResponse, + SearchResponse, + SearchResultItem, ) from .settings import ApiSettings @@ -84,16 +95,26 @@ def _require_loaded_artifact( def create_app( *, artifact_path: str | Path | None = None, + realtime_repository: ArticleRepository | None = None, ) -> FastAPI: """Create an isolated NewsLens ASGI application.""" - configured_artifact_path = _resolve_artifact_path(artifact_path) + settings = ApiSettings.from_environment() + configured_artifact_path = ( + Path(artifact_path).expanduser() + if artifact_path is not None + else settings.artifact_path + ) + configured_realtime_repository = realtime_repository + candidate_limit = settings.realtime_candidate_limit @asynccontextmanager async def lifespan( application: FastAPI, ) -> AsyncIterator[None]: loaded_artifact: LoadedArtifact | None = None + live_repository = configured_realtime_repository + owns_live_repository = False if configured_artifact_path is not None: loaded_artifact = load_artifact(configured_artifact_path) @@ -106,12 +127,21 @@ async def lifespan( "The configured artifact does not contain a NewsLens fallback recommender." ) + if live_repository is None and settings.realtime_database_url is not None: + live_repository = PostgresArticleRepository(settings.realtime_database_url) + owns_live_repository = True + application.state.loaded_artifact = loaded_artifact + application.state.realtime_repository = live_repository try: yield finally: application.state.loaded_artifact = None + application.state.realtime_repository = None + + if owns_live_repository and live_repository is not None: + live_repository.close() application = FastAPI( title="NewsLens API", @@ -121,6 +151,7 @@ async def lifespan( ) install_request_observability(application) + ranker = FreshnessRanker() @application.get( "/health", @@ -157,6 +188,31 @@ def readiness( artifact_version=(loaded.metadata.artifact_version), ) + @application.get( + "/realtime/ready", + response_model=SearchReadinessResponse, + tags=["service"], + summary="Check streamed-article search readiness", + responses={ + status.HTTP_503_SERVICE_UNAVAILABLE: { + "description": "The real-time article store is not ready." + } + }, + ) + def realtime_readiness(request: Request) -> SearchReadinessResponse: + repository = cast( + ArticleRepository | None, + getattr(request.app.state, "realtime_repository", None), + ) + + if repository is None or not repository.ping(): + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail="Real-time article store is not ready.", + ) + + return SearchReadinessResponse(status="ready", realtime_store_ready=True) + @application.get( "/model-info", response_model=ModelInfoResponse, @@ -259,6 +315,77 @@ def recommend( recommendations=items, ) + @application.get( + "/search", + response_model=SearchResponse, + tags=["search"], + summary="Search newly ingested articles with freshness-aware ranking", + responses={ + status.HTTP_503_SERVICE_UNAVAILABLE: { + "description": "The real-time article store is not ready." + } + }, + ) + def search( + request: Request, + q: str = Query(min_length=1, max_length=500), + top_k: int = Query(default=10, ge=1, le=100), + ) -> SearchResponse: + repository = cast( + ArticleRepository | None, + getattr(request.app.state, "realtime_repository", None), + ) + + if repository is None: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail="Real-time article store is not ready.", + ) + + intent = understand_query(q) + search_started_at = perf_counter() + candidates = repository.list_candidates(intent, limit=candidate_limit) + ranked = ranker.rank(intent, candidates, top_k=top_k) + search_ms = (perf_counter() - search_started_at) * 1_000 + request_id = get_request_id(request) + + LOGGER.info( + "search_completed request_id=%s candidate_count=%d returned_count=%d " + "freshness_intent=%s search_ms=%.3f", + request_id, + len(candidates), + len(ranked), + intent.prefers_freshness, + search_ms, + ) + + return SearchResponse( + request_id=request_id, + intent=QueryIntentResponse( + normalized_query=intent.normalized_query, + category=intent.category, + entity=intent.entity, + prefers_freshness=intent.prefers_freshness, + ), + candidate_count=len(candidates), + returned_count=len(ranked), + search_ms=search_ms, + results=tuple( + SearchResultItem( + article_id=item.document.article_id, + title=item.document.title, + category=item.document.category, + published_at=item.document.published_at.isoformat(), + score=item.score, + relevance_score=item.relevance_score, + freshness_score=item.freshness_score, + popularity_score=item.popularity_score, + index_freshness_ms=item.index_freshness_ms, + ) + for item in ranked + ), + ) + return application diff --git a/src/newslens/api/schemas.py b/src/newslens/api/schemas.py index 0acbd26..e898750 100644 --- a/src/newslens/api/schemas.py +++ b/src/newslens/api/schemas.py @@ -119,3 +119,52 @@ class RecommendationResponse(BaseModel): returned_count: int inference_ms: float = Field(ge=0.0) recommendations: tuple[RecommendationItem, ...] + + +class SearchReadinessResponse(BaseModel): + """Readiness information for the streamed-article search path.""" + + model_config = ConfigDict(frozen=True) + + status: Literal["ready"] + realtime_store_ready: Literal[True] + + +class QueryIntentResponse(BaseModel): + """Parsed query information returned for inspectability.""" + + model_config = ConfigDict(frozen=True) + + normalized_query: str + category: str | None + entity: str | None + prefers_freshness: bool + + +class SearchResultItem(BaseModel): + """One freshness-aware search result.""" + + model_config = ConfigDict(frozen=True) + + article_id: str + title: str + category: str + published_at: str + score: float + relevance_score: float + freshness_score: float + popularity_score: float + index_freshness_ms: float = Field(ge=0.0) + + +class SearchResponse(BaseModel): + """Real-time search response with query and ranking diagnostics.""" + + model_config = ConfigDict(frozen=True) + + request_id: str + intent: QueryIntentResponse + candidate_count: int + returned_count: int + search_ms: float = Field(ge=0.0) + results: tuple[SearchResultItem, ...] diff --git a/src/newslens/api/settings.py b/src/newslens/api/settings.py index 1c4d4e5..537ba1b 100644 --- a/src/newslens/api/settings.py +++ b/src/newslens/api/settings.py @@ -7,6 +7,8 @@ from pathlib import Path ARTIFACT_PATH_ENVIRONMENT_VARIABLE = "NEWSLENS_ARTIFACT_PATH" +REALTIME_DATABASE_URL_ENVIRONMENT_VARIABLE = "NEWSLENS_REALTIME_DATABASE_URL" +REALTIME_CANDIDATE_LIMIT_ENVIRONMENT_VARIABLE = "NEWSLENS_REALTIME_CANDIDATE_LIMIT" class ApiSettingsError(ValueError): @@ -18,6 +20,8 @@ class ApiSettings: """Configuration required to start the NewsLens API.""" artifact_path: Path | None = None + realtime_database_url: str | None = None + realtime_candidate_limit: int = 200 @classmethod def from_environment(cls) -> ApiSettings: @@ -25,12 +29,46 @@ def from_environment(cls) -> ApiSettings: raw_artifact_path = os.getenv(ARTIFACT_PATH_ENVIRONMENT_VARIABLE) - if raw_artifact_path is None: - return cls() + raw_database_url = os.getenv(REALTIME_DATABASE_URL_ENVIRONMENT_VARIABLE) + raw_candidate_limit = os.getenv(REALTIME_CANDIDATE_LIMIT_ENVIRONMENT_VARIABLE) - normalized_path = raw_artifact_path.strip() + artifact_path: Path | None = None - if not normalized_path: - raise ApiSettingsError(f"{ARTIFACT_PATH_ENVIRONMENT_VARIABLE} cannot be empty.") + if raw_artifact_path is not None: + normalized_path = raw_artifact_path.strip() - return cls(artifact_path=Path(normalized_path).expanduser()) + if not normalized_path: + raise ApiSettingsError(f"{ARTIFACT_PATH_ENVIRONMENT_VARIABLE} cannot be empty.") + + artifact_path = Path(normalized_path).expanduser() + + database_url: str | None = None + + if raw_database_url is not None: + database_url = raw_database_url.strip() + + if not database_url: + raise ApiSettingsError( + f"{REALTIME_DATABASE_URL_ENVIRONMENT_VARIABLE} cannot be empty." + ) + + candidate_limit = 200 + + if raw_candidate_limit is not None: + try: + candidate_limit = int(raw_candidate_limit) + except ValueError as error: + raise ApiSettingsError( + f"{REALTIME_CANDIDATE_LIMIT_ENVIRONMENT_VARIABLE} must be an integer." + ) from error + + if candidate_limit <= 0 or candidate_limit > 10_000: + raise ApiSettingsError( + f"{REALTIME_CANDIDATE_LIMIT_ENVIRONMENT_VARIABLE} must be between 1 and 10000." + ) + + return cls( + artifact_path=artifact_path, + realtime_database_url=database_url, + realtime_candidate_limit=candidate_limit, + ) diff --git a/src/newslens/realtime/__init__.py b/src/newslens/realtime/__init__.py new file mode 100644 index 0000000..49b4c27 --- /dev/null +++ b/src/newslens/realtime/__init__.py @@ -0,0 +1,20 @@ +"""Real-time query understanding, storage, and freshness-aware ranking.""" + +from .query import QueryIntent, understand_query +from .ranking import FreshnessRanker, RankedArticle, SearchDocument +from .repository import ( + ArticleRepository, + InMemoryArticleRepository, + PostgresArticleRepository, +) + +__all__ = [ + "ArticleRepository", + "FreshnessRanker", + "InMemoryArticleRepository", + "PostgresArticleRepository", + "QueryIntent", + "RankedArticle", + "SearchDocument", + "understand_query", +] diff --git a/src/newslens/realtime/query.py b/src/newslens/realtime/query.py new file mode 100644 index 0000000..349660c --- /dev/null +++ b/src/newslens/realtime/query.py @@ -0,0 +1,161 @@ +"""Deterministic query understanding for the real-time search path.""" + +from __future__ import annotations + +import re +from dataclasses import dataclass + +TOKEN_PATTERN = re.compile(r"[A-Za-z0-9]+(?:'[A-Za-z0-9]+)?") + +FRESHNESS_TERMS = frozenset( + { + "breaking", + "current", + "latest", + "live", + "new", + "recent", + "today", + "tonight", + } +) + +CATEGORY_TERMS: dict[str, frozenset[str]] = { + "business": frozenset( + { + "business", + "earnings", + "finance", + "market", + "markets", + "shares", + "stock", + "stocks", + } + ), + "sports": frozenset( + { + "baseball", + "basketball", + "football", + "game", + "match", + "score", + "soccer", + "sports", + } + ), + "technology": frozenset( + { + "ai", + "app", + "chip", + "software", + "tech", + "technology", + } + ), + "science": frozenset( + { + "climate", + "mars", + "nasa", + "research", + "science", + "space", + } + ), + "health": frozenset( + { + "health", + "hospital", + "medicine", + "vaccine", + } + ), +} + +ENTITY_STOPWORDS = frozenset( + { + "a", + "an", + "and", + "for", + "in", + "of", + "on", + "the", + "to", + } +) | FRESHNESS_TERMS | frozenset().union(*CATEGORY_TERMS.values()) + +ENTITY_ALLOWLIST = frozenset({"mars", "nasa"}) + + +@dataclass(frozen=True, slots=True) +class QueryIntent: + """Structured query information used by retrieval and ranking.""" + + normalized_query: str + tokens: tuple[str, ...] + category: str | None + entity: str | None + prefers_freshness: bool + + +def _detect_category(tokens: tuple[str, ...]) -> str | None: + token_set = set(tokens) + matches = [ + ( + category, + len(token_set & terms), + min(index for index, token in enumerate(tokens) if token in terms), + ) + for category, terms in CATEGORY_TERMS.items() + if token_set & terms + ] + + if not matches: + return None + + return max(matches, key=lambda item: (item[1], -item[2], item[0]))[0] + + +def _detect_entity(query: str) -> str | None: + candidates: list[str] = [] + + for token in TOKEN_PATTERN.findall(query): + normalized = token.lower() + + if normalized in ENTITY_STOPWORDS and normalized not in ENTITY_ALLOWLIST: + continue + + if token.isupper() or token[:1].isupper(): + candidates.append(token) + + if not candidates: + return None + + return " ".join(candidates[:3]) + + +def understand_query(query: str) -> QueryIntent: + """Normalize a query and extract category, entity, and freshness intent.""" + + normalized_query = " ".join(query.strip().split()) + + if not normalized_query: + raise ValueError("query cannot be empty") + + tokens = tuple(token.lower() for token in TOKEN_PATTERN.findall(normalized_query)) + + if not tokens: + raise ValueError("query must contain at least one searchable token") + + return QueryIntent( + normalized_query=normalized_query, + tokens=tokens, + category=_detect_category(tokens), + entity=_detect_entity(normalized_query), + prefers_freshness=bool(set(tokens) & FRESHNESS_TERMS), + ) diff --git a/src/newslens/realtime/ranking.py b/src/newslens/realtime/ranking.py new file mode 100644 index 0000000..47526b8 --- /dev/null +++ b/src/newslens/realtime/ranking.py @@ -0,0 +1,169 @@ +"""Transparent lexical and freshness-aware ranking for streamed articles.""" + +from __future__ import annotations + +import math +import re +from dataclasses import dataclass +from datetime import UTC, datetime + +from .query import FRESHNESS_TERMS, QueryIntent + +TOKEN_PATTERN = re.compile(r"[A-Za-z0-9]+(?:'[A-Za-z0-9]+)?") + + +@dataclass(frozen=True, slots=True) +class SearchDocument: + """One article available to the real-time search ranker.""" + + article_id: str + title: str + body: str + category: str + published_at: datetime + produced_at: datetime + indexed_at: datetime + popularity: int = 0 + + +@dataclass(frozen=True, slots=True) +class RankedArticle: + """Ranked article with inspectable score components.""" + + document: SearchDocument + score: float + relevance_score: float + freshness_score: float + popularity_score: float + + @property + def index_freshness_ms(self) -> float: + """Return event-produced to indexed latency in milliseconds.""" + + return max( + 0.0, + (self.document.indexed_at - self.document.produced_at).total_seconds() + * 1_000, + ) + + +class FreshnessRanker: + """Rank documents with lexical relevance and an explicit freshness boost.""" + + def __init__( + self, + *, + relevance_weight: float = 0.75, + freshness_weight: float = 0.20, + popularity_weight: float = 0.05, + freshness_half_life_hours: float = 24.0, + ) -> None: + weights = (relevance_weight, freshness_weight, popularity_weight) + + if any(weight < 0.0 for weight in weights): + raise ValueError("ranking weights cannot be negative") + + if not math.isclose(sum(weights), 1.0, abs_tol=1e-9): + raise ValueError("ranking weights must sum to 1.0") + + if freshness_half_life_hours <= 0.0: + raise ValueError("freshness_half_life_hours must be positive") + + self.relevance_weight = relevance_weight + self.freshness_weight = freshness_weight + self.popularity_weight = popularity_weight + self.freshness_half_life_hours = freshness_half_life_hours + + @staticmethod + def _relevance(intent: QueryIntent, document: SearchDocument) -> float: + query_tokens = tuple( + token for token in intent.tokens if token not in FRESHNESS_TERMS + ) + + if not query_tokens: + query_tokens = intent.tokens + + title_tokens = { + token.lower() for token in TOKEN_PATTERN.findall(document.title) + } + body_tokens = {token.lower() for token in TOKEN_PATTERN.findall(document.body)} + document_tokens = title_tokens | body_tokens | {document.category.lower()} + + matched_weight = sum( + 2.0 if token in title_tokens else 1.0 + for token in query_tokens + if token in document_tokens + ) + maximum_weight = 2.0 * len(query_tokens) + + if maximum_weight == 0.0: + return 0.0 + + score = matched_weight / maximum_weight + + if intent.category and document.category.lower() == intent.category: + score = min(1.0, score + 0.15) + + if intent.entity and intent.entity.lower() in document.title.lower(): + score = min(1.0, score + 0.20) + + return score + + def _freshness(self, document: SearchDocument, *, now: datetime) -> float: + age_hours = max(0.0, (now - document.published_at).total_seconds() / 3_600) + return math.exp(-math.log(2.0) * age_hours / self.freshness_half_life_hours) + + def rank( + self, + intent: QueryIntent, + documents: tuple[SearchDocument, ...], + *, + top_k: int, + now: datetime | None = None, + ) -> tuple[RankedArticle, ...]: + """Rank candidates and return score components for inspection.""" + + if top_k <= 0: + raise ValueError("top_k must be positive") + + ranking_time = now or datetime.now(UTC) + max_popularity = max((document.popularity for document in documents), default=0) + scored: list[RankedArticle] = [] + + for document in documents: + relevance = self._relevance(intent, document) + freshness = self._freshness(document, now=ranking_time) + popularity = ( + math.log1p(max(0, document.popularity)) / math.log1p(max_popularity) + if max_popularity > 0 + else 0.0 + ) + active_freshness_weight = ( + self.freshness_weight if intent.prefers_freshness else 0.0 + ) + active_relevance_weight = self.relevance_weight + ( + self.freshness_weight - active_freshness_weight + ) + score = ( + active_relevance_weight * relevance + + active_freshness_weight * freshness + + self.popularity_weight * popularity + ) + scored.append( + RankedArticle( + document=document, + score=score, + relevance_score=relevance, + freshness_score=freshness, + popularity_score=popularity, + ) + ) + + scored.sort( + key=lambda item: ( + -item.score, + -item.document.published_at.timestamp(), + item.document.article_id, + ) + ) + return tuple(scored[:top_k]) diff --git a/src/newslens/realtime/repository.py b/src/newslens/realtime/repository.py new file mode 100644 index 0000000..64af531 --- /dev/null +++ b/src/newslens/realtime/repository.py @@ -0,0 +1,153 @@ +"""Repositories for real-time article search.""" + +from __future__ import annotations + +from collections.abc import Iterable +from typing import Protocol, runtime_checkable + +from .query import QueryIntent +from .ranking import SearchDocument + + +@runtime_checkable +class ArticleRepository(Protocol): + """Storage contract used by the real-time search API.""" + + def ping(self) -> bool: + """Return whether the backing store is reachable.""" + + def list_candidates( + self, + intent: QueryIntent, + *, + limit: int, + ) -> tuple[SearchDocument, ...]: + """Return a bounded candidate set for ranking.""" + + def close(self) -> None: + """Release repository resources.""" + + +class InMemoryArticleRepository: + """Deterministic repository for tests and local ranking experiments.""" + + def __init__(self, documents: Iterable[SearchDocument]) -> None: + self._documents = tuple(documents) + + def ping(self) -> bool: + return True + + def list_candidates( + self, + intent: QueryIntent, + *, + limit: int, + ) -> tuple[SearchDocument, ...]: + if limit <= 0: + raise ValueError("limit must be positive") + + candidates = self._documents + + if intent.category is not None: + category_matches = tuple( + document + for document in candidates + if document.category.lower() == intent.category + ) + if category_matches: + candidates = category_matches + + return tuple( + sorted( + candidates, + key=lambda document: ( + -document.published_at.timestamp(), + document.article_id, + ), + )[:limit] + ) + + def close(self) -> None: + return None + + +class PostgresArticleRepository: + """Read streamed articles from the PostgreSQL ingestion store.""" + + def __init__(self, database_url: str) -> None: + if not database_url.strip(): + raise ValueError("database_url cannot be empty") + + import psycopg + + self._connection = psycopg.connect(database_url, autocommit=True) + self._database_error = psycopg.Error + + def ping(self) -> bool: + try: + with self._connection.cursor() as cursor: + cursor.execute("SELECT 1") + return cursor.fetchone() == (1,) + except self._database_error: + return False + + def list_candidates( + self, + intent: QueryIntent, + *, + limit: int, + ) -> tuple[SearchDocument, ...]: + if limit <= 0: + raise ValueError("limit must be positive") + + entity_pattern = f"%{intent.entity}%" if intent.entity else None + + with self._connection.cursor() as cursor: + cursor.execute( + """ + SELECT + article_id, + title, + body, + category, + published_at, + produced_at, + indexed_at, + popularity + FROM realtime_articles + WHERE (CAST(%s AS TEXT) IS NULL OR lower(category) = %s) + AND ( + CAST(%s AS TEXT) IS NULL + OR title ILIKE %s + OR body ILIKE %s + ) + ORDER BY published_at DESC, article_id ASC + LIMIT %s + """, + ( + intent.category, + intent.category, + entity_pattern, + entity_pattern, + entity_pattern, + limit, + ), + ) + rows = cursor.fetchall() + + return tuple( + SearchDocument( + article_id=row[0], + title=row[1], + body=row[2], + category=row[3], + published_at=row[4], + produced_at=row[5], + indexed_at=row[6], + popularity=row[7], + ) + for row in rows + ) + + def close(self) -> None: + self._connection.close() diff --git a/tests/test_api.py b/tests/test_api.py index cf20401..e8934bf 100644 --- a/tests/test_api.py +++ b/tests/test_api.py @@ -2,6 +2,7 @@ from __future__ import annotations +from datetime import UTC, datetime, timedelta from pathlib import Path import pandas as pd @@ -14,6 +15,7 @@ ArtifactNotFoundError, export_fallback_artifact, ) +from newslens.realtime import InMemoryArticleRepository, SearchDocument pytestmark = pytest.mark.filterwarnings( "ignore:Setting the shape on a NumPy array has been deprecated:DeprecationWarning" @@ -171,11 +173,75 @@ def test_openapi_schema_lists_service_endpoints() -> None: assert set(response.json()["paths"]) == { "/health", "/ready", + "/realtime/ready", "/model-info", "/recommend", + "/search", } +def test_realtime_readiness_is_honest_without_store() -> None: + with TestClient(create_app()) as client: + response = client.get("/realtime/ready") + + assert response.status_code == 503 + assert response.json() == {"detail": "Real-time article store is not ready."} + + +def test_search_requires_realtime_store() -> None: + with TestClient(create_app()) as client: + response = client.get("/search", params={"q": "latest Apple earnings"}) + + assert response.status_code == 503 + + +def test_search_returns_intent_and_freshness_diagnostics() -> None: + now = datetime.now(UTC) + repository = InMemoryArticleRepository( + ( + SearchDocument( + article_id="N-new", + title="Apple reports latest earnings", + body="Quarterly market results.", + category="business", + published_at=now - timedelta(hours=1), + produced_at=now - timedelta(seconds=1), + indexed_at=now - timedelta(milliseconds=800), + popularity=10, + ), + SearchDocument( + article_id="N-old", + title="Apple reports earnings", + body="An older quarterly market report.", + category="business", + published_at=now - timedelta(days=5), + produced_at=now - timedelta(days=5, seconds=1), + indexed_at=now - timedelta(days=5), + popularity=100, + ), + ) + ) + + with TestClient(create_app(realtime_repository=repository)) as client: + readiness = client.get("/realtime/ready") + response = client.get( + "/search", + params={"q": "latest Apple earnings", "top_k": 2}, + ) + + assert readiness.status_code == 200 + assert response.status_code == 200 + body = response.json() + assert body["intent"] == { + "normalized_query": "latest Apple earnings", + "category": "business", + "entity": "Apple", + "prefers_freshness": True, + } + assert body["results"][0]["article_id"] == "N-new" + assert body["results"][0]["index_freshness_ms"] == pytest.approx(200.0) + + def test_unknown_route_returns_not_found() -> None: with TestClient(create_app()) as client: response = client.get("/missing") diff --git a/tests/test_realtime_query.py b/tests/test_realtime_query.py new file mode 100644 index 0000000..eab7f67 --- /dev/null +++ b/tests/test_realtime_query.py @@ -0,0 +1,29 @@ +"""Tests for deterministic real-time query understanding.""" + +from __future__ import annotations + +import pytest + +from newslens.realtime import understand_query + + +def test_understand_query_extracts_freshness_category_and_entity() -> None: + intent = understand_query("latest Apple earnings") + + assert intent.normalized_query == "latest Apple earnings" + assert intent.category == "business" + assert intent.entity == "Apple" + assert intent.prefers_freshness is True + + +def test_understand_query_handles_plain_topic_search() -> None: + intent = understand_query("mars rover mission") + + assert intent.category == "science" + assert intent.entity is None + assert intent.prefers_freshness is False + + +def test_understand_query_rejects_empty_text() -> None: + with pytest.raises(ValueError, match="empty"): + understand_query(" ") diff --git a/tests/test_realtime_ranking.py b/tests/test_realtime_ranking.py new file mode 100644 index 0000000..ee6cca9 --- /dev/null +++ b/tests/test_realtime_ranking.py @@ -0,0 +1,70 @@ +"""Tests for transparent freshness-aware ranking.""" + +from __future__ import annotations + +from datetime import UTC, datetime, timedelta + +import pytest + +from newslens.realtime import FreshnessRanker, SearchDocument, understand_query + + +def article( + article_id: str, + *, + title: str, + age_hours: int, + popularity: int = 0, +) -> SearchDocument: + now = datetime(2026, 8, 23, 12, tzinfo=UTC) + produced_at = now - timedelta(hours=age_hours, milliseconds=120) + return SearchDocument( + article_id=article_id, + title=title, + body="Apple quarterly earnings and market coverage.", + category="business", + published_at=now - timedelta(hours=age_hours), + produced_at=produced_at, + indexed_at=produced_at + timedelta(milliseconds=120), + popularity=popularity, + ) + + +def test_freshness_intent_prefers_newer_equally_relevant_article() -> None: + now = datetime(2026, 8, 23, 12, tzinfo=UTC) + ranker = FreshnessRanker() + ranked = ranker.rank( + understand_query("latest Apple earnings"), + ( + article("old", title="Apple earnings report", age_hours=72), + article("new", title="Apple earnings report", age_hours=1), + ), + top_k=2, + now=now, + ) + + assert [item.document.article_id for item in ranked] == ["new", "old"] + assert ranked[0].freshness_score > ranked[1].freshness_score + + +def test_score_components_and_index_freshness_are_exposed() -> None: + now = datetime(2026, 8, 23, 12, tzinfo=UTC) + result = FreshnessRanker().rank( + understand_query("Apple earnings"), + (article("N1", title="Apple earnings report", age_hours=2, popularity=10),), + top_k=1, + now=now, + )[0] + + assert 0.0 <= result.score <= 1.0 + assert result.relevance_score > 0.0 + assert result.index_freshness_ms == pytest.approx(120.0) + + +def test_ranker_rejects_invalid_weights() -> None: + with pytest.raises(ValueError, match="sum"): + FreshnessRanker( + relevance_weight=0.5, + freshness_weight=0.5, + popularity_weight=0.5, + )