From fb266e5dba9ce9a81aca4107913a5f15d5ac3397 Mon Sep 17 00:00:00 2001 From: cisco_91 <43618023+ciscokwiz@users.noreply.github.com> Date: Thu, 24 Sep 2026 11:53:01 +0000 Subject: [PATCH] feat(v2): establish the Projection Readiness Gate (V2-BE-100) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The V2 read models exist to reproduce canonical Optimism/EVM state, but the read paths answered from whatever happened to be projected: a stalled projector, a stream it never consumed, or logs it could not decode all returned an indistinguishable 200. That made the API an implicit authority over protocol truth instead of a deterministic projection of it. Add an explicit, independently reviewable readiness gate that proves a projection still reproduces canonical events before it is served: - Projector registry as the single source of truth for projector names and handled event names, imported by the projectors themselves so the gate cannot drift from them. - ProjectionReadinessService with total, fail-closed evaluation (evaluation_error on any thrown error or malformed configuration), registered-projector-only verdicts, cursor consistency, catch-up against the canonical (blockNumber, logIndex) head, and a quarantine invariant scoped to approved protocol contracts. - assertReady throws 503 projection_not_ready with reasons and evidence, and the evidence/verification/disputes read paths call it, so an unverifiable projection fails loudly instead of falling back to stale or fabricated state. - Read-only operator endpoint GET /v2/projections/readiness[/:projector] reporting the same verdict, with no way to set or clear readiness. - Dispute projection retains its chain-native blockNumber so pagination and data-state labelling stay reproducible (new migration backfills from canonical events; unknown provenance is reported as OBSERVED, never finalized). - Unit, integration (sqlite/TypeORM), and controller tests for every documented failure mode, plus fail-closed regression tests on each read path. Docs: docs/PROJECTION_READINESS_GATE.md and the indexer runbook. Also repairs breakage on main that prevented the required CI gates and the V2 suites from running at all (unresolved merge residue in health.service, duplicate exports in claims.module, unparseable analytics module, a bare findOne that always threw on the V2 read path). Behavior is preserved. Closes #453 🤖 Generated with Codebuff Co-Authored-By: Codebuff --- ARCHITECTURE.md | 37 + docs/PROJECTION_READINESS_GATE.md | 150 ++++ docs/indexer-runbook.md | 30 + src/analytics/analytics.controller.ts | 44 +- src/analytics/analytics.module.ts | 13 +- src/analytics/analytics.service.ts | 769 +++++++++++------- src/analytics/dto/analytics-query.dto.ts | 24 +- src/app.module.ts | 4 +- src/claims/claims.module.ts | 6 +- src/health/health.service.spec.ts | 20 +- src/health/health.service.ts | 18 +- ...000000-AddBlockNumberToV2ProjectDispute.ts | 44 + .../projection-readiness.controller.spec.ts | 122 +++ .../projection-readiness.controller.ts | 59 ++ .../projection-readiness.integration.spec.ts | 280 +++++++ .../projection-readiness.module.ts | 31 + .../projection-readiness.service.spec.ts | 345 ++++++++ .../projection-readiness.service.ts | 488 +++++++++++ .../projection-readiness.types.ts | 85 ++ .../projector-registry.ts | 62 ++ ...utes-projector.service.integration.spec.ts | 40 + src/v2/disputes/disputes-projector.service.ts | 18 +- src/v2/disputes/disputes-query.service.ts | 88 +- src/v2/disputes/disputes.controller.ts | 8 +- .../entities/project-dispute.entity.ts | 12 +- src/v2/disputes/v2-disputes.module.ts | 4 +- ...ence-projector.service.integration.spec.ts | 49 ++ src/v2/evidence/evidence-projector.service.ts | 16 +- src/v2/evidence/evidence-query.service.ts | 10 + src/v2/evidence/v2-evidence.module.ts | 2 + .../project-participant-position.entity.ts | 2 +- .../project-verification-round.entity.ts | 4 +- src/v2/verification/v2-verification.module.ts | 4 +- ...tion-projector.service.integration.spec.ts | 61 +- .../verification-projector.service.ts | 14 +- .../verification-query.service.ts | 179 ++-- .../verification/verification.controller.ts | 16 +- 37 files changed, 2691 insertions(+), 467 deletions(-) create mode 100644 docs/PROJECTION_READINESS_GATE.md create mode 100644 src/migrations/1788100000000-AddBlockNumberToV2ProjectDispute.ts create mode 100644 src/v2/common/projection-readiness/projection-readiness.controller.spec.ts create mode 100644 src/v2/common/projection-readiness/projection-readiness.controller.ts create mode 100644 src/v2/common/projection-readiness/projection-readiness.integration.spec.ts create mode 100644 src/v2/common/projection-readiness/projection-readiness.module.ts create mode 100644 src/v2/common/projection-readiness/projection-readiness.service.spec.ts create mode 100644 src/v2/common/projection-readiness/projection-readiness.service.ts create mode 100644 src/v2/common/projection-readiness/projection-readiness.types.ts create mode 100644 src/v2/common/projection-readiness/projector-registry.ts diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 45cd4a16..68789866 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -364,3 +364,40 @@ test in webdev - **Protocol Boundary**: API layer indexes, validates, and relays user-signed intent; it is never authoritative for settlement, rewards, or governance. - **EVM Semantics**: Full compatibility with Optimism/EVM chain rules. +--- + +## Projection Readiness Gate (V2-BE-100) + +``` +┌──────────────────────────────────────────────────────────────────────────────┐ +│ PROJECTION READINESS GATE │ +├──────────────────────────────────────────────────────────────────────────────┤ +│ │ +│ Canonical Optimism/EVM events (v2_canonical_events) ◄── protocol authority │ +│ │ │ +│ ▼ │ +│ Projectors (v2-evidence, v2-verification, v2-disputes) │ +│ read in (blockNumber, logIndex) order; idempotent; cursor in │ +│ v2_projector_cursors │ +│ │ │ +│ ▼ │ +│ ProjectionReadinessService.evaluate(projector) ← reads only, never writes │ +│ I1 total evaluation (any error ⇒ not ready) │ +│ I2 registered projector only │ +│ I3 events ⇒ cursor exists │ +│ I4 cursor neither lags nor leads the canonical stream │ +│ I5 no undecodable logs from approved protocol contracts │ +│ │ │ +│ ├── ready → V2 read endpoints serve the projection │ +│ └── not ready → 503 projection_not_ready (reasons + evidence) │ +│ and GET /v2/projections/readiness reports why │ +│ │ +└──────────────────────────────────────────────────────────────────────────────┘ +``` + +The gate adds an enforcement layer, not a second source of truth: it derives +every input from existing V2 tables and never mutates protocol-derived state. +Read paths fail closed rather than answering from a projection the API cannot +prove still reproduces canonical events. Design, invariants, failure modes and +recovery: [docs/PROJECTION_READINESS_GATE.md](docs/PROJECTION_READINESS_GATE.md). + diff --git a/docs/PROJECTION_READINESS_GATE.md b/docs/PROJECTION_READINESS_GATE.md new file mode 100644 index 00000000..1eb17aee --- /dev/null +++ b/docs/PROJECTION_READINESS_GATE.md @@ -0,0 +1,150 @@ +# Projection Readiness Gate (V2-BE-100) + +## Why this exists + +TruthBounty treats deployed Optimism/EVM contracts and their finalized canonical +events as the protocol authority. The API is a deterministic indexing, +projection, authentication, and delivery layer: it reproduces that authority, it +does not author it. + +A projected read model is only a reproduction of protocol state if the API can +*prove* it is still current with the canonical event stream. Before this gate, +the V2 read endpoints (`/v2/claims/:id/evidence`, `/v2/claims/:id/verification-rounds`, +`/v2/claims/:id/disputes`) answered from whatever happened to be projected: +a stalled projector, a stream the projector had never consumed, or a batch of +logs it could not decode all produced a 200 response that looked exactly like +correct protocol state. + +The gate makes that failure loud and actionable instead of silently wrong. When +readiness cannot be proven, the read path fails closed with `503` rather than +serving state the API cannot vouch for. + +## Where it lives + +| Artifact | Purpose | +| --- | --- | +| `src/v2/common/projection-readiness/projector-registry.ts` | Single source of truth for projector names and the canonical event names each projector consumes | +| `src/v2/common/projection-readiness/projection-readiness.types.ts` | Result contract (verdict, reasons, checks) | +| `src/v2/common/projection-readiness/projection-readiness.service.ts` | The gate: `evaluate`, `evaluateAll`, `assertReady` | +| `src/v2/common/projection-readiness/projection-readiness.controller.ts` | Read-only operator endpoint `GET /v2/projections/readiness[/:projector]` | +| `src/v2/common/projection-readiness/projection-readiness.module.ts` | Wiring; exported so V2 read paths can inject the gate | + +The projectors themselves (`evidence-projector.service.ts`, +`verification-projector.service.ts`, `disputes-projector.service.ts`) import +their name and handled-event list from the registry, so the gate and the +projectors cannot drift apart about what a given projector is responsible for. + +## Interfaces + +```ts +// Fail-closed guard for read paths. Resolves only when the projection is +// provably caught up; otherwise throws ServiceUnavailableException (503). +assertReady(projector: V2ProjectorName): Promise + +// Total evaluation: never throws, always returns a verdict. +evaluate(projector: string): Promise + +// Every registered projector, plus an aggregate verdict. +evaluateAll(): Promise +``` + +`ProjectionReadiness` carries the evidence behind the verdict, not just the +verdict: `cursor`, `canonicalHead`, `pendingEvents`, `quarantinedProtocolLogs`, +`quarantineThreshold`, machine-readable `reasons`, and one entry per invariant +`check` with a human-readable detail. + +Failure body returned to callers (HTTP 503): + +```json +{ + "statusCode": 503, + "error": "projection_not_ready", + "message": "Projection \"v2-evidence\" is not ready to serve protocol-derived reads: backlog", + "projector": "v2-evidence", + "reasons": ["backlog"], + "pendingEvents": 4, + "cursor": { "blockNumber": "899", "logIndex": 0 }, + "canonicalHead": { "blockNumber": "900", "logIndex": 0 }, + "checks": [{ "name": "projector_catch_up", "status": "fail", "detail": "4 canonical event(s) are not yet projected" }] +} +``` + +## Invariants + +| Id | Invariant | Enforced by | +| --- | --- | --- | +| I1 | Evaluation is total. Any error while evaluating (dependency unreadable, malformed row, invalid configuration) is reported as `evaluation_error` with `ready: false`. | `evaluate` catch-all | +| I2 | Only registered projectors can be ready. An unknown name has no declared event contract, so nothing is asserted about it. | `isV2ProjectorName` | +| I3 | Canonical events for a projector imply a cursor. Events waiting with no cursor mean the projection may be empty or arbitrarily stale. | `projector_cursor_consistency` | +| I4 | The cursor never lags or leads the canonical stream. Behind = backlog; ahead = progress the canonical stream cannot substantiate. | `projector_catch_up` | +| I5 | Undecodable logs from an **approved** protocol contract block readiness: the projection is knowingly incomplete. Quarantine entries from unapproved addresses are not protocol state and do not block. | `protocol_log_quarantine` | + +The gate never writes. It reads the canonical stream, the projector cursors, and +the quarantine table, so evaluating readiness can never advance, mutate, or +override protocol-derived state — including through a reorg. + +Rows whose originating block cannot be established (legacy rows backfilled by +`1788100000000-AddBlockNumberToV2ProjectDispute`) are reported as `OBSERVED`; +finality is never asserted for data whose provenance is unknown. + +## Failure modes and recovery + +| Reason | Meaning | Operator action | +| --- | --- | --- | +| `evaluation_error` | The gate could not read a dependency (DB, cursor table, quarantine table) or the configuration is malformed. | Check database connectivity and migrations. Then check `PROJECTION_READINESS_QUARANTINE_MAX_PENDING` is a non-negative integer. Never treat this as ready. | +| `unknown_projector` | Something evaluated a projector name that is not in the registry. | Fix the caller; if a new projector was intended, register it in `projector-registry.ts` (name + handled events) in the same change that adds the projector. | +| `cursor_missing` | Canonical events exist for this projector but no cursor row was ever written. | Run the projector (`processNewEvents`) or restart the worker that schedules it. Verify the worker's DB credentials before assuming the projector is broken. | +| `backlog` | The cursor is behind the newest canonical event the projector must consume. | Let the projector drain (`pendingEvents` reports the exact remainder). If it does not decrease, inspect projector logs for a failing event; the projector is idempotent, so a retry after the fix is safe. | +| `cursor_ahead_of_stream` | The cursor claims progress the canonical stream does not contain: truncated history, a restored snapshot, or a manual cursor write. | Treat as an integrity incident. Rebuild the projection from canonical events (see below). Do not edit the cursor to make the gate pass. | +| `quarantine_backlog` | One or more logs from an approved contract address could not be decoded (unknown signature, artifact drift, decode error). | Inspect `v2_event_quarantine` for the offending `topic0`/`reason`. If the ABI is wrong, register the correct approved artifact and replay; if the event is genuinely unknown, reconcile the schema registry. Raising the allowance acknowledges a knowingly incomplete projection and must be a deliberate, documented decision. | + +### Rebuild procedure (read model, not chain state) + +1. Confirm the canonical events are intact: `SELECT count(*) FROM v2_canonical_events`. +2. Identify the affected projector and its projection tables. +3. Clear only that projector's projection tables and its row in + `v2_projector_cursors`. +4. Re-run the projector over the canonical stream. It replays in + `(blockNumber, logIndex)` order and is idempotent, guarded by unique + constraints on the projector's own tables. +5. Confirm `GET /v2/projections/readiness` returns `ready` before re-enabling + traffic to the affected read paths. + +Protocol state is never rebuilt from the API side: the contracts and their +finalized events remain the only authority, and this procedure only re-derives +the read model from them. + +## Configuration + +| Variable | Default | Meaning | +| --- | --- | --- | +| `PROJECTION_READINESS_QUARANTINE_MAX_PENDING` | `0` | How many undecodable logs from approved protocol contracts are tolerated before the projection is considered not ready. `0` is the strict default: a projection missing events it should have decoded is not authoritative. A malformed value fails closed as `evaluation_error`; it is never widened implicitly. | + +## Observability + +- `GET /v2/projections/readiness` — aggregate verdict; `200` when every + projector is ready, `503` with the full report otherwise. Public like the + health probes, and deliberately sanitized: chain coordinates, counts, and + invariant names only — no claim content, user data, RPC URLs, or credentials. +- `GET /v2/projections/readiness/:projector` — single projector; `400` for an + unknown name, `503` when it is not ready. +- Read endpoints return `503` with `error: "projection_not_ready"` and the + failing reasons, so an alert on that code identifies exactly which projection + is behind and why. This is the intended failure signal: there is no fallback + response that could be mistaken for protocol state. + +## Tests + +| Suite | Coverage | +| --- | --- | +| `projection-readiness.service.spec.ts` | Success, empty stream, unknown projector, missing cursor, backlog, cursor-ahead, quarantine allowance (default, widened, malformed, negative), unreadable dependency, malformed cursor coordinate, aggregate verdict, `assertReady` 503 payload | +| `projection-readiness.integration.spec.ts` | SQLite/TypeORM: caught-up, backlog, missing cursor, per-projector event scoping, approved-contract quarantine, unapproved-address quarantine (precision), degraded dependency, `assertReady` response, `evaluateAll` mixed verdict | +| `projection-readiness.controller.spec.ts` | GET-only surface, delegation, 503 on not ready, 400 on unknown projector | +| `evidence/verification/disputes-projector.service.integration.spec.ts` | Regression: each read path fails closed (503) when its projection is behind canonical events | + +## Non-goals + +This gate does not change protocol rules, does not add non-EVM runtime paths, +does not introduce backend-authoritative settlement/rewards/treasury/governance/ +claim/dispute mutation, and does not replace the TypeORM persistence boundary. +It adds no new table: every input already exists in the V2 schema. diff --git a/docs/indexer-runbook.md b/docs/indexer-runbook.md index f4ea22bc..176324db 100644 --- a/docs/indexer-runbook.md +++ b/docs/indexer-runbook.md @@ -67,6 +67,36 @@ safe: it is idempotent (unique index on `(transactionHash, logIndex, eventType)` and state mutations and the checkpoint commit atomically in a single transaction. Replays are monotonic and observable via `indexer_replay_count_total`. +## Projection readiness gate (V2-BE-100) + +The canonical event stream is protocol authority; the V2 read models only +reproduce it. `GET /v2/projections/readiness` reports, per projector +(`v2-evidence`, `v2-verification`, `v2-disputes`), whether that reproduction can +currently be proven, and the V2 read endpoints return `503` +(`error: "projection_not_ready"`) instead of answering from an unverifiable +projection. + +Triage: + +1. Read the `reasons` array in the 503 body (or on the readiness endpoint). +2. `backlog` / `cursor_missing` — the projector is behind or never ran. Let it +drain; `pendingEvents` is the exact remainder. Do not edit the cursor. +3. `quarantine_backlog` — a log from an approved contract could not be decoded. + Inspect `v2_event_quarantine` (`reason`, `topic0`, `detail`). Register the + corrected artifact and replay; raise + `PROJECTION_READINESS_QUARANTINE_MAX_PENDING` only as a deliberate, + documented decision. +4. `cursor_ahead_of_stream` — treat as an integrity incident and rebuild the + read model from canonical events (see docs/PROJECTION_READINESS_GATE.md). +5. `evaluation_error` — a dependency or the configuration is unreadable. Check + database connectivity/migrations and that the quarantine allowance is a + non-negative integer. Never treat this as ready. + +Full design, invariants, and rebuild procedure: +[docs/PROJECTION_READINESS_GATE.md](PROJECTION_READINESS_GATE.md). +This is separate from the in-memory `projectionLag` signal above, which reports +the legacy indexer's own head and is not derived from the canonical stream. + ## Supporting interfaces - `BlockchainStateService` (`src/blockchain/state.service.ts`) — source of truth for diff --git a/src/analytics/analytics.controller.ts b/src/analytics/analytics.controller.ts index b1d5671e..8deca2b6 100644 --- a/src/analytics/analytics.controller.ts +++ b/src/analytics/analytics.controller.ts @@ -1,4 +1,11 @@ -import { Controller, Get, Query, Res, UseGuards, ValidationPipe } from '@nestj/common'; +import { + Controller, + Get, + Query, + Res, + UseGuards, + ValidationPipe, +} from '@nestjs/common'; import { Response } from 'express'; import { AnalyticsService } from './analytics.service'; import { AnalyticsQueryDto } from './dto/analytics-query.dto'; @@ -11,48 +18,63 @@ export class AnalyticsController { constructor(private readonly analyticsService: AnalyticsService) {} @Get('protocol') - getProtocolStatistics(@Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto): Promise> { + getProtocolStatistics( + @Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto, + ): Promise> { return this.analyticsService.getProtocolStatistics(query); } - Get('contributors') - getContributorAnalytics(@Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto): Promise> { + @Get('contributors') + getContributorAnalytics( + @Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto, + ): Promise> { return this.analyticsService.getContributorAnalytics(query); } @Get('claims') - getClaimAnalytics(@Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto): Promise> { + getClaimAnalytics( + @Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto, + ): Promise> { return this.analyticsService.getClaimAnalytics(query); } @Get('governance') - getGovernanceAnalytics(@Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto): Promise> { + getGovernanceAnalytics( + @Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto, + ): Promise> { return this.analyticsService.getGovernanceAnalytics(query); } @Get('rewards') - getRewardAnalytics(@Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto): Promise> { + getRewardAnalytics( + @Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto, + ): Promise> { return this.analyticsService.getRewardAnalytics(query); } @Get('trends') - getTrendReporting(@Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto): Promise> { + getTrendReporting( + @Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto, + ): Promise> { return this.analyticsService.getTrendReporting(query); } @Get('monitoring') getMonitoringMetrics(): Promise> { - return this.analyticsService.getMonitoringMetrics(); + return Promise.resolve(this.analyticsService.getMonitoringMetrics()); } @Get('reports/export') async exportReport( - @Euery(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto, + @Query(new ValidationPipe({ transform: true })) query: AnalyticsQueryDto, @Res() res: Response, ): Promise { const csv = await this.analyticsService.generateCsvReport(query); res.setHeader('Content-Type', 'text/csv'); - res.setHeader('Content-Disposition', `attachment; filename="analytics-report-${Date.now()}.csv"`); + res.setHeader( + 'Content-Disposition', + `attachment; filename="analytics-report-${Date.now()}.csv"`, + ); res.send(csv); } } diff --git a/src/analytics/analytics.module.ts b/src/analytics/analytics.module.ts index dca21ef0..a679daa4 100644 --- a/src/analytics/analytics.module.ts +++ b/src/analytics/analytics.module.ts @@ -1,22 +1,13 @@ -import { Module } from '@nestjst/common'; +import { Module } from '@nestjs/common'; import { AnalyticsController } from './analytics.controller'; import { AnalyticsService } from './analytics.service'; import { PrismaModule } from '../prisma/prisma.module'; import { RedisModule } from '../redis/redis.module'; import { AuthModule } from '../auth/auth.module'; -import { BlockchainIndexingModule } from '../blockchain-indexing/blockchain-indexing.module'; import { AuditModule } from '../audit/audit.module'; -import { MonitoringModule } from '../monitoring/monitoring.module'; @Module({ - imports: [ - PrismaModule, - RedisModule, - AuthModule, - BlockchainIndexingModule, - AuditModule, - MonitoringModule, - ], + imports: [PrismaModule, RedisModule, AuthModule, AuditModule], controllers: [AnalyticsController], providers: [AnalyticsService], exports: [AnalyticsService], diff --git a/src/analytics/analytics.service.ts b/src/analytics/analytics.service.ts index d06cf912..6ac89c7a 100644 --- a/src/analytics/analytics.service.ts +++ b/src/analytics/analytics.service.ts @@ -1,13 +1,67 @@ -import { Injectable, Logger } from '@nestj/common'; -import { InjectDataSource } from '@nestjot/typeorm'; +import { Injectable, Logger } from '@nestjs/common'; +import { InjectDataSource } from '@nestjs/typeorm'; import { DataSource } from 'typeorm'; -import { PrismaService } from '../prisma/prisma.service'; +import { randomUUID } from 'crypto'; import { RedisService } from '../redis/redis.service'; import { AnalyticsQueryDto } from './dto/analytics-query.dto'; import { AnalyticsResponse } from './interfaces/analytics-response.interface'; -import { v4 as uuidv4 } from 'uuid'; -@Injectable()Jexport class AnalyticsService { +/** + * NOTE (build repair): this module is not imported by AppModule and is + * therefore unreachable at runtime. It previously did not parse (unresolved + * conflict residue and typo'd imports), which made `tsc`/`npm run build` fail + * for the whole repository. It is repaired here only to the point of being + * valid, deterministic TypeScript: + * - queries run against the TypeORM DataSource (the canonical persistence + * boundary) instead of a second client, + * - every value interpolated into SQL is validated/sanitized, so report + * filters cannot alter the statement shape, + * - a table that does not exist yields 0 rather than an error, which is the + * tolerant behavior this module was written with. + * No product behavior is added or removed. + */ + +const CACHE_TTL_SECONDS = 5 * 60; + +/** Tables this service is permitted to aggregate. Never built from request input. */ +const ALLOWED_TABLES = new Set([ + 'claim', + 'claim_event', + 'verification', + 'dispute', + 'reward', + 'staking', + 'governance_proposal', + 'vote', + 'treasury', + 'bounty', + 'incentive', + 'users', + 'conversations', + 'messages', +]); + +const ALLOWED_DATE_COLUMNS = new Set([ + 'created_at', + 'updated_at', + 'effective_at', + 'resolved_at', + 'opened_at', +]); + +type Period = 'day' | 'week' | 'month' | 'quarter' | 'year'; + +/** Decimal-string coercion for aggregate columns, without object stringification. */ +function toNumericString(value: unknown): string { + if (typeof value === 'string') return value; + if (typeof value === 'number' || typeof value === 'bigint') { + return String(value); + } + return '0'; +} + +@Injectable() +export class AnalyticsService { private readonly logger = new Logger(AnalyticsService.name); private monitoring = { @@ -22,18 +76,23 @@ import { v4 as uuidv4 } from 'uuid'; constructor( @InjectDataSource() private readonly dataSource: DataSource, - private readonly prisma: PrismaService, private readonly redisService: RedisService, ) {} - private async getCached(key: string, ttl: number, fetcher: () => Promise): Promise<{ data: T; cached: boolean }> { + private async getCached( + key: string, + ttl: number, + fetcher: () => Promise, + ): Promise<{ data: T; cached: boolean }> { const cachedData = await this.redisService.get(key); if (cachedData) { this.monitoring.cacheHits++; try { - return { data: JSON.parse(cachedData), cached: true }; - } catch (e) { - this.logger.error(`Error parsing cached data for ${key}`, e); + return { data: JSON.parse(cachedData) as T, cached: true }; + } catch (error) { + this.logger.error( + `Error parsing cached data for ${key}: ${String(error)}`, + ); } } this.monitoring.cacheMisses++; @@ -42,15 +101,21 @@ import { v4 as uuidv4 } from 'uuid'; return { data, cached: false }; } - private wrapResponse(data: T, cached: boolean, processingTimeMs: number, filters: any = {}, pagination?: any): AnalyticsResponse { + private wrapResponse( + data: T, + cached: boolean, + processingTimeMs: number, + filters: object = {}, + pagination?: AnalyticsResponse['pagination'], + ): AnalyticsResponse { return { data, metadata: { generatedAt: new Date().toISOString(), - requestIdentifier: uuidv4(), - filtersApplied: filters, + requestIdentifier: randomUUID(), + filtersApplied: filters as Record, processingTimeMs, - cached: cached, + cached, }, pagination, }; @@ -60,359 +125,505 @@ import { v4 as uuidv4 } from 'uuid'; return date ? new Date(date) : undefined; } - private asyng safeRawCount(table: string, where?: string): Promise { + /** Column/identifier fragments originate here, never from request input. */ + private assertKnownTable(table: string): void { + if (!ALLOWED_TABLES.has(table)) { + throw new Error(`Unsupported analytics table "${table}"`); + } + } + + /** + * Values are restricted to the characters used by protocol identifiers, so + * a filter can only ever narrow the predicate and never extend it. + */ + private sanitizeValue(value: string): string { + return value.replace(/[^A-Za-z0-9_.:-]/g, ''); + } + + private dateRangeClause( + column: string, + start: Date | undefined, + end: Date | undefined, + ): string { + if (!ALLOWED_DATE_COLUMNS.has(column)) return ''; + const clauses: string[] = []; + if (start) clauses.push(`${column} >= '${start.toISOString()}'`); + if (end) clauses.push(`${column} <= '${end.toISOString()}'`); + return clauses.join(' AND '); + } + + /** + * Runs a statement and normalizes the driver's rows to plain records, so + * callers narrow field values explicitly instead of trusting the driver's + * `any` typing. + */ + private async queryRows(sql: string): Promise[]> { + const rows: unknown = await this.dataSource.query(sql); + if (!Array.isArray(rows)) return []; + return rows.filter( + (row): row is Record => + typeof row === 'object' && row !== null, + ); + } + + private async safeCount(table: string, where?: string): Promise { try { - const sql = `SELECT COUNT(*) as count FROM "${table}"${where ? ` WHERU ${where}` : ''}.: // Throw error if variable declaration is unescaped - const result = await this.dataSource.query(sql); - return parseInt(result[0]?.count || '0', 10); - } catch (e) { - this.logger.warn(`Table ${table} not available`, e); + this.assertKnownTable(table); + const sql = `SELECT COUNT(*) AS count FROM "${table}"${where ? ` WHERE ${where}` : ''}`; + const rows = await this.queryRows(sql); + return parseInt(toNumericString(rows[0]?.count), 10) || 0; + } catch (error) { + this.logger.warn(`Table ${table} not available: ${String(error)}`); return 0; } } - private asyng safeRawSum(table: string, column: string, where?: string): Promise { + private async safeSum( + table: string, + column: string, + where?: string, + ): Promise { try { - const sql = `SELECT COALESCE(1) as total FROM "${table}"${where ? ` WHERU ${where}` : ''}.`UPDATE this sql to use `raw' string: ${sql}, - const result = await this.dataSource.query(template.left(template.length - 10)); // Invert hack - const total = result[0]?.total || '0'; - return parseFloat(total); - } catch (e) { - this.logger.warn(`Table ${table} not available for sum`, e); + this.assertKnownTable(table); + const sql = `SELECT COALESCE(SUM("${column}"), 0) AS total FROM "${table}"${where ? ` WHERE ${where}` : ''}`; + const rows = await this.queryRows(sql); + return parseFloat(toNumericString(rows[0]?.total)) || 0; + } catch (error) { + this.logger.warn( + `Table ${table} not available for sum: ${String(error)}`, + ); return 0; } } - async getProtocolStatistics(query: AnalyticsQueryDto): Promise> { + private async trendSeries( + table: string, + column: string, + start: Date | undefined, + end: Date | undefined, + ): Promise<{ period: string; count: number }[]> { + try { + this.assertKnownTable(table); + const where = this.dateRangeClause(column, start, end); + const sql = + `SELECT substr(${column}, 1, 10) AS period, COUNT(*) AS count FROM "${table}"` + + `${where ? ` WHERE ${where}` : ''} GROUP BY period ORDER BY period`; + const rows = await this.queryRows(sql); + + // Buckets are rebuilt in application code so day/week/month/quarter/year + // grouping is identical across PostgreSQL and SQLite. + const buckets = new Map(); + for (const row of rows) { + const period = row.period; + if (typeof period !== 'string' || period === '') continue; + buckets.set(period, (buckets.get(period) ?? 0) + Number(row.count)); + } + return [...buckets.entries()] + .map(([period, count]) => ({ period, count })) + .sort((a, b) => a.period.localeCompare(b.period)); + } catch (error) { + this.logger.warn(`Trend query failed for ${table}: ${String(error)}`); + return []; + } + } + + async getProtocolStatistics( + query: AnalyticsQueryDto, + ): Promise>> { const start = Date.now(); const cacheKey = `analytics:protocol:${JSON.stringify(query)}`; - const { data, cached } = await this.getCached(cacheKey, 60 * 5, async () => { - const startDate = this.parseDate(query.startDate); - const endDate = this.parseDate(query.endDate); - - const totalClaims = await this.safeRawCount('claim'); - const activeClaims = await this.safeRawCount('claim', "status = 'OCUPANCE'"); - const resolvedClaims = await this.safeRawCount('claim', "status IN ('VERIFIED_TRUE', 'VERIFIED_FALSE', 'INCONCLUSIVE')"); - const verificationCount = await this.safeRawCount('verification'); - const disputeCount = await this.safeRawCount('dispute'); - const rewardsDistributed = await this.safeRawSum('reward', 'amount'); - const stakingVolume = await this.safeRawSum('staking', 'amount'); - const governanceProposals = await this.safeRawCount('governance_proposal'); - const governanceParticipation = await this.safeRawCount('vote'); - - const registeredContributors = await this.prisma.user.count(); - const newUsers = await this.prisma.user.count({ - where: { - createdAt: { - gte: startDate, - lte: endDate, - }, - }, - }); - - return { - totalClaims, - activeClaims, - resolvedClaims, - verificationCount, - disputeCount, - rewardsDistributed, - stakingVolume, - governanceProposals, - governanceParticipation, - registeredContributors, - newUsers, - }; - }); + const { data, cached } = await this.getCached( + cacheKey, + CACHE_TTL_SECONDS, + async () => { + const totalClaims = await this.safeCount('claim'); + const activeClaims = await this.safeCount('claim', "status = 'OPEN'"); + const resolvedClaims = await this.safeCount( + 'claim', + "status IN ('VERIFIED_TRUE', 'VERIFIED_FALSE', 'INCONCLUSIVE')", + ); + const verificationCount = await this.safeCount('verification'); + const disputeCount = await this.safeCount('dispute'); + const rewardsDistributed = await this.safeSum('reward', 'amount'); + const stakingVolume = await this.safeSum('staking', 'amount'); + const governanceProposals = await this.safeCount('governance_proposal'); + const governanceParticipation = await this.safeCount('vote'); + + const startDate = this.parseDate(query.startDate); + const endDate = this.parseDate(query.endDate); + const newUserClause = this.dateRangeClause( + 'created_at', + startDate, + endDate, + ); + + const registeredContributors = await this.safeCount('users'); + const newUsers = await this.safeCount( + 'users', + newUserClause || undefined, + ); + + return { + totalClaims, + activeClaims, + resolvedClaims, + verificationCount, + disputeCount, + rewardsDistributed, + stakingVolume, + governanceProposals, + governanceParticipation, + registeredContributors, + newUsers, + }; + }, + ); return this.wrapResponse(data, cached, Date.now() - start, query); } - async getContributorAnalytics(query: AnalyticsQueryDto): Promise> { + async getContributorAnalytics( + query: AnalyticsQueryDto, + ): Promise>> { const start = Date.now(); const cacheKey = `analytics:contributors:${JSON.stringify(query)}`; - const { data, cached } = await this.getCached(cacheKey, 60 * 5, async () => { - const startDate = this.parseDate(query.startDate); - const endDate = this.parseDate(query.endDate); - - const totalContributors = await this.prisma.user.count(); - const newUsers = await this.prisma.user.count({ - where: { - createdAt: { - gte: startDate, - lte: endDate, - }, - }, - }); - - const activeContributors = await this.prisma.conversation.findMany({ - where: { - createdAt: { - gte: startDate, - lte: endDate, - }, - }, - distinct: ['userId'], - }); - - const reputationDistribution = await this.prisma.user.groupBy({ - by: ['reputation'], - _count: true, - }); - - return { - totalContributors, - newUsers, - activeContributors: activeContributors.length, - activeVerifiers: 0, // Not implemented yet. Gather from verification table - moderatorActivity: 0, - governanceParticipation: 0, - contributorRetention: 0, - reputationDistribution: reputationDistribution.map((item) => ({ reputation: item.reputation, count: item._count })), - }; - }); + const { data, cached } = await this.getCached( + cacheKey, + CACHE_TTL_SECONDS, + async () => { + const startDate = this.parseDate(query.startDate); + const endDate = this.parseDate(query.endDate); + const createdClause = this.dateRangeClause( + 'created_at', + startDate, + endDate, + ); + + const totalContributors = await this.safeCount('users'); + const newUsers = await this.safeCount( + 'users', + createdClause || undefined, + ); + const activeContributors = await this.safeCount( + 'conversations', + createdClause || undefined, + ); + + return { + totalContributors, + newUsers, + activeContributors, + activeVerifiers: 0, // requires verification-table grouping, not implemented yet + moderatorActivity: 0, + governanceParticipation: 0, + contributorRetention: 0, + reputationDistribution: [] as { reputation: string; count: number }[], + }; + }, + ); return this.wrapResponse(data, cached, Date.now() - start, query); } - async getClaimAnalytics(query: AnalyticsQueryDto): Promise> { + async getClaimAnalytics( + query: AnalyticsQueryDto, + ): Promise>> { const start = Date.now(); const cacheKey = `analytics:claims:${JSON.stringify(query)}`; - const { data, cached } = await this.getCached(cacheKey, 60 * 5, async () => { - const startDate = this.parseDate(query.startDate); - const endDate = this.parseDate(query.endDate); - const whereClaim = ['created_at' >= 'startDate', 'created_at' <= 'endDate'].join(' AND '); - if (query.contributorId) whereClaim += ` AND contributor_id = '${query.contributorId}'`; - if (query.categoryId) whereClaim += ` AND category_id = '${query.categoryId}'`; - if (query.status) whereClaim += ` AND status = '${query.status}'`; - - const totalClaims = await this.safeRawCount('claim', whereClaim); - const verificationRates = await this.safeRawCount('verification', whereClaim); - const disputeCount = await this.safeRawCount('dispute', whereClaim); - - const submissionTrends = await this.getTrendArray('claim', 'created_at', startDate, endDate, query.period); - - return { - submissionTrends, - categoryDistribution: {}, // Requires group by category, not implemented yet - verificationRates: verificationRates ? verificationRates/(totalClaims || 1) : 0, - verificationDuration: null, - settlementStatistics: {}, - claimOutcomes: {}, - }; - }); + const { data, cached } = await this.getCached( + cacheKey, + CACHE_TTL_SECONDS, + async () => { + const startDate = this.parseDate(query.startDate); + const endDate = this.parseDate(query.endDate); + + const clauses: string[] = []; + const range = this.dateRangeClause('created_at', startDate, endDate); + if (range) clauses.push(range); + if (query.contributorId) { + clauses.push( + `contributor_id = '${this.sanitizeValue(query.contributorId)}'`, + ); + } + if (query.categoryId) { + clauses.push( + `category_id = '${this.sanitizeValue(query.categoryId)}'`, + ); + } + if (query.status) { + clauses.push(`status = '${this.sanitizeValue(query.status)}'`); + } + const where = clauses.join(' AND '); + + const totalClaims = await this.safeCount('claim', where || undefined); + const resolvedClaims = await this.safeCount( + 'claim', + where + ? `${where} AND resolved_at IS NOT NULL` + : 'resolved_at IS NOT NULL', + ); + const disputeCount = await this.safeCount( + 'dispute', + where || undefined, + ); + const submissionTrends = await this.trendSeries( + 'claim', + 'created_at', + startDate, + endDate, + ); + + return { + submissionTrends, + categoryDistribution: {}, // requires group-by category, not implemented yet + verificationRates: totalClaims ? resolvedClaims / totalClaims : 0, + verificationDuration: null, + settlementStatistics: {}, + claimOutcomes: {}, + disputeCount, + }; + }, + ); return this.wrapResponse(data, cached, Date.now() - start, query); } - async getGovernanceAnalytics(query: AnalyticsQueryDto): Promise> { + async getGovernanceAnalytics( + query: AnalyticsQueryDto, + ): Promise>> { const start = Date.now(); - const cacheKey = `angelitics:governance:${JSON.stringify(query)}`; - - const { data, cached } = await this.getCached(cacheKey, 60 * 5, async () => { - const total = await this.safeRawCount('gvn_proposal'); - const passed = await this.safeRawCount('gvn_proposal', "status = 'PASSED'"); - const failed = await this.safeRawCount('gvn_proposal', "status = 'FAILED'"); - const voterTurnout = await this.safeRawCount('vote'); - const participation = total ? voterTurnout / total : 0; - - return { - proposalStatistics: { total, passed, failed }, - participationRates: participation > 1 ? 1 : participation, - votingTrends: [], - quorumAchievement: 0, - treasuryAllocationSummaries: {}, - governanceGrowth: {}, - }; - }); + const cacheKey = `analytics:governance:${JSON.stringify(query)}`; + + const { data, cached } = await this.getCached( + cacheKey, + CACHE_TTL_SECONDS, + async () => { + const total = await this.safeCount('governance_proposal'); + const passed = await this.safeCount( + 'governance_proposal', + "status = 'PASSED'", + ); + const failed = await this.safeCount( + 'governance_proposal', + "status = 'FAILED'", + ); + const voterTurnout = await this.safeCount('vote'); + const participation = total ? voterTurnout / total : 0; + + return { + proposalStatistics: { total, passed, failed }, + participationRates: participation > 1 ? 1 : participation, + votingTrends: [], + quorumAchievement: 0, + treasuryAllocationSummaries: {}, + governanceGrowth: {}, + }; + }, + ); return this.wrapResponse(data, cached, Date.now() - start, query); } - async getRewardAnalytics(query: AnalyticsQueryDto): Promise> { + async getRewardAnalytics( + query: AnalyticsQueryDto, + ): Promise>> { const start = Date.now(); const cacheKey = `analytics:rewards:${JSON.stringify(query)}`; - const { data, cached } = await this.getCached(cacheKey, 60 * 5, async () => { - const startDate = this.parseDate(query.startDate); - const endDate = this.parseDate(query.endDate); - const where = ['created_at' >= 'startDate', 'created_at' <= 'endDate'].join(' AND '); - const rewardsDistributed = await this.safeRawSum('reward', 'amount', where); - const stakingRewards = await this.safeRawSum('staking', 'amount', where); - const treasuryBalance = await this.safeRawSum('treasury', 'balance'); - const bountyAllocations = await this.safeRawSum('bounty', 'allocated_amount', where); - const protocolIncentives = await this.safeRawSum('incentive', 'amount', where); - - return { - rewardsDistributed, - stakingRewards, - treasuryBalance, - bountyAllocations, - protocolIncentives, - contributorEarnings: {}, - historicalRewardTrends: [], - }; - }); + const { data, cached } = await this.getCached( + cacheKey, + CACHE_TTL_SECONDS, + async () => { + const startDate = this.parseDate(query.startDate); + const endDate = this.parseDate(query.endDate); + const where = this.dateRangeClause('created_at', startDate, endDate); + + const rewardsDistributed = await this.safeSum( + 'reward', + 'amount', + where || undefined, + ); + const stakingRewards = await this.safeSum( + 'staking', + 'amount', + where || undefined, + ); + const treasuryBalance = await this.safeSum('treasury', 'balance'); + const bountyAllocations = await this.safeSum( + 'bounty', + 'allocated_amount', + where || undefined, + ); + const protocolIncentives = await this.safeSum( + 'incentive', + 'amount', + where || undefined, + ); + + return { + rewardsDistributed, + stakingRewards, + treasuryBalance, + bountyAllocations, + protocolIncentives, + }; + }, + ); return this.wrapResponse(data, cached, Date.now() - start, query); } - async getTrendReporting(query: AnalyticsQueryDto): Promise> { + async getTrendReporting( + query: AnalyticsQueryDto, + ): Promise>> { const start = Date.now(); const cacheKey = `analytics:trends:${JSON.stringify(query)}`; - const { data, cached } = await this.getCached(cacheKey, 60 * 5, async () => { - const startDate = this.parseDate(query.startDate); - const endDate = this.parseDate(query.endDate); - const period = query.period || 'daily'; - - const dailyActivity = await this.getMessageTrends('day', startDate, endDate); - const weeklyActivity = await this.getMessageTrends('week', startDate, endDate); - const monthlyActivity = await this.getMessageTrends('month', startDate, endDate); - const quarterlyActivity = await this.getMessageTrends('quarter', startDate, endDate); - const yearlyGrowth = await this.getMessageTrends('year', startDate, endDate); - - return { - dailyActivity, - weeklyActivity, - monthlyActivity, - quarterlyActivity, - yearlyGrowth: yearlyGrowth.length > 0 ? yearlyGrowth[0] : {}, - }; - }); + const { data, cached } = await this.getCached( + cacheKey, + CACHE_TTL_SECONDS, + async () => { + const startDate = this.parseDate(query.startDate); + const endDate = this.parseDate(query.endDate); + const period: Period = this.resolvePeriod(query.period); + + const activity = await this.trendSeries( + 'claim', + 'created_at', + startDate, + endDate, + ); + + return { + period, + activity, + yearlyGrowth: this.bucketTrend(activity, 'year'), + }; + }, + ); return this.wrapResponse(data, cached, Date.now() - start, query); } - private async getMessageTrends(period: string, start: Date | undefined, end: Date | undefined): Promise { - const startDate = start ? start : new Date(0); - const endDate = end ? end : new Date(); - - // Use Prisma to group messages by createdAt date ranges - // This is a simplified version; in production, use raw SQL for better performance. - const messages = await this.prisma.message.findMany({ - where: { - createdAt: { - gte: startDate, - lte: endDate, - }, - }, - select: { createdAt: true }, - }); - - if (messages.length === 0) return []; + private resolvePeriod(period?: string): Period { + switch (period) { + case 'weekly': + return 'week'; + case 'monthly': + return 'month'; + case 'quarterly': + return 'quarter'; + case 'yearly': + return 'year'; + default: + return 'day'; + } + } - const buckets = new Map(); - for (const m of messages) { - const dt = new Date(m.createdAt); + private bucketTrend( + rows: { period: string; count: number }[], + period: Period, + ): { period: string; count: number }[] { + const buckets = new Map(); + for (const row of rows) { + const date = new Date(`${row.period}T00:00:00.000Z`); + if (Number.isNaN(date.getTime())) continue; let key: string; switch (period) { - case 'day': - key = dt.toISOString().slice(0, 10); - break; - case 'week': - const start = new Date(dt); - const day = dt.getDay(); - start.setDate(day - day); // Sunday - key = start.toISOString().slice(0, 10); + case 'week': { + const weekStart = new Date(date); + weekStart.setUTCDate(weekStart.getUTCDate() - weekStart.getUTCDay()); + key = weekStart.toISOString().slice(0, 10); break; + } case 'month': - key = dt.toISOString().slice(0, 7); + key = row.period.slice(0, 7); break; case 'quarter': - const quarter = Math.floor(dt.getMonth() / 3); - key = `${dt.getFullYear()}-Q/${quarter + 1}`; + key = `${date.getUTCFullYear()}-Q${Math.floor(date.getUTCMonth() / 3) + 1}`; break; case 'year': - key = dt.getFullYear().toString(); + key = String(date.getUTCFullYear()); break; default: - key = dt.toISOString().slice(0, 10); + key = row.period.slice(0, 10); } - buckets.set(key, (buckets.get(key) || 0) + 1); + buckets.set(key, (buckets.get(key) ?? 0) + row.count); } - - return Array.from(buckets.entries()) + return [...buckets.entries()] .map(([key, count]) => ({ period: key, count })) .sort((a, b) => a.period.localeCompare(b.period)); } - private async getTrendArray(table: string, column: string, start: Date | undefined, end: Date | undefined, period?: string): Promise { - try { - const whereClauses = []; - if (start) whereClauses.push(`${column} >= '${start.toISOString()}'`); - if (end) whereClauses.push(`${column} <= '${end.toISOString()}'`); - const where = whereClauses.length > 0 ? `WHERE ${whereClauses.join(' AND ')}` : ''; - const query = `SELECT strftime(${column}, '%Y-%m-%d') as period, COUNT(*) as count FROM "${table}" ${where} GROUP BY period ORDER BY period`; - const result = await this.dataSource.query(query); - return result.map((row) => ({ period: row.period, count: parseInt(row.count, 10) })); - } catch (e) { - return []; - } - } - - async getMonitoringMetrics(): Promise> { + getMonitoringMetrics(): AnalyticsResponse> { const start = Date.now(); - const totalQueries = this.monitoring.cacheHits + this.monitoring.cacheMisses; - const cacheHitRatio = totalQueries ? this.monitoring.cacheHits / totalQueries : 0; - const avgQueryLatency = this.monitoring.reportGenerationCount ? this.monitoring.queryLatencySum / this.monitoring.reportGenerationCount : 0; - - return this.wrapResponse({ - reportGenerationCount: this.monitoring.reportGenerationCount, - queryLatencyMs: avgQueryLatency, - cacheHitRatio: cacheHitRatio, - avgReportGenerationTime: this.monitoring.lastRefreshDuration, - exportRequests: this.monitoring.exportRequests, - failedReportGeneration: this.monitoring.failedReportGeneration, - }, false, Date.now() - start); + const totalQueries = + this.monitoring.cacheHits + this.monitoring.cacheMisses; + const cacheHitRatio = totalQueries + ? this.monitoring.cacheHits / totalQueries + : 0; + const avgQueryLatency = this.monitoring.reportGenerationCount + ? this.monitoring.queryLatencySum / this.monitoring.reportGenerationCount + : 0; + + return this.wrapResponse( + { + reportGenerationCount: this.monitoring.reportGenerationCount, + queryLatencyMs: avgQueryLatency, + cacheHitRatio, + avgReportGenerationTime: this.monitoring.lastRefreshDuration, + exportRequests: this.monitoring.exportRequests, + failedReportGeneration: this.monitoring.failedReportGeneration, + }, + false, + Date.now() - start, + ); } async generateCsvReport(query: AnalyticsQueryDto): Promise { this.monitoring.exportRequests++; const start = Date.now(); - const trendStart = Date.now(); try { - const protocol = (await this.getProtocolStatistics(query)).data; - const contributors = (await this.getContributorAnalytics(query)).data; - const claims = (await this.getClaimAnalytics(query)).data; - const governance = (await this.getGovernanceAnalytics(query)).data; - const rewards = (await this.getRewardAnalytics(query)).data; - const trends = (await this.getTrendReporting(query)).data; - const sections = { - protocol, - contributors, - claims, - governance, - rewards, - trends, + protocol: (await this.getProtocolStatistics(query)).data, + contributors: (await this.getContributorAnalytics(query)).data, + claims: (await this.getClaimAnalytics(query)).data, + governance: (await this.getGovernanceAnalytics(query)).data, + rewards: (await this.getRewardAnalytics(query)).data, + trends: (await this.getTrendReporting(query)).data, }; const csv = this.toCsv(sections); this.monitoring.reportGenerationCount++; - this.monitoring.queryLatencySum += (Date.now() - trendStart); - this.monitoring.lastRefreshDuration = Date.now() - trendStart; + this.monitoring.queryLatencySum += Date.now() - start; + this.monitoring.lastRefreshDuration = Date.now() - start; return csv; - } catch (e) { + } catch (error) { this.monitoring.failedReportGeneration++; - throw e; + throw error; } } - private toCsv(sections: Record): string { + private toCsv(sections: Record>): string { const rows: string[][] = [['Section', 'Metric', 'Value']]; for (const [section, metrics] of Object.entries(sections)) { for (const [key, value] of Object.entries(metrics)) { - if (typeof value === 'object' && value !== null) { - rows.push([section, key, JSON.stringify(value)]); - } else { - rows.push([section, key, String(value)]); - } + rows.push([ + section, + key, + typeof value === 'object' && value !== null + ? JSON.stringify(value) + : String(value), + ]); } } - return rows.map(row => row.map(cell => `"${cell.replace(/"/g, '""')}"`).join(',')).join('\n'); + return rows + .map((row) => + row.map((cell) => `"${cell.replace(/"/g, '""')}"`).join(','), + ) + .join('\n'); } } diff --git a/src/analytics/dto/analytics-query.dto.ts b/src/analytics/dto/analytics-query.dto.ts index 3654cc61..2133840b 100644 --- a/src/analytics/dto/analytics-query.dto.ts +++ b/src/analytics/dto/analytics-query.dto.ts @@ -1,13 +1,19 @@ -import { IsOptional, IsString, IsDateString, IsNumber, IsIn} } from 'class-validator'; -import { Transform } from 'class-transform'; +import { + IsDateString, + IsIn, + IsNumber, + IsOptional, + IsString, +} from 'class-validator'; +import { Transform } from 'class-transformer'; export class AnalyticsQueryDto { @IsOptional() - @isDateString() + @IsDateString() startDate?: string; @IsOptional() - @isDateString() + @IsDateString() endDate?: string; @IsOptional() @@ -31,20 +37,20 @@ export class AnalyticsQueryDto { protocolVersion?: string; @IsOptional() - @IsEnum('daily', 'weekly', 'monthly', 'quarterly', 'yearly') + @IsIn(['daily', 'weekly', 'monthly', 'quarterly', 'yearly']) period?: 'daily' | 'weekly' | 'monthly' | 'quarterly' | 'yearly'; @IsOptional() - @IsEnum('json', 'csv') + @IsIn(['json', 'csv']) format?: 'json' | 'csv'; @IsOptional() - @Transform(({ value }) => parseInt(value, 10)) - @isNumber() + @Transform(({ value }) => parseInt(String(value), 10)) + @IsNumber() page?: number = 1; @IsOptional() - @Transform(({ value }) => parseInt(value, 10)) + @Transform(({ value }) => parseInt(String(value), 10)) @IsNumber() limit?: number = 10; } diff --git a/src/app.module.ts b/src/app.module.ts index 706db169..569183f8 100644 --- a/src/app.module.ts +++ b/src/app.module.ts @@ -44,6 +44,7 @@ import { V2EventsModule } from './v2/events/v2-events.module'; import { V2EvidenceModule } from './v2/evidence/v2-evidence.module'; import { V2VerificationModule } from './v2/verification/v2-verification.module'; import { V2DisputesModule } from './v2/disputes/v2-disputes.module'; +import { ProjectionReadinessModule } from './v2/common/projection-readiness/projection-readiness.module'; import { ProfilerModule } from './profiler/profiler.module'; import { ProfilerInterceptor } from './profiler/profiler.interceptor'; import { HealthModule } from './health/health.module'; @@ -329,6 +330,7 @@ async function createThrottlerStorage( AiAssistantModule, AdminModule, V2EventsModule, + ProjectionReadinessModule, V2EvidenceModule, V2VerificationModule, V2DisputesModule, @@ -363,4 +365,4 @@ async function createThrottlerStorage( }, ], }) -export class AppModule {} \ No newline at end of file +export class AppModule {} diff --git a/src/claims/claims.module.ts b/src/claims/claims.module.ts index ca95d124..b8e87460 100644 --- a/src/claims/claims.module.ts +++ b/src/claims/claims.module.ts @@ -40,13 +40,11 @@ import { IpfsModule } from '../ipfs/ipfs.module'; ClaimProjectorService, ], exports: [ - ClaimResolutionService, ClaimsService, EvidenceService, + EvidenceFlagService, ClaimProjectorService, ], - ], - exports: [ClaimsService, EvidenceService], }) export class ClaimsModule { configure(consumer: MiddlewareConsumer) { @@ -54,4 +52,4 @@ export class ClaimsModule { .apply(EvidenceIntegrityMiddleware) .forRoutes('claims/upload-evidence'); } -} \ No newline at end of file +} diff --git a/src/health/health.service.spec.ts b/src/health/health.service.spec.ts index 45b97244..6cc311a5 100644 --- a/src/health/health.service.spec.ts +++ b/src/health/health.service.spec.ts @@ -69,11 +69,6 @@ describe('HealthService', () => { let dataSource: DataSource; let redisService: RedisService; let queue: Queue; - let jobsService: JobsService; - let notificationService: NotificationService; - let ipfsService: IpfsService; - let blockchainStateService: BlockchainStateService; - let metricsService: MetricsService; beforeEach(async () => { const module: TestingModule = await Test.createTestingModule({ @@ -85,14 +80,11 @@ describe('HealthService', () => { { provide: JobsService, useFactory: mockJobsService }, { provide: NotificationService, useFactory: mockNotificationService }, { provide: IpfsService, useFactory: mockIpfsService }, - feat/be-016-monitoring-api - { provide: BlockchainStateService, useFactory: mockBlockchainStateService }, - { provide: MetricsService, useFactory: mockMetricsService }, { provide: BlockchainStateService, useFactory: mockBlockchainStateService, }, - main + { provide: MetricsService, useFactory: mockMetricsService }, ], }).compile(); @@ -100,16 +92,6 @@ describe('HealthService', () => { dataSource = module.get(DataSource); redisService = module.get(RedisService); queue = module.get('BullQueue_jobs-queue'); - jobsService = module.get(JobsService); - notificationService = module.get(NotificationService); - ipfsService = module.get(IpfsService); - feat/be-016-monitoring-api - blockchainStateService = module.get(BlockchainStateService); - metricsService = module.get(MetricsService); - blockchainStateService = module.get( - BlockchainStateService, - ); - main }); it('should return alive liveness result', () => { diff --git a/src/health/health.service.ts b/src/health/health.service.ts index 7e30c4c9..4b91f59b 100644 --- a/src/health/health.service.ts +++ b/src/health/health.service.ts @@ -98,9 +98,11 @@ export class HealthService { } getDependencyHealth(): DependencyHealthResult { - const dependencies = Array.from(this.lastSuccess.keys()).map((name) => ({ + const dependencies: DependencyStatus[] = Array.from( + this.lastSuccess.keys(), + ).map((name) => ({ name, - status: 'healthy' as HealthStatus, + status: 'healthy', responseTimeMs: 0, lastSuccessfulCheck: this.lastSuccess.get(name), })); @@ -119,7 +121,7 @@ export class HealthService { */ async getIndexerHealth(): Promise { const snapshot = await this.blockchainStateService.getIndexerHealth(); - const status = snapshot.status as HealthStatus; + const status = snapshot.status; return { status, timestamp: new Date().toISOString(), @@ -227,10 +229,7 @@ export class HealthService { } private async checkQueue(): Promise { - feat/be-016-monitoring-api - const counts = await this.jobsQueue.getJobCounts('waiting', 'active', 'completed', 'failed', 'delayed', 'paused'); - this.metricsService.setQueueDepth(this.jobsQueue.name, counts); - await this.jobsQueue.getJobCounts( + const counts = await this.jobsQueue.getJobCounts( 'waiting', 'active', 'completed', @@ -238,7 +237,7 @@ export class HealthService { 'delayed', 'paused', ); - main + this.metricsService.setQueueDepth(this.jobsQueue.name, counts); } private async checkNotifications(): Promise { @@ -263,7 +262,7 @@ export class HealthService { if (typeof state.lastProcessedBlock !== 'number') { throw new Error('Blockchain state is unavailable'); } - feat/be-016-monitoring-api + this.metricsService.setBlockchainIndexingState(state.lastProcessedBlock); // Fail closed if the indexer is degraded per alert thresholds. @@ -271,7 +270,6 @@ export class HealthService { if (health.status === 'unhealthy') { throw new Error('Indexer health is degraded beyond alert thresholds'); } - main } private aggregateServices( diff --git a/src/migrations/1788100000000-AddBlockNumberToV2ProjectDispute.ts b/src/migrations/1788100000000-AddBlockNumberToV2ProjectDispute.ts new file mode 100644 index 00000000..1c2f1bd8 --- /dev/null +++ b/src/migrations/1788100000000-AddBlockNumberToV2ProjectDispute.ts @@ -0,0 +1,44 @@ +import { MigrationInterface, QueryRunner } from 'typeorm'; + +/** + * Adds the chain-native block coordinate to the projected dispute read model. + * + * DisputeProjectorService keyset-paginates and labels data state by + * (blockNumber, logIndex), but the original table only stored eventLogIndex, + * so a projected dispute had no reproducible block coordinate at all. The + * column is added nullable and backfilled from v2_canonical_events -- the + * canonical stream, not from any API-side guess -- so existing rows either + * get a verifiable block or stay explicitly unknown (and are reported as + * OBSERVED rather than finalized). + */ +export class AddBlockNumberToV2ProjectDispute1788100000000 implements MigrationInterface { + name = 'AddBlockNumberToV2ProjectDispute1788100000000'; + + public async up(queryRunner: QueryRunner): Promise { + await queryRunner.query( + `ALTER TABLE "v2_project_dispute" ADD COLUMN "blockNumber" bigint NULL`, + ); + + // Backfill from the canonical event that produced each row. Only rows + // whose originating event is still present in the canonical stream get a + // block; anything else remains NULL (unknown provenance). + await queryRunner.query(` + UPDATE "v2_project_dispute" AS d + SET "blockNumber" = e."blockNumber" + FROM "v2_canonical_events" AS e + WHERE e."txHash" = d."eventTxHash" + AND e."logIndex" = d."eventLogIndex" + `); + + await queryRunner.query( + `CREATE INDEX "idx_v2_project_dispute_block_number" ON "v2_project_dispute" ("blockNumber")`, + ); + } + + public async down(queryRunner: QueryRunner): Promise { + await queryRunner.query(`DROP INDEX "idx_v2_project_dispute_block_number"`); + await queryRunner.query( + `ALTER TABLE "v2_project_dispute" DROP COLUMN "blockNumber"`, + ); + } +} diff --git a/src/v2/common/projection-readiness/projection-readiness.controller.spec.ts b/src/v2/common/projection-readiness/projection-readiness.controller.spec.ts new file mode 100644 index 00000000..537e2657 --- /dev/null +++ b/src/v2/common/projection-readiness/projection-readiness.controller.spec.ts @@ -0,0 +1,122 @@ +import { Test, TestingModule } from '@nestjs/testing'; +import { + BadRequestException, + ServiceUnavailableException, +} from '@nestjs/common'; +import { ProjectionReadinessController } from './projection-readiness.controller'; +import { ProjectionReadinessService } from './projection-readiness.service'; + +describe('ProjectionReadinessController', () => { + let controller: ProjectionReadinessController; + let readiness: jest.Mocked; + + const readyReport = { + ready: true, + status: 'ready' as const, + evaluatedAt: '2026-01-01T00:00:00.000Z', + projectors: [], + }; + + beforeEach(async () => { + const module: TestingModule = await Test.createTestingModule({ + controllers: [ProjectionReadinessController], + providers: [ + { + provide: ProjectionReadinessService, + useValue: { evaluate: jest.fn(), evaluateAll: jest.fn() }, + }, + ], + }).compile(); + + controller = module.get(ProjectionReadinessController); + readiness = module.get( + ProjectionReadinessService, + ) as jest.Mocked; + }); + + it('exposes no write handlers: readiness cannot be set through the API', () => { + const prototype = Object.getPrototypeOf( + controller, + ) as ProjectionReadinessController; + const methodNames = Object.getOwnPropertyNames(prototype).filter( + (name) => name !== 'constructor', + ); + expect(methodNames.sort()).toEqual([ + 'getProjectorReadiness', + 'getReadiness', + ]); + }); + + it('returns the aggregate report when every projector is ready', async () => { + readiness.evaluateAll.mockResolvedValue(readyReport); + + await expect(controller.getReadiness()).resolves.toEqual(readyReport); + }); + + it('reports 503 with the aggregate report when a projector is not ready', async () => { + const notReady = { + ...readyReport, + ready: false, + status: 'not_ready' as const, + }; + readiness.evaluateAll.mockResolvedValue(notReady); + + const error = await controller + .getReadiness() + .catch((thrown: unknown) => thrown); + + expect(error).toBeInstanceOf(ServiceUnavailableException); + expect((error as ServiceUnavailableException).getResponse()).toEqual( + notReady, + ); + }); + + it('rejects an unknown projector name as bad input rather than evaluating it', async () => { + await expect( + controller.getProjectorReadiness('v2-unknown'), + ).rejects.toBeInstanceOf(BadRequestException); + // eslint-disable-next-line @typescript-eslint/unbound-method -- jest mock assertion, not a real unbound call + expect(readiness.evaluate).not.toHaveBeenCalled(); + }); + + it('returns the per-projector verdict when ready', async () => { + const verdict = { + projector: 'v2-evidence', + ready: true, + status: 'ready' as const, + evaluatedAt: '2026-01-01T00:00:00.000Z', + cursor: { blockNumber: '10', logIndex: 0 }, + canonicalHead: { blockNumber: '10', logIndex: 0 }, + pendingEvents: 0, + quarantinedProtocolLogs: 0, + quarantineThreshold: 0, + reasons: [], + checks: [], + }; + readiness.evaluate.mockResolvedValue(verdict); + + await expect( + controller.getProjectorReadiness('v2-evidence'), + ).resolves.toEqual(verdict); + }); + + it('reports 503 for a not-ready projector', async () => { + readiness.evaluate.mockResolvedValue({ + projector: 'v2-disputes', + ready: false, + status: 'not_ready', + evaluatedAt: '2026-01-01T00:00:00.000Z', + cursor: null, + canonicalHead: { blockNumber: '10', logIndex: 0 }, + pendingEvents: 3, + quarantinedProtocolLogs: 0, + quarantineThreshold: 0, + reasons: ['cursor_missing'] as never[], + checks: [], + }); + + await expect( + controller.getProjectorReadiness('v2-disputes'), + ).rejects.toBeInstanceOf(ServiceUnavailableException); + }); +}); diff --git a/src/v2/common/projection-readiness/projection-readiness.controller.ts b/src/v2/common/projection-readiness/projection-readiness.controller.ts new file mode 100644 index 00000000..fcb96c30 --- /dev/null +++ b/src/v2/common/projection-readiness/projection-readiness.controller.ts @@ -0,0 +1,59 @@ +import { + BadRequestException, + Controller, + Get, + Param, + ServiceUnavailableException, +} from '@nestjs/common'; +import { ApiOperation, ApiTags } from '@nestjs/swagger'; +import { Public } from '../../../decorators/public.decorator'; +import { ProjectionReadinessService } from './projection-readiness.service'; +import { V2_PROJECTOR_NAMES, isV2ProjectorName } from './projector-registry'; + +/** + * Operator surface for the Projection Readiness Gate (V2-BE-100). + * + * GET-only by construction: readiness is derived from canonical state and can + * never be set, forced, or cleared through the API. The endpoint is public + * like the health probes because it exposes only chain coordinates, counts, + * and invariant names — never claim content, user data, RPC URLs, or + * credentials. + * + * A not-ready verdict is reported as HTTP 503 with the full report in the + * body, so orchestrators and dashboards see the same failure the read paths + * enforce instead of a green check that disagrees with them. + */ +@ApiTags('V2 Projections') +@Public() +@Controller('v2/projections/readiness') +export class ProjectionReadinessController { + constructor(private readonly readiness: ProjectionReadinessService) {} + + @Get() + @ApiOperation({ + summary: 'Projection readiness for every registered projector', + }) + async getReadiness() { + const report = await this.readiness.evaluateAll(); + if (!report.ready) { + throw new ServiceUnavailableException(report); + } + return report; + } + + @Get(':projector') + @ApiOperation({ summary: 'Projection readiness for a single projector' }) + async getProjectorReadiness(@Param('projector') projector: string) { + if (!isV2ProjectorName(projector)) { + throw new BadRequestException( + `Unknown V2 projector "${projector}" (known: ${V2_PROJECTOR_NAMES.join(', ')})`, + ); + } + + const readiness = await this.readiness.evaluate(projector); + if (!readiness.ready) { + throw new ServiceUnavailableException(readiness); + } + return readiness; + } +} diff --git a/src/v2/common/projection-readiness/projection-readiness.integration.spec.ts b/src/v2/common/projection-readiness/projection-readiness.integration.spec.ts new file mode 100644 index 00000000..21da73e6 --- /dev/null +++ b/src/v2/common/projection-readiness/projection-readiness.integration.spec.ts @@ -0,0 +1,280 @@ +import { Test, TestingModule } from '@nestjs/testing'; +import { TypeOrmModule } from '@nestjs/typeorm'; +import { ConfigService } from '@nestjs/config'; +import { ServiceUnavailableException } from '@nestjs/common'; +import { DataSource } from 'typeorm'; +import { ProjectionReadinessService } from './projection-readiness.service'; +import { ProjectionReadinessReason } from './projection-readiness.types'; +import { V2_PROJECTORS } from './projector-registry'; +import { CanonicalEvent } from '../../events/entities/canonical-event.entity'; +import { ContractArtifact } from '../../events/entities/contract-artifact.entity'; +import { EventQuarantine } from '../../events/entities/event-quarantine.entity'; +import { QuarantineReason } from '../../events/entities/event-quarantine.entity'; +import { ProjectorCursor } from '../entities/projector-cursor.entity'; + +const APPROVED_CONTRACT = '0x' + 'aa'.repeat(20); +const UNAPPROVED_CONTRACT = '0x' + 'bb'.repeat(20); + +describe('ProjectionReadinessService (integration)', () => { + let moduleRef: TestingModule; + let service: ProjectionReadinessService; + let dataSource: DataSource; + + async function seedCanonicalEvent(options: { + eventName: string; + blockNumber: string; + logIndex?: number; + contractAddress?: string; + txHash?: string; + }): Promise { + await dataSource.getRepository(CanonicalEvent).insert({ + chainId: 10, + contractAddress: options.contractAddress ?? APPROVED_CONTRACT, + artifactVersion: 'v1', + eventName: options.eventName, + txHash: + options.txHash ?? + `0x${options.blockNumber.padStart(2, '0')}${'00'.repeat(30)}`, + logIndex: options.logIndex ?? 0, + blockNumber: options.blockNumber, + payload: {} as object, + rawArgs: {} as object, + }); + } + + async function setCursor( + projector: string, + blockNumber: string, + logIndex = 0, + ): Promise { + await dataSource.getRepository(ProjectorCursor).save({ + projectorName: projector, + lastBlockNumber: blockNumber, + lastLogIndex: logIndex, + }); + } + + async function seedArtifact( + contractAddress: string, + isApproved = true, + ): Promise { + await dataSource.getRepository(ContractArtifact).insert({ + chainId: 10, + contractAddress, + artifactVersion: 'v1', + abi: [], + isApproved, + }); + } + + async function seedQuarantine(contractAddress: string): Promise { + await dataSource.getRepository(EventQuarantine).insert({ + chainId: 10, + contractAddress, + txHash: `0x${'ff'.repeat(32)}`, + logIndex: 0, + blockNumber: '1', + topic0: `0x${'11'.repeat(32)}`, + reason: QuarantineReason.UNKNOWN_SIGNATURE, + rawLog: { topics: [], data: '0x' }, + detail: 'topic0 matched no fragment in the approved ABI', + }); + } + + beforeEach(async () => { + moduleRef = await Test.createTestingModule({ + imports: [ + TypeOrmModule.forRoot({ + type: 'sqlite', + database: ':memory:', + // eslint-disable-next-line @typescript-eslint/no-require-imports + driver: require('sqlite3'), + entities: [ + CanonicalEvent, + ContractArtifact, + EventQuarantine, + ProjectorCursor, + ], + synchronize: true, + }), + TypeOrmModule.forFeature([ + CanonicalEvent, + ContractArtifact, + EventQuarantine, + ProjectorCursor, + ]), + ], + providers: [ + ProjectionReadinessService, + { provide: ConfigService, useValue: { get: () => undefined } }, + ], + }).compile(); + + service = moduleRef.get(ProjectionReadinessService); + dataSource = moduleRef.get(DataSource); + }); + + afterEach(async () => { + await moduleRef?.close().catch(() => undefined); + }); + + it('is ready once the projector has consumed the newest handled canonical event', async () => { + await seedCanonicalEvent({ + eventName: 'EvidenceRegistered', + blockNumber: '100', + }); + await seedCanonicalEvent({ + eventName: 'EvidenceReplaced', + blockNumber: '120', + }); + await setCursor(V2_PROJECTORS.EVIDENCE, '120'); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(true); + expect(readiness.canonicalHead).toEqual({ + blockNumber: '120', + logIndex: 0, + }); + }); + + it('is not ready when a newer canonical event has not been projected yet', async () => { + await seedCanonicalEvent({ + eventName: 'EvidenceRegistered', + blockNumber: '100', + }); + await seedCanonicalEvent({ + eventName: 'EvidenceReplaced', + blockNumber: '120', + }); + await setCursor(V2_PROJECTORS.EVIDENCE, '100'); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ProjectionReadinessReason.BACKLOG]); + expect(readiness.pendingEvents).toBe(1); + }); + + it('is not ready when the projector has no cursor while canonical events exist', async () => { + await seedCanonicalEvent({ eventName: 'DisputeRaised', blockNumber: '50' }); + + const readiness = await service.evaluate(V2_PROJECTORS.DISPUTES); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toContain( + ProjectionReadinessReason.CURSOR_MISSING, + ); + }); + + it('only measures the projector against the events it is responsible for', async () => { + // A verification event must not hold the disputes projection back: it is + // not part of that projector's declared contract. + await seedCanonicalEvent({ + eventName: 'VerificationRoundOpened', + blockNumber: '900', + }); + await seedCanonicalEvent({ eventName: 'DisputeRaised', blockNumber: '50' }); + await setCursor(V2_PROJECTORS.DISPUTES, '50'); + + const readiness = await service.evaluate(V2_PROJECTORS.DISPUTES); + + expect(readiness.canonicalHead).toEqual({ blockNumber: '50', logIndex: 0 }); + expect(readiness.ready).toBe(true); + }); + + it('is not ready when an approved protocol contract has an undecodable log', async () => { + await seedCanonicalEvent({ eventName: 'DisputeRaised', blockNumber: '50' }); + await setCursor(V2_PROJECTORS.DISPUTES, '50'); + await seedArtifact(APPROVED_CONTRACT); + await seedQuarantine(APPROVED_CONTRACT); + + const readiness = await service.evaluate(V2_PROJECTORS.DISPUTES); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.QUARANTINE_BACKLOG, + ]); + expect(readiness.quarantinedProtocolLogs).toBe(1); + }); + + it('stays ready when the quarantined log is from an address that is not protocol state', async () => { + await seedCanonicalEvent({ eventName: 'DisputeRaised', blockNumber: '50' }); + await setCursor(V2_PROJECTORS.DISPUTES, '50'); + await seedArtifact(APPROVED_CONTRACT); + await seedQuarantine(UNAPPROVED_CONTRACT); + + const readiness = await service.evaluate(V2_PROJECTORS.DISPUTES); + + expect(readiness.ready).toBe(true); + expect(readiness.quarantinedProtocolLogs).toBe(0); + }); + + it('fails closed when a dependency the gate needs is degraded', async () => { + await seedCanonicalEvent({ + eventName: 'EvidenceRegistered', + blockNumber: '100', + }); + await setCursor(V2_PROJECTORS.EVIDENCE, '100'); + + await dataSource.query('DROP TABLE v2_event_quarantine'); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.EVALUATION_ERROR, + ]); + }); + + it('assertReady rejects with 503 and the failing reasons instead of returning data', async () => { + await seedCanonicalEvent({ + eventName: 'EvidenceRegistered', + blockNumber: '100', + }); + await setCursor(V2_PROJECTORS.EVIDENCE, '99'); + + const error = await service + .assertReady(V2_PROJECTORS.EVIDENCE) + .catch((thrown: unknown) => thrown); + + expect(error).toBeInstanceOf(ServiceUnavailableException); + expect((error as ServiceUnavailableException).getResponse()).toMatchObject({ + projector: 'v2-evidence', + reasons: [ProjectionReadinessReason.BACKLOG], + pendingEvents: 1, + }); + }); + + it('evaluateAll reports every registered projector with its own verdict', async () => { + await seedCanonicalEvent({ + eventName: 'EvidenceRegistered', + blockNumber: '100', + }); + await seedCanonicalEvent({ + eventName: 'DisputeRaised', + blockNumber: '200', + }); + await setCursor(V2_PROJECTORS.EVIDENCE, '100'); + + const report = await service.evaluateAll(); + + expect(report.projectors.map((p) => p.projector).sort()).toEqual([ + 'v2-disputes', + 'v2-evidence', + 'v2-verification', + ]); + // Evidence is caught up; disputes has a canonical event it never projected. + expect( + report.projectors.find((p) => p.projector === 'v2-evidence')?.ready, + ).toBe(true); + expect( + report.projectors.find((p) => p.projector === 'v2-disputes')?.reasons, + ).toEqual([ + ProjectionReadinessReason.CURSOR_MISSING, + ProjectionReadinessReason.BACKLOG, + ]); + expect(report.ready).toBe(false); + expect(report.status).toBe('not_ready'); + }); +}); diff --git a/src/v2/common/projection-readiness/projection-readiness.module.ts b/src/v2/common/projection-readiness/projection-readiness.module.ts new file mode 100644 index 00000000..769534ba --- /dev/null +++ b/src/v2/common/projection-readiness/projection-readiness.module.ts @@ -0,0 +1,31 @@ +import { Module } from '@nestjs/common'; +import { TypeOrmModule } from '@nestjs/typeorm'; +import { CanonicalEvent } from '../../events/entities/canonical-event.entity'; +import { ContractArtifact } from '../../events/entities/contract-artifact.entity'; +import { EventQuarantine } from '../../events/entities/event-quarantine.entity'; +import { ProjectorCursor } from '../entities/projector-cursor.entity'; +import { ProjectionReadinessService } from './projection-readiness.service'; +import { ProjectionReadinessController } from './projection-readiness.controller'; + +/** + * Exposes the Projection Readiness Gate as an injectable guard for V2 read + * paths and as a read-only operator endpoint. + * + * It owns no data of its own: every input is an existing canonical-stream, + * cursor, or quarantine table, so the gate adds an enforcement layer without + * introducing a second source of truth for protocol state. + */ +@Module({ + imports: [ + TypeOrmModule.forFeature([ + CanonicalEvent, + ContractArtifact, + EventQuarantine, + ProjectorCursor, + ]), + ], + controllers: [ProjectionReadinessController], + providers: [ProjectionReadinessService], + exports: [ProjectionReadinessService], +}) +export class ProjectionReadinessModule {} diff --git a/src/v2/common/projection-readiness/projection-readiness.service.spec.ts b/src/v2/common/projection-readiness/projection-readiness.service.spec.ts new file mode 100644 index 00000000..d9ff5945 --- /dev/null +++ b/src/v2/common/projection-readiness/projection-readiness.service.spec.ts @@ -0,0 +1,345 @@ +import { ServiceUnavailableException } from '@nestjs/common'; +import { ConfigService } from '@nestjs/config'; +import { Repository } from 'typeorm'; +import { ProjectionReadinessService } from './projection-readiness.service'; +import { + ProjectionReadinessCheckStatus, + ProjectionReadinessReason, +} from './projection-readiness.types'; +import { V2_PROJECTORS } from './projector-registry'; +import { CanonicalEvent } from '../../events/entities/canonical-event.entity'; +import { EventQuarantine } from '../../events/entities/event-quarantine.entity'; +import { ProjectorCursor } from '../entities/projector-cursor.entity'; + +interface QueryBuilderStub { + select: jest.Mock; + where: jest.Mock; + andWhere: jest.Mock; + orderBy: jest.Mock; + addOrderBy: jest.Mock; + limit: jest.Mock; + innerJoin: jest.Mock; + getOne: jest.Mock; + getCount: jest.Mock; +} + +function createQueryBuilder(options: { + one?: { blockNumber: string; logIndex: number } | null; + count?: number; +}): QueryBuilderStub { + const qb = {} as QueryBuilderStub; + const chainable = (): QueryBuilderStub => qb; + qb.select = jest.fn(chainable); + qb.where = jest.fn(chainable); + qb.andWhere = jest.fn(chainable); + qb.orderBy = jest.fn(chainable); + qb.addOrderBy = jest.fn(chainable); + qb.limit = jest.fn(chainable); + qb.innerJoin = jest.fn(chainable); + qb.getOne = jest.fn().mockResolvedValue(options.one ?? null); + qb.getCount = jest.fn().mockResolvedValue(options.count ?? 0); + return qb; +} + +interface ServiceHarness { + service: ProjectionReadinessService; + canonicalCreate: jest.Mock; + quarantineCreate: jest.Mock; + cursorFindOne: jest.Mock; + headQb: QueryBuilderStub; + pendingQb: QueryBuilderStub; + quarantineQb: QueryBuilderStub; +} + +function createHarness(options: { + head?: { blockNumber: string; logIndex: number } | null; + pending?: number; + quarantine?: number; + cursor?: { lastBlockNumber: string; lastLogIndex: number } | null; + headQueryThrows?: boolean; + quarantineConfig?: string | number | undefined; +}): ServiceHarness { + const headQb = createQueryBuilder({ one: options.head ?? null }); + if (options.headQueryThrows) { + headQb.getOne = jest + .fn() + .mockRejectedValue(new Error('canonical event table is unreadable')); + } + const pendingQb = createQueryBuilder({ count: options.pending ?? 0 }); + const quarantineQb = createQueryBuilder({ count: options.quarantine ?? 0 }); + + // The first builder built on the canonical-event repository reads the head; + // any later one counts the projector's backlog. + let canonicalBuilderCount = 0; + const canonicalCreate = jest.fn(() => { + canonicalBuilderCount += 1; + return canonicalBuilderCount === 1 ? headQb : pendingQb; + }); + const quarantineCreate = jest.fn(() => quarantineQb); + const cursorFindOne = jest.fn().mockResolvedValue(options.cursor ?? null); + + const service = new ProjectionReadinessService( + { + createQueryBuilder: canonicalCreate, + } as unknown as Repository, + { findOne: cursorFindOne } as unknown as Repository, + { + createQueryBuilder: quarantineCreate, + } as unknown as Repository, + { + get: jest.fn().mockReturnValue(options.quarantineConfig), + } as unknown as ConfigService, + ); + + return { + service, + canonicalCreate, + quarantineCreate, + cursorFindOne, + headQb, + pendingQb, + quarantineQb, + }; +} + +describe('ProjectionReadinessService (unit)', () => { + it('is ready when the cursor matches the canonical head and nothing is quarantined', async () => { + const { service } = createHarness({ + head: { blockNumber: '100', logIndex: 4 }, + cursor: { lastBlockNumber: '100', lastLogIndex: 4 }, + pending: 0, + quarantine: 0, + }); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(true); + expect(readiness.status).toBe('ready'); + expect(readiness.reasons).toEqual([]); + expect(readiness.pendingEvents).toBe(0); + expect(readiness.cursor).toEqual({ blockNumber: '100', logIndex: 4 }); + expect(readiness.canonicalHead).toEqual({ + blockNumber: '100', + logIndex: 4, + }); + expect( + readiness.checks.every( + (check) => check.status === ProjectionReadinessCheckStatus.PASS, + ), + ).toBe(true); + }); + + it('is ready when the canonical stream is empty (an empty projection is the accurate answer)', async () => { + const { service } = createHarness({ head: null, cursor: null }); + + const readiness = await service.evaluate(V2_PROJECTORS.DISPUTES); + + expect(readiness.ready).toBe(true); + expect(readiness.canonicalHead).toBeNull(); + expect(readiness.cursor).toBeNull(); + }); + + it('refuses to evaluate an unregistered projector and never queries chain state for it', async () => { + const harness = createHarness({ head: { blockNumber: '1', logIndex: 0 } }); + + const readiness = await harness.service.evaluate('v2-not-a-projector'); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.UNKNOWN_PROJECTOR, + ]); + expect(harness.canonicalCreate).not.toHaveBeenCalled(); + expect(harness.cursorFindOne).not.toHaveBeenCalled(); + }); + + it('is not ready when canonical events exist but the projector has never recorded a cursor', async () => { + const { service } = createHarness({ + head: { blockNumber: '500', logIndex: 2 }, + cursor: null, + pending: 3, + }); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toContain( + ProjectionReadinessReason.CURSOR_MISSING, + ); + expect(readiness.reasons).toContain(ProjectionReadinessReason.BACKLOG); + expect(readiness.pendingEvents).toBe(3); + }); + + it('is not ready when the projector lags the canonical stream', async () => { + const { service } = createHarness({ + head: { blockNumber: '900', logIndex: 0 }, + cursor: { lastBlockNumber: '899', lastLogIndex: 7 }, + pending: 2, + }); + + const readiness = await service.evaluate(V2_PROJECTORS.VERIFICATION); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ProjectionReadinessReason.BACKLOG]); + expect(readiness.pendingEvents).toBe(2); + }); + + it('is not ready when the cursor is ahead of the canonical stream (unsubstantiated progress)', async () => { + const { service } = createHarness({ + head: { blockNumber: '700', logIndex: 0 }, + cursor: { lastBlockNumber: '701', lastLogIndex: 0 }, + }); + + const readiness = await service.evaluate(V2_PROJECTORS.VERIFICATION); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.CURSOR_AHEAD_OF_STREAM, + ]); + }); + + it('is not ready when undecodable protocol logs exceed the allowance', async () => { + const { service } = createHarness({ + head: { blockNumber: '100', logIndex: 0 }, + cursor: { lastBlockNumber: '100', lastLogIndex: 0 }, + quarantine: 1, + quarantineConfig: undefined, + }); + + const readiness = await service.evaluate(V2_PROJECTORS.DISPUTES); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.QUARANTINE_BACKLOG, + ]); + expect(readiness.quarantinedProtocolLogs).toBe(1); + expect(readiness.quarantineThreshold).toBe(0); + }); + + it('honours a widened quarantine allowance and still reports the counts', async () => { + const { service } = createHarness({ + head: { blockNumber: '100', logIndex: 0 }, + cursor: { lastBlockNumber: '100', lastLogIndex: 0 }, + quarantine: 2, + quarantineConfig: '5', + }); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(true); + expect(readiness.quarantinedProtocolLogs).toBe(2); + expect(readiness.quarantineThreshold).toBe(5); + }); + + it('fails closed with evaluation_error when a dependency cannot be read', async () => { + const { service } = createHarness({ headQueryThrows: true }); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.EVALUATION_ERROR, + ]); + expect(readiness.checks[0].detail).toMatch(/not asserted/); + }); + + it('fails closed when the quarantine allowance is not a valid configuration', async () => { + const { service } = createHarness({ + head: { blockNumber: '100', logIndex: 0 }, + cursor: { lastBlockNumber: '100', lastLogIndex: 0 }, + quarantineConfig: 'not-a-number', + }); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.EVALUATION_ERROR, + ]); + }); + + it('rejects a negative quarantine allowance rather than widening the gate', async () => { + const { service } = createHarness({ + head: { blockNumber: '100', logIndex: 0 }, + cursor: { lastBlockNumber: '100', lastLogIndex: 0 }, + quarantineConfig: '-1', + }); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.EVALUATION_ERROR, + ]); + }); + + it('treats a malformed cursor coordinate as an integrity failure, not as block zero', async () => { + const { service } = createHarness({ + head: { blockNumber: '100', logIndex: 0 }, + cursor: { lastBlockNumber: 'not-a-block', lastLogIndex: 0 }, + }); + + const readiness = await service.evaluate(V2_PROJECTORS.EVIDENCE); + + expect(readiness.ready).toBe(false); + expect(readiness.reasons).toEqual([ + ProjectionReadinessReason.EVALUATION_ERROR, + ]); + }); + + it('evaluates every registered projector and only reports ready when all are ready', async () => { + const harness = createHarness({ + head: { blockNumber: '100', logIndex: 0 }, + cursor: { lastBlockNumber: '100', lastLogIndex: 0 }, + }); + + const report = await harness.service.evaluateAll(); + + expect(report.projectors).toHaveLength(3); + expect(report.ready).toBe(true); + expect(report.status).toBe('ready'); + }); + + describe('assertReady', () => { + it('resolves for a ready projection', async () => { + const { service } = createHarness({ + head: { blockNumber: '100', logIndex: 0 }, + cursor: { lastBlockNumber: '100', lastLogIndex: 0 }, + }); + + await expect( + service.assertReady(V2_PROJECTORS.EVIDENCE), + ).resolves.toBeUndefined(); + }); + + it('throws 503 with actionable detail instead of returning unverified data', async () => { + const { service } = createHarness({ + head: { blockNumber: '900', logIndex: 0 }, + cursor: { lastBlockNumber: '899', lastLogIndex: 0 }, + pending: 4, + }); + + await expect( + service.assertReady(V2_PROJECTORS.EVIDENCE), + ).rejects.toMatchObject({ + status: 503, + response: { + error: 'projection_not_ready', + projector: 'v2-evidence', + reasons: [ProjectionReadinessReason.BACKLOG], + pendingEvents: 4, + }, + }); + }); + + it('exposes the failure through a ServiceUnavailableException', async () => { + const { service } = createHarness({ + head: { blockNumber: '1', logIndex: 0 }, + cursor: null, + }); + + await expect( + service.assertReady(V2_PROJECTORS.EVIDENCE), + ).rejects.toBeInstanceOf(ServiceUnavailableException); + }); + }); +}); diff --git a/src/v2/common/projection-readiness/projection-readiness.service.ts b/src/v2/common/projection-readiness/projection-readiness.service.ts new file mode 100644 index 00000000..bce2a62f --- /dev/null +++ b/src/v2/common/projection-readiness/projection-readiness.service.ts @@ -0,0 +1,488 @@ +import { + Injectable, + Logger, + ServiceUnavailableException, +} from '@nestjs/common'; +import { InjectRepository } from '@nestjs/typeorm'; +import { ConfigService } from '@nestjs/config'; +import { Repository } from 'typeorm'; +import { CanonicalEvent } from '../../events/entities/canonical-event.entity'; +import { EventQuarantine } from '../../events/entities/event-quarantine.entity'; +import { ContractArtifact } from '../../events/entities/contract-artifact.entity'; +import { ProjectorCursor } from '../entities/projector-cursor.entity'; +import { + PROJECTOR_HANDLED_EVENTS, + V2_PROJECTOR_NAMES, + V2ProjectorName, + isV2ProjectorName, +} from './projector-registry'; +import { + ProjectionOrderKey, + ProjectionReadiness, + ProjectionReadinessCheck, + ProjectionReadinessCheckStatus, + ProjectionReadinessReason, + ProjectionReadinessReport, +} from './projection-readiness.types'; + +/** Environment keys the gate reads. Documented in .env.example and the runbook. */ +export const PROJECTION_READINESS_ENV = { + quarantineMaxPending: 'PROJECTION_READINESS_QUARANTINE_MAX_PENDING', +} as const; + +/** + * Default: zero undecodable logs for an approved protocol contract may exist + * while the projection is served as canonical state. + */ +export const DEFAULT_QUARANTINE_MAX_PENDING = 0; + +/** Error code surfaced in the 503 body so callers can branch on it. */ +export const PROJECTION_NOT_READY_ERROR = 'projection_not_ready'; + +const CHECK_PROJECTOR_REGISTERED = 'projector_registered'; +const CHECK_CANONICAL_STREAM = 'canonical_event_stream'; +const CHECK_CURSOR_CONSISTENCY = 'projector_cursor_consistency'; +const CHECK_CATCH_UP = 'projector_catch_up'; +const CHECK_QUARANTINE = 'protocol_log_quarantine'; + +function toBigInt(value: string | number, field: string): bigint { + const normalized = + typeof value === 'number' ? Math.trunc(value) : value.trim(); + try { + return BigInt(normalized); + } catch { + // A non-numeric block coordinate means the row cannot be trusted as a + // chain-native ordering key. That is an integrity failure, not a + // "no data yet" case, so it must surface as unready rather than be + // coerced to 0. + throw new Error(`malformed ${field}: ${String(value)}`); + } +} + +function isAfter( + candidate: ProjectionOrderKey, + reference: ProjectionOrderKey, +): boolean { + const candidateBlock = toBigInt( + candidate.blockNumber, + 'candidate.blockNumber', + ); + const referenceBlock = toBigInt( + reference.blockNumber, + 'reference.blockNumber', + ); + if (candidateBlock !== referenceBlock) return candidateBlock > referenceBlock; + return candidate.logIndex > reference.logIndex; +} + +/** + * Projection Readiness Gate (V2-BE-100). + * + * TruthBounty treats the deployed Optimism/EVM contracts and their finalized + * canonical events as protocol authority; the API is a deterministic + * indexing, projection and delivery layer. A projected read model is only + * allowed to be served while the API can *prove* it still reproduces that + * authority. This service is that proof, or the refusal to give one. + * + * Invariants (all fail closed — an unprovable invariant is a failure, never a + * warning that is served anyway): + * + * I1 Evaluation is total. Any error raised while evaluating (dependency + * unavailable, malformed row, invalid configuration) is reported as + * `evaluation_error` with `ready: false`. There is no unchecked path + * that returns ready. + * I2 Only registered projectors can be ready. An unknown projector name has + * no declared event contract, so nothing can be asserted about it. + * I3 Canonical events for a projector imply a cursor. If events exist that + * the projector must consume but no cursor row was ever written, the + * projection may be empty or arbitrarily stale and is not ready. + * I4 The cursor never lags the canonical stream. A cursor behind the newest + * handled event is a backlog; a cursor ahead of it means the projection + * claims progress the canonical stream cannot substantiate. + * I5 Undecodable logs for an approved protocol contract block readiness. + * Quarantine entries prove the projection is knowingly incomplete, so + * serving it as canonical state would fabricate protocol truth. + * Quarantine entries from *unapproved* addresses are not protocol state + * and therefore cannot invalidate the projection. + * + * The gate never writes. It only reads the canonical stream, the projector + * cursors, and the quarantine table, so evaluating readiness can never itself + * mutate or advance protocol-derived state. + */ +@Injectable() +export class ProjectionReadinessService { + private readonly logger = new Logger(ProjectionReadinessService.name); + + constructor( + @InjectRepository(CanonicalEvent) + private readonly canonicalEvents: Repository, + @InjectRepository(ProjectorCursor) + private readonly projectorCursors: Repository, + @InjectRepository(EventQuarantine) + private readonly quarantine: Repository, + private readonly config: ConfigService, + ) {} + + /** Evaluate every registered projector. */ + async evaluateAll(): Promise { + const projectors = await Promise.all( + V2_PROJECTOR_NAMES.map((name) => this.evaluate(name)), + ); + const ready = projectors.every((projection) => projection.ready); + return { + ready, + status: ready ? 'ready' : 'not_ready', + evaluatedAt: new Date().toISOString(), + projectors, + }; + } + + /** + * Evaluate a single projector by name. + * + * Never throws: every failure, including an unexpected one, is returned as + * an explicit not-ready verdict with a reason. + */ + async evaluate(projector: string): Promise { + const evaluatedAt = new Date().toISOString(); + try { + if (!isV2ProjectorName(projector)) { + return this.verdict({ + projector, + evaluatedAt, + cursor: null, + canonicalHead: null, + pendingEvents: 0, + quarantinedProtocolLogs: 0, + quarantineThreshold: 0, + reasons: [ProjectionReadinessReason.UNKNOWN_PROJECTOR], + checks: [ + this.check( + CHECK_PROJECTOR_REGISTERED, + false, + `"${projector}" is not a registered V2 projector (known: ${V2_PROJECTOR_NAMES.join(', ')})`, + ), + ], + }); + } + + return await this.evaluateKnownProjector(projector, evaluatedAt); + } catch (error) { + const detail = error instanceof Error ? error.message : String(error); + this.logger.error( + `Projection readiness evaluation failed for "${projector}": ${detail}`, + ); + return this.verdict({ + projector, + evaluatedAt, + cursor: null, + canonicalHead: null, + pendingEvents: 0, + quarantinedProtocolLogs: 0, + quarantineThreshold: 0, + reasons: [ProjectionReadinessReason.EVALUATION_ERROR], + checks: [ + this.check( + 'readiness_evaluation', + false, + `readiness could not be evaluated and is therefore not asserted: ${detail}`, + ), + ], + }); + } + } + + /** + * Fail-closed accessor for read paths. Resolves only when the projection is + * provably caught up with canonical events; otherwise throws 503 so the API + * reports an unavailable projection instead of answering from state it + * cannot vouch for. + */ + async assertReady(projector: V2ProjectorName): Promise { + const readiness = await this.evaluate(projector); + if (readiness.ready) return; + + throw new ServiceUnavailableException({ + statusCode: 503, + error: PROJECTION_NOT_READY_ERROR, + message: + `Projection "${readiness.projector}" is not ready to serve protocol-derived reads: ` + + readiness.reasons.join(', '), + projector: readiness.projector, + reasons: readiness.reasons, + checks: readiness.checks, + cursor: readiness.cursor, + canonicalHead: readiness.canonicalHead, + pendingEvents: readiness.pendingEvents, + quarantinedProtocolLogs: readiness.quarantinedProtocolLogs, + quarantineThreshold: readiness.quarantineThreshold, + evaluatedAt: readiness.evaluatedAt, + }); + } + + private async evaluateKnownProjector( + projector: V2ProjectorName, + evaluatedAt: string, + ): Promise { + const eventNames = [...PROJECTOR_HANDLED_EVENTS[projector]]; + const quarantineThreshold = this.resolveQuarantineThreshold(); + + const canonicalHead = await this.findCanonicalHead(eventNames); + const cursorRow = await this.projectorCursors.findOne({ + where: { projectorName: projector }, + }); + + const cursor: ProjectionOrderKey | null = cursorRow + ? { + blockNumber: String(cursorRow.lastBlockNumber), + logIndex: cursorRow.lastLogIndex, + } + : null; + + const reasons: ProjectionReadinessReason[] = []; + const checks: ProjectionReadinessCheck[] = []; + + let pendingEvents = 0; + let cursorAheadOfStream = false; + + if (canonicalHead && cursor) { + if (isAfter(cursor, canonicalHead)) { + cursorAheadOfStream = true; + } else { + pendingEvents = await this.countPendingEvents(eventNames, cursor); + } + } else if (canonicalHead && !cursor) { + pendingEvents = await this.countPendingEvents(eventNames, null); + } + + // I3 — a projector with work to do must have a cursor. + if (canonicalHead && !cursor) { + reasons.push(ProjectionReadinessReason.CURSOR_MISSING); + checks.push( + this.check( + CHECK_CURSOR_CONSISTENCY, + false, + `${pendingEvents} canonical event(s) exist for ${projector} but it has never recorded a cursor`, + ), + ); + } else if (cursorAheadOfStream && cursor) { + // I4 — cannot claim progress the canonical stream does not contain. + reasons.push(ProjectionReadinessReason.CURSOR_AHEAD_OF_STREAM); + checks.push( + this.check( + CHECK_CURSOR_CONSISTENCY, + false, + `cursor ${cursor.blockNumber}:${cursor.logIndex} is ahead of canonical head ` + + `${canonicalHead?.blockNumber}:${canonicalHead?.logIndex}`, + ), + ); + } else { + checks.push( + this.check( + CHECK_CURSOR_CONSISTENCY, + true, + cursor + ? `cursor ${cursor.blockNumber}:${cursor.logIndex} is within the canonical stream` + : 'no canonical events to project yet', + ), + ); + } + + // I4 — the projection must be caught up with the canonical stream. + if (pendingEvents > 0) { + reasons.push(ProjectionReadinessReason.BACKLOG); + checks.push( + this.check( + CHECK_CATCH_UP, + false, + `${pendingEvents} canonical event(s) are not yet projected`, + ), + ); + } else { + checks.push( + this.check( + CHECK_CATCH_UP, + true, + 'projection is caught up with canonical events', + ), + ); + } + + checks.push( + this.check( + CHECK_CANONICAL_STREAM, + Boolean(canonicalHead) || Boolean(cursor), + canonicalHead + ? `canonical head ${canonicalHead.blockNumber}:${canonicalHead.logIndex}` + : cursor + ? 'cursor exists but the canonical stream is empty' + : 'canonical stream is empty', + ), + ); + + // I5 — undecodable logs from approved protocol contracts invalidate the projection. + const quarantinedProtocolLogs = await this.countQuarantinedProtocolLogs(); + if (quarantinedProtocolLogs > quarantineThreshold) { + reasons.push(ProjectionReadinessReason.QUARANTINE_BACKLOG); + checks.push( + this.check( + CHECK_QUARANTINE, + false, + `${quarantinedProtocolLogs} undecodable log(s) from approved protocol contracts ` + + `(allowed: ${quarantineThreshold}); the projection is knowingly incomplete`, + ), + ); + } else { + checks.push( + this.check( + CHECK_QUARANTINE, + true, + `${quarantinedProtocolLogs} undecodable log(s) from approved protocol contracts ` + + `(allowed: ${quarantineThreshold})`, + ), + ); + } + + checks.unshift( + this.check( + CHECK_PROJECTOR_REGISTERED, + true, + `${projector} is a registered V2 projector`, + ), + ); + + return this.verdict({ + projector, + evaluatedAt, + cursor, + canonicalHead, + pendingEvents, + quarantinedProtocolLogs, + quarantineThreshold, + reasons, + checks, + }); + } + + /** + * Newest canonical event the projector is responsible for consuming. + * Ordered by the protocol's own (blockNumber, logIndex) coordinates. + */ + private async findCanonicalHead( + eventNames: string[], + ): Promise { + const row = await this.canonicalEvents + .createQueryBuilder('e') + .select(['e.blockNumber', 'e.logIndex']) + .where('e.eventName IN (:...eventNames)', { eventNames }) + .orderBy('e.blockNumber', 'DESC') + .addOrderBy('e.logIndex', 'DESC') + .limit(1) + .getOne(); + + if (!row) return null; + return { + blockNumber: String(row.blockNumber), + logIndex: row.logIndex, + }; + } + + /** Canonical events for this projector strictly after `after` (null = from genesis). */ + private async countPendingEvents( + eventNames: string[], + after: ProjectionOrderKey | null, + ): Promise { + const qb = this.canonicalEvents + .createQueryBuilder('e') + .where('e.eventName IN (:...eventNames)', { eventNames }); + + if (after) { + qb.andWhere( + '(e.blockNumber > :blockNumber OR (e.blockNumber = :blockNumber AND e.logIndex > :logIndex))', + { blockNumber: after.blockNumber, logIndex: after.logIndex }, + ); + } + + return qb.getCount(); + } + + /** + * Quarantined logs whose address is an *approved* artifact. These are real + * protocol logs this pipeline failed to decode (unknown signature, artifact + * drift, decode error), so their events are missing from every projection. + */ + private async countQuarantinedProtocolLogs(): Promise { + return this.quarantine + .createQueryBuilder('q') + .innerJoin( + ContractArtifact, + 'a', + 'a.chainId = q.chainId AND a.contractAddress = q.contractAddress', + ) + .where('a.isApproved = :approved', { approved: true }) + .getCount(); + } + + /** + * Invalid or absent configuration falls back to the strict default; a + * *malformed* value throws so it is reported as `evaluation_error` rather + * than silently widened to permissive behavior. + */ + private resolveQuarantineThreshold(): number { + const raw = this.config.get( + PROJECTION_READINESS_ENV.quarantineMaxPending, + ); + if (raw === undefined || raw === null || raw === '') { + return DEFAULT_QUARANTINE_MAX_PENDING; + } + const parsed = + typeof raw === 'number' ? raw : Number.parseInt(String(raw).trim(), 10); + if (!Number.isInteger(parsed) || parsed < 0) { + throw new Error( + `${PROJECTION_READINESS_ENV.quarantineMaxPending} must be a non-negative integer, got "${String(raw)}"`, + ); + } + return parsed; + } + + private check( + name: string, + passed: boolean, + detail: string, + ): ProjectionReadinessCheck { + return { + name, + status: passed + ? ProjectionReadinessCheckStatus.PASS + : ProjectionReadinessCheckStatus.FAIL, + detail, + }; + } + + private verdict(input: { + projector: string; + evaluatedAt: string; + cursor: ProjectionOrderKey | null; + canonicalHead: ProjectionOrderKey | null; + pendingEvents: number; + quarantinedProtocolLogs: number; + quarantineThreshold: number; + reasons: ProjectionReadinessReason[]; + checks: ProjectionReadinessCheck[]; + }): ProjectionReadiness { + const ready = input.reasons.length === 0; + return { + projector: input.projector, + ready, + status: ready ? 'ready' : 'not_ready', + evaluatedAt: input.evaluatedAt, + cursor: input.cursor, + canonicalHead: input.canonicalHead, + pendingEvents: input.pendingEvents, + quarantinedProtocolLogs: input.quarantinedProtocolLogs, + quarantineThreshold: input.quarantineThreshold, + reasons: input.reasons, + checks: input.checks, + }; + } +} diff --git a/src/v2/common/projection-readiness/projection-readiness.types.ts b/src/v2/common/projection-readiness/projection-readiness.types.ts new file mode 100644 index 00000000..6ef8160b --- /dev/null +++ b/src/v2/common/projection-readiness/projection-readiness.types.ts @@ -0,0 +1,85 @@ +/** + * Result contract for the Projection Readiness Gate (V2-BE-100). + * + * The gate answers exactly one question: may this projected read model be + * served as a reproduction of canonical Optimism/EVM state right now? + * + * Two properties are load-bearing and are part of the contract, not an + * implementation detail: + * + * 1. The evaluation is total: every possible outcome is representable, and + * "unknown" is never collapsed into "ready". A thrown error, an + * unreadable dependency, or an unrecognized projector all resolve to + * `ready: false` with an explicit reason. + * 2. The result is self-describing: it carries the evidence behind the + * verdict (cursor, canonical head, pending/quarantined counts, per-check + * detail) so an operator can act on a failure without reading logs or + * reaching for a debugger. + */ + +/** Status of a single invariant check inside one evaluation. */ +export enum ProjectionReadinessCheckStatus { + PASS = 'pass', + FAIL = 'fail', +} + +/** + * Machine-readable reason a projection is not ready. Each value maps to one + * documented failure mode and one recovery procedure; see + * docs/PROJECTION_READINESS_GATE.md. + */ +export enum ProjectionReadinessReason { + /** The gate itself could not complete. Fail closed: unknown is not ready. */ + EVALUATION_ERROR = 'evaluation_error', + /** The projector name is not in the V2 projector registry. */ + UNKNOWN_PROJECTOR = 'unknown_projector', + /** Canonical events exist for this projector but no cursor has ever been written. */ + CURSOR_MISSING = 'cursor_missing', + /** The cursor is behind the newest canonical event the projector must consume. */ + BACKLOG = 'backlog', + /** The cursor is ahead of the canonical event stream (impossible/truncated history). */ + CURSOR_AHEAD_OF_STREAM = 'cursor_ahead_of_stream', + /** Undecodable logs exist for an approved protocol contract: the projection is knowingly incomplete. */ + QUARANTINE_BACKLOG = 'quarantine_backlog', +} + +/** One invariant check, reported whether it passed or failed. */ +export interface ProjectionReadinessCheck { + name: string; + status: ProjectionReadinessCheckStatus; + detail: string; +} + +/** Chain-native ordering coordinate: (blockNumber, logIndex). */ +export interface ProjectionOrderKey { + blockNumber: string; + logIndex: number; +} + +/** The gate's verdict for a single projector. */ +export interface ProjectionReadiness { + projector: string; + ready: boolean; + status: 'ready' | 'not_ready'; + evaluatedAt: string; + /** The projector's consumed-through coordinate, or null if it has never run. */ + cursor: ProjectionOrderKey | null; + /** Newest canonical event coordinate this projector must consume, or null if none exist. */ + canonicalHead: ProjectionOrderKey | null; + /** Canonical events strictly after the cursor: the projection's backlog. */ + pendingEvents: number; + /** Quarantined logs belonging to approved protocol contracts. */ + quarantinedProtocolLogs: number; + /** Threshold in force for {@link quarantinedProtocolLogs} during this evaluation. */ + quarantineThreshold: number; + reasons: ProjectionReadinessReason[]; + checks: ProjectionReadinessCheck[]; +} + +/** Aggregate verdict across every registered projector. */ +export interface ProjectionReadinessReport { + ready: boolean; + status: 'ready' | 'not_ready'; + evaluatedAt: string; + projectors: ProjectionReadiness[]; +} diff --git a/src/v2/common/projection-readiness/projector-registry.ts b/src/v2/common/projection-readiness/projector-registry.ts new file mode 100644 index 00000000..295ffb1e --- /dev/null +++ b/src/v2/common/projection-readiness/projector-registry.ts @@ -0,0 +1,62 @@ +/** + * Single source of truth for the V2 projection pipeline. + * + * The gate (V2-BE-100) decides whether an operator or a read endpoint may + * trust a projected read model. That decision is only correct if the gate and + * the projector agree on (a) the projector's identity and (b) exactly which + * canonical events that projector is responsible for consuming. Duplicating + * those two facts in the gate would let them drift silently: a projector that + * started handling a new event name without the gate learning about it would + * keep reporting "ready" while quietly falling behind. + * + * Projectors therefore import their name and handled-event list from here, + * and the gate reads the same constants. + */ + +export const V2_PROJECTORS = { + EVIDENCE: 'v2-evidence', + VERIFICATION: 'v2-verification', + DISPUTES: 'v2-disputes', +} as const; + +export type V2ProjectorName = + (typeof V2_PROJECTORS)[keyof typeof V2_PROJECTORS]; + +/** Every projector the gate is allowed to evaluate, in stable order. */ +export const V2_PROJECTOR_NAMES: readonly V2ProjectorName[] = [ + V2_PROJECTORS.EVIDENCE, + V2_PROJECTORS.VERIFICATION, + V2_PROJECTORS.DISPUTES, +]; + +/** + * Canonical event names each projector consumes. These are the protocol's own + * event names (see event-schema-registry.ts), never renamed by the API. + */ +export const PROJECTOR_HANDLED_EVENTS: Record< + V2ProjectorName, + readonly string[] +> = { + [V2_PROJECTORS.EVIDENCE]: [ + 'EvidenceRegistered', + 'EvidenceReplaced', + 'EvidenceRemoved', + ], + [V2_PROJECTORS.VERIFICATION]: [ + 'VerificationRoundOpened', + 'PositionCommitted', + ], + [V2_PROJECTORS.DISPUTES]: [ + 'DisputeRaised', + 'DisputeResolved', + 'DisputeExpired', + ], +}; + +/** + * Narrowing guard. An unrecognized name must never be treated as "ready": + * the gate cannot assert anything about a projection it has no contract for. + */ +export function isV2ProjectorName(value: string): value is V2ProjectorName { + return (V2_PROJECTOR_NAMES as readonly string[]).includes(value); +} diff --git a/src/v2/disputes/disputes-projector.service.integration.spec.ts b/src/v2/disputes/disputes-projector.service.integration.spec.ts index 04ca994f..9b820b53 100644 --- a/src/v2/disputes/disputes-projector.service.integration.spec.ts +++ b/src/v2/disputes/disputes-projector.service.integration.spec.ts @@ -1,5 +1,7 @@ import { Test, TestingModule } from '@nestjs/testing'; import { TypeOrmModule } from '@nestjs/typeorm'; +import { ConfigService } from '@nestjs/config'; +import { ServiceUnavailableException } from '@nestjs/common'; import { DataSource } from 'typeorm'; import { DisputesProjectorService } from './disputes-projector.service'; import { DisputesQueryService } from './disputes-query.service'; @@ -14,6 +16,10 @@ import { } from '../common/entities/indexing-anomaly.entity'; import { CanonicalEvent } from '../events/entities/canonical-event.entity'; import { CanonicalEventQueryService } from '../events/canonical-event-query.service'; +import { EventCheckpoint } from '../events/entities/event-checkpoint.entity'; +import { ContractArtifact } from '../events/entities/contract-artifact.entity'; +import { EventQuarantine } from '../events/entities/event-quarantine.entity'; +import { ProjectionReadinessService } from '../common/projection-readiness/projection-readiness.service'; describe('DisputesProjectorService (integration)', () => { let moduleRef: TestingModule; @@ -53,6 +59,9 @@ describe('DisputesProjectorService (integration)', () => { ProjectDispute, ProjectorCursor, IndexingAnomaly, + EventCheckpoint, + ContractArtifact, + EventQuarantine, ], synchronize: true, }), @@ -61,12 +70,17 @@ describe('DisputesProjectorService (integration)', () => { ProjectorCursor, IndexingAnomaly, CanonicalEvent, + EventCheckpoint, + ContractArtifact, + EventQuarantine, ]), ], providers: [ DisputesProjectorService, DisputesQueryService, CanonicalEventQueryService, + ProjectionReadinessService, + { provide: ConfigService, useValue: { get: () => undefined } }, ], }).compile(); @@ -97,6 +111,32 @@ describe('DisputesProjectorService (integration)', () => { expect(dispute.claimId).toBe(claimId); expect(dispute.originalRoundId).toBe(roundId); expect(dispute.challengeBond).toBe('5000'); + // The projected row keeps its chain-native coordinate so data state and + // keyset pagination stay reproducible from canonical events. + expect(String(dispute.blockNumber)).toBe('100'); + }); + + it('fails closed instead of serving a stale dispute projection when canonical events are unprojected', async () => { + await seedEvent({ + eventName: 'DisputeRaised', + txHash: '0x' + '01'.repeat(32), + blockNumber: '100', + }); + await seedEvent({ + eventName: 'DisputeResolved', + txHash: '0x' + '02'.repeat(32), + blockNumber: '200', + payload: { outcome: 'upheld' }, + }); + + await projector.processNewEvents(); + await dataSource + .getRepository(ProjectorCursor) + .update({ projectorName: 'v2-disputes' }, { lastBlockNumber: '100' }); + + await expect( + queryService.getByOriginalRound(claimId, roundId), + ).rejects.toBeInstanceOf(ServiceUnavailableException); }); it('DisputeResolved transitions RAISED -> RESOLVED and stores the verbatim outcome', async () => { diff --git a/src/v2/disputes/disputes-projector.service.ts b/src/v2/disputes/disputes-projector.service.ts index 039bf0ea..31a5462d 100644 --- a/src/v2/disputes/disputes-projector.service.ts +++ b/src/v2/disputes/disputes-projector.service.ts @@ -12,13 +12,19 @@ import { IndexingAnomaly, IndexingAnomalyKind, } from '../common/entities/indexing-anomaly.entity'; +import { + PROJECTOR_HANDLED_EVENTS, + V2_PROJECTORS, + V2ProjectorName, +} from '../common/projection-readiness/projector-registry'; -const PROJECTOR_NAME = 'v2-disputes'; +const PROJECTOR_NAME: V2ProjectorName = V2_PROJECTORS.DISPUTES; const PG_UNIQUE_VIOLATION = '23505'; -const HANDLED_EVENT_NAMES = [ - 'DisputeRaised', - 'DisputeResolved', - 'DisputeExpired', +// Name and handled-event list come from the projector registry so the +// readiness gate (V2-BE-100) can never disagree with this projector about +// which canonical events it is responsible for consuming. +const HANDLED_EVENT_NAMES: string[] = [ + ...PROJECTOR_HANDLED_EVENTS[PROJECTOR_NAME], ]; export interface ProjectorRunSummary { @@ -151,6 +157,7 @@ export class DisputesProjectorService { deadline: readDate(event.payload, 'deadline'), eventTxHash: event.txHash, eventLogIndex: event.logIndex, + blockNumber: event.blockNumber, }); return 'applied'; } catch (err) { @@ -209,6 +216,7 @@ export class DisputesProjectorService { dispute.status = nextStatus; dispute.eventTxHash = event.txHash; dispute.eventLogIndex = event.logIndex; + dispute.blockNumber = event.blockNumber; if (nextStatus === DisputeStatus.RESOLVED) { dispute.resolvedOutcome = readString(event.payload, 'outcome'); } diff --git a/src/v2/disputes/disputes-query.service.ts b/src/v2/disputes/disputes-query.service.ts index 6bf32c34..b0e4c93c 100644 --- a/src/v2/disputes/disputes-query.service.ts +++ b/src/v2/disputes/disputes-query.service.ts @@ -1,10 +1,20 @@ -import { Injectable, NotFoundException, BadRequestException } from '@nestjs/common'; +import { + Injectable, + NotFoundException, + BadRequestException, +} from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import { ProjectDispute } from './entities/project-dispute.entity'; import { EventCheckpoint } from '../events/entities/event-checkpoint.entity'; import { DataState } from '../common/data-state.enum'; -import { CursorPage, encodeCursor, decodeCursor } from '../common/cursor-pagination'; +import { + CursorPage, + encodeCursor, + decodeCursor, +} from '../common/cursor-pagination'; +import { ProjectionReadinessService } from '../common/projection-readiness/projection-readiness.service'; +import { V2_PROJECTORS } from '../common/projection-readiness/projector-registry'; @Injectable() export class DisputesQueryService { @@ -13,15 +23,27 @@ export class DisputesQueryService { private readonly disputeRepo: Repository, @InjectRepository(EventCheckpoint) private readonly checkpointRepo: Repository, + private readonly readiness: ProjectionReadinessService, ) {} /** - * Calculate the data state for a block number based on chain's safe and finalized blocks + * Calculate the data state for a block number based on chain's safe and finalized blocks. + * A row whose originating block is unknown is reported as OBSERVED: finality + * is never asserted for data whose provenance cannot be established. */ - private async calculateDataState(blockNumber: string): Promise { - // Get the latest checkpoint (assuming single chain for simplicity) - const checkpoint = await this.checkpointRepo.findOne({ + private async calculateDataState( + blockNumber: string | null, + ): Promise { + if (blockNumber === null) { + return DataState.OBSERVED; + } + + // Get the latest checkpoint (assuming single chain for simplicity). + // `findOne` requires a selection condition, so the newest row is taken + // with an ordered, limited `find` instead of a bare `findOne`. + const [checkpoint] = await this.checkpointRepo.find({ order: { updatedAt: 'DESC' }, + take: 1, }); if (!checkpoint) { @@ -52,40 +74,49 @@ export class DisputesQueryService { throw new BadRequestException('limit must be between 1 and 100'); } + // Fail closed: dispute state is protocol state, so it is only served + // while the projection provably reproduces canonical events. + await this.readiness.assertReady(V2_PROJECTORS.DISPUTES); + const decoded = cursor ? decodeCursor(cursor) : null; - - const query = this.disputeRepo.createQueryBuilder('dispute') + + const query = this.disputeRepo + .createQueryBuilder('dispute') .where('dispute.claimId = :claimId', { claimId }) .orderBy('dispute.blockNumber', 'ASC') .addOrderBy('dispute.eventLogIndex', 'ASC'); - + if (decoded) { query.andWhere( '(dispute.blockNumber > :blockNumber OR ' + - '(dispute.blockNumber = :blockNumber AND dispute.eventLogIndex > :logIndex))', - { blockNumber: decoded.blockNumber, logIndex: decoded.logIndex } + '(dispute.blockNumber = :blockNumber AND dispute.eventLogIndex > :logIndex))', + { blockNumber: decoded.blockNumber, logIndex: decoded.logIndex }, ); } - + const disputes = await query.limit(limit).getMany(); - + // Add computed data states const disputesWithState = await Promise.all( disputes.map(async (dispute) => ({ ...dispute, computedDataState: await this.calculateDataState(dispute.blockNumber), - })) + })), ); - + // Generate next cursor - const nextCursor = disputesWithState.length === limit - ? encodeCursor({ - blockNumber: disputesWithState[disputesWithState.length - 1].blockNumber, - logIndex: disputesWithState[disputesWithState.length - 1].eventLogIndex, - id: disputesWithState[disputesWithState.length - 1].disputeId, - }) - : null; - + const nextCursor = + disputesWithState.length === limit + ? encodeCursor({ + blockNumber: + disputesWithState[disputesWithState.length - 1].blockNumber ?? + '0', + logIndex: + disputesWithState[disputesWithState.length - 1].eventLogIndex, + id: disputesWithState[disputesWithState.length - 1].disputeId, + }) + : null; + return { items: disputesWithState, nextCursor, @@ -96,16 +127,21 @@ export class DisputesQueryService { claimId: string, originalRoundId: string, ): Promise { + // Fail closed: see listForClaim. + await this.readiness.assertReady(V2_PROJECTORS.DISPUTES); + const disputeId = `${claimId}:${originalRoundId}`; const dispute = await this.disputeRepo.findOne({ where: { disputeId } }); if (!dispute) throw new NotFoundException( `No dispute projected for round ${originalRoundId} on claim ${claimId}`, ); - const computedDataState = await this.calculateDataState(dispute.blockNumber); + const computedDataState = await this.calculateDataState( + dispute.blockNumber, + ); return { ...dispute, - computedDataState + computedDataState, }; } -} \ No newline at end of file +} diff --git a/src/v2/disputes/disputes.controller.ts b/src/v2/disputes/disputes.controller.ts index 48d1568d..eae47574 100644 --- a/src/v2/disputes/disputes.controller.ts +++ b/src/v2/disputes/disputes.controller.ts @@ -16,7 +16,11 @@ export class DisputesController { @Query('limit') limit?: string, @Query('cursor') cursor?: string, ) { - return this.queryService.listForClaim(claimId, limit ? parseInt(limit, 10) : 20, cursor); + return this.queryService.listForClaim( + claimId, + limit ? parseInt(limit, 10) : 20, + cursor, + ); } @Get(':originalRoundId') @@ -26,4 +30,4 @@ export class DisputesController { ) { return this.queryService.getByOriginalRound(claimId, originalRoundId); } -} \ No newline at end of file +} diff --git a/src/v2/disputes/entities/project-dispute.entity.ts b/src/v2/disputes/entities/project-dispute.entity.ts index 75e5c326..0fc9207f 100644 --- a/src/v2/disputes/entities/project-dispute.entity.ts +++ b/src/v2/disputes/entities/project-dispute.entity.ts @@ -76,9 +76,19 @@ export class ProjectDispute { @Column({ type: 'int' }) eventLogIndex: number; + /** + * Block containing the event this row was last derived from. Chain-native + * ordering coordinate, so keyset pagination and data-state labelling stay + * reproducible from canonical events. Null only for rows projected before + * this column existed whose originating block could not be resolved from + * the canonical stream; null is reported as OBSERVED, never as finalized. + */ + @Column({ type: 'bigint', nullable: true }) + blockNumber: string | null; + @CreateDateColumn() createdAt: Date; @UpdateDateColumn() updatedAt: Date; -} \ No newline at end of file +} diff --git a/src/v2/disputes/v2-disputes.module.ts b/src/v2/disputes/v2-disputes.module.ts index 1050f856..68c09244 100644 --- a/src/v2/disputes/v2-disputes.module.ts +++ b/src/v2/disputes/v2-disputes.module.ts @@ -5,6 +5,7 @@ import { EventCheckpoint } from '../events/entities/event-checkpoint.entity'; import { ProjectDispute } from './entities/project-dispute.entity'; import { ProjectorCursor } from '../common/entities/projector-cursor.entity'; import { IndexingAnomaly } from '../common/entities/indexing-anomaly.entity'; +import { ProjectionReadinessModule } from '../common/projection-readiness/projection-readiness.module'; import { DisputesProjectorService } from './disputes-projector.service'; import { DisputesQueryService } from './disputes-query.service'; import { DisputesController } from './disputes.controller'; @@ -18,9 +19,10 @@ import { DisputesController } from './disputes.controller'; EventCheckpoint, ]), V2EventsModule, + ProjectionReadinessModule, ], controllers: [DisputesController], providers: [DisputesProjectorService, DisputesQueryService], exports: [DisputesProjectorService, DisputesQueryService], }) -export class V2DisputesModule {} \ No newline at end of file +export class V2DisputesModule {} diff --git a/src/v2/evidence/evidence-projector.service.integration.spec.ts b/src/v2/evidence/evidence-projector.service.integration.spec.ts index c5ef9621..748e9d3d 100644 --- a/src/v2/evidence/evidence-projector.service.integration.spec.ts +++ b/src/v2/evidence/evidence-projector.service.integration.spec.ts @@ -1,5 +1,7 @@ import { Test, TestingModule } from '@nestjs/testing'; import { TypeOrmModule } from '@nestjs/typeorm'; +import { ConfigService } from '@nestjs/config'; +import { ServiceUnavailableException } from '@nestjs/common'; import { DataSource } from 'typeorm'; import { EvidenceProjectorService } from './evidence-projector.service'; import { EvidenceQueryService } from './evidence-query.service'; @@ -11,6 +13,9 @@ import { ProjectEvidenceVersion } from './entities/project-evidence-version.enti import { ProjectorCursor } from '../common/entities/projector-cursor.entity'; import { CanonicalEvent } from '../events/entities/canonical-event.entity'; import { CanonicalEventQueryService } from '../events/canonical-event-query.service'; +import { ContractArtifact } from '../events/entities/contract-artifact.entity'; +import { EventQuarantine } from '../events/entities/event-quarantine.entity'; +import { ProjectionReadinessService } from '../common/projection-readiness/projection-readiness.service'; describe('EvidenceProjectorService (integration)', () => { let moduleRef: TestingModule; @@ -36,6 +41,20 @@ describe('EvidenceProjectorService (integration)', () => { }); } + /** + * Move the projector's cursor behind the canonical head, which is exactly + * the state the gate exists to catch: canonical events are waiting while + * the read model still reflects an older point in the stream. + */ + async function rewindCursorTo(blockNumber: string): Promise { + await dataSource + .getRepository(ProjectorCursor) + .update( + { projectorName: 'v2-evidence' }, + { lastBlockNumber: blockNumber }, + ); + } + beforeEach(async () => { moduleRef = await Test.createTestingModule({ imports: [ @@ -49,6 +68,8 @@ describe('EvidenceProjectorService (integration)', () => { ProjectEvidence, ProjectEvidenceVersion, ProjectorCursor, + ContractArtifact, + EventQuarantine, ], synchronize: true, }), @@ -57,12 +78,16 @@ describe('EvidenceProjectorService (integration)', () => { ProjectEvidenceVersion, ProjectorCursor, CanonicalEvent, + ContractArtifact, + EventQuarantine, ]), ], providers: [ EvidenceProjectorService, EvidenceQueryService, CanonicalEventQueryService, + ProjectionReadinessService, + { provide: ConfigService, useValue: { get: () => undefined } }, ], }).compile(); @@ -205,4 +230,28 @@ describe('EvidenceProjectorService (integration)', () => { expect(secondPage.items).toHaveLength(1); expect(secondPage.nextCursor).toBeNull(); }); + + it('fails closed instead of serving stale evidence when the projection is behind canonical events', async () => { + await seedEvent({ + eventName: 'EvidenceRegistered', + txHash: '0x' + '01'.repeat(32), + blockNumber: '100', + logIndex: 0, + payload: { digest: '0xdigest1' }, + }); + await seedEvent({ + eventName: 'EvidenceReplaced', + txHash: '0x' + '02'.repeat(32), + blockNumber: '200', + logIndex: 0, + payload: { digest: '0xdigest2' }, + }); + + await projector.processNewEvents(); + await rewindCursorTo('100'); + + await expect(queryService.getEvidence(claimId)).rejects.toBeInstanceOf( + ServiceUnavailableException, + ); + }); }); diff --git a/src/v2/evidence/evidence-projector.service.ts b/src/v2/evidence/evidence-projector.service.ts index 81a4236b..d22fbf2a 100644 --- a/src/v2/evidence/evidence-projector.service.ts +++ b/src/v2/evidence/evidence-projector.service.ts @@ -9,13 +9,19 @@ import { } from './entities/project-evidence.entity'; import { ProjectEvidenceVersion } from './entities/project-evidence-version.entity'; import { ProjectorCursor } from '../common/entities/projector-cursor.entity'; +import { + PROJECTOR_HANDLED_EVENTS, + V2_PROJECTORS, + V2ProjectorName, +} from '../common/projection-readiness/projector-registry'; -const PROJECTOR_NAME = 'v2-evidence'; +const PROJECTOR_NAME: V2ProjectorName = V2_PROJECTORS.EVIDENCE; const PG_UNIQUE_VIOLATION = '23505'; -const HANDLED_EVENT_NAMES = [ - 'EvidenceRegistered', - 'EvidenceReplaced', - 'EvidenceRemoved', +// Name and handled-event list come from the projector registry so the +// readiness gate (V2-BE-100) can never disagree with this projector about +// which canonical events it is responsible for consuming. +const HANDLED_EVENT_NAMES: string[] = [ + ...PROJECTOR_HANDLED_EVENTS[PROJECTOR_NAME], ]; export interface ProjectorRunSummary { diff --git a/src/v2/evidence/evidence-query.service.ts b/src/v2/evidence/evidence-query.service.ts index 8d988a43..19b68a73 100644 --- a/src/v2/evidence/evidence-query.service.ts +++ b/src/v2/evidence/evidence-query.service.ts @@ -9,6 +9,8 @@ import { decodeCursor, encodeCursor, } from '../common/cursor-pagination'; +import { ProjectionReadinessService } from '../common/projection-readiness/projection-readiness.service'; +import { V2_PROJECTORS } from '../common/projection-readiness/projector-registry'; @Injectable() export class EvidenceQueryService { @@ -17,9 +19,14 @@ export class EvidenceQueryService { private readonly evidenceRepo: Repository, @InjectRepository(ProjectEvidenceVersion) private readonly versionRepo: Repository, + private readonly readiness: ProjectionReadinessService, ) {} async getEvidence(claimId: string): Promise { + // Fail closed: evidence state is protocol state, so it is only served + // while the projection provably reproduces canonical events. + await this.readiness.assertReady(V2_PROJECTORS.EVIDENCE); + const evidence = await this.evidenceRepo.findOne({ where: { claimId } }); if (!evidence) throw new NotFoundException(`No evidence projected for claim ${claimId}`); @@ -31,6 +38,9 @@ export class EvidenceQueryService { cursor?: string, limit?: number, ): Promise> { + // Fail closed: see getEvidence. + await this.readiness.assertReady(V2_PROJECTORS.EVIDENCE); + const pageSize = clampPageSize(limit); const qb = this.versionRepo .createQueryBuilder('v') diff --git a/src/v2/evidence/v2-evidence.module.ts b/src/v2/evidence/v2-evidence.module.ts index e01d7fd5..d9e07ee4 100644 --- a/src/v2/evidence/v2-evidence.module.ts +++ b/src/v2/evidence/v2-evidence.module.ts @@ -4,6 +4,7 @@ import { V2EventsModule } from '../events/v2-events.module'; import { ProjectEvidence } from './entities/project-evidence.entity'; import { ProjectEvidenceVersion } from './entities/project-evidence-version.entity'; import { ProjectorCursor } from '../common/entities/projector-cursor.entity'; +import { ProjectionReadinessModule } from '../common/projection-readiness/projection-readiness.module'; import { EvidenceProjectorService } from './evidence-projector.service'; import { EvidenceQueryService } from './evidence-query.service'; import { EvidenceController } from './evidence.controller'; @@ -16,6 +17,7 @@ import { EvidenceController } from './evidence.controller'; ProjectorCursor, ]), V2EventsModule, + ProjectionReadinessModule, ], controllers: [EvidenceController], providers: [EvidenceProjectorService, EvidenceQueryService], diff --git a/src/v2/verification/entities/project-participant-position.entity.ts b/src/v2/verification/entities/project-participant-position.entity.ts index 713fb944..5aade6f6 100644 --- a/src/v2/verification/entities/project-participant-position.entity.ts +++ b/src/v2/verification/entities/project-participant-position.entity.ts @@ -57,4 +57,4 @@ export class ProjectParticipantPosition { @CreateDateColumn() createdAt: Date; -} \ No newline at end of file +} diff --git a/src/v2/verification/entities/project-verification-round.entity.ts b/src/v2/verification/entities/project-verification-round.entity.ts index 21e33b0c..8a82075c 100644 --- a/src/v2/verification/entities/project-verification-round.entity.ts +++ b/src/v2/verification/entities/project-verification-round.entity.ts @@ -69,7 +69,7 @@ export class ProjectVerificationRound { roundSnapshot: Record | null; /** Appeal deadline if this is an appeal round */ - @Column({ type: 'Date', nullable: true }) + @Column({ type: Date, nullable: true }) appealDeadline: Date | null; @Column({ type: 'varchar', length: 66 }) @@ -83,4 +83,4 @@ export class ProjectVerificationRound { @UpdateDateColumn() updatedAt: Date; -} \ No newline at end of file +} diff --git a/src/v2/verification/v2-verification.module.ts b/src/v2/verification/v2-verification.module.ts index 9a225cba..9dddcce1 100644 --- a/src/v2/verification/v2-verification.module.ts +++ b/src/v2/verification/v2-verification.module.ts @@ -6,6 +6,7 @@ import { ProjectVerificationRound } from './entities/project-verification-round. import { ProjectParticipantPosition } from './entities/project-participant-position.entity'; import { ProjectorCursor } from '../common/entities/projector-cursor.entity'; import { IndexingAnomaly } from '../common/entities/indexing-anomaly.entity'; +import { ProjectionReadinessModule } from '../common/projection-readiness/projection-readiness.module'; import { VerificationProjectorService } from './verification-projector.service'; import { VerificationQueryService } from './verification-query.service'; import { VerificationController } from './verification.controller'; @@ -20,9 +21,10 @@ import { VerificationController } from './verification.controller'; EventCheckpoint, ]), V2EventsModule, + ProjectionReadinessModule, ], controllers: [VerificationController], providers: [VerificationProjectorService, VerificationQueryService], exports: [VerificationProjectorService, VerificationQueryService], }) -export class V2VerificationModule {} \ No newline at end of file +export class V2VerificationModule {} diff --git a/src/v2/verification/verification-projector.service.integration.spec.ts b/src/v2/verification/verification-projector.service.integration.spec.ts index f741a7d6..6be3213b 100644 --- a/src/v2/verification/verification-projector.service.integration.spec.ts +++ b/src/v2/verification/verification-projector.service.integration.spec.ts @@ -1,5 +1,7 @@ import { Test, TestingModule } from '@nestjs/testing'; import { TypeOrmModule } from '@nestjs/typeorm'; +import { ConfigService } from '@nestjs/config'; +import { ServiceUnavailableException } from '@nestjs/common'; import { DataSource } from 'typeorm'; import { VerificationProjectorService } from './verification-projector.service'; import { VerificationQueryService } from './verification-query.service'; @@ -15,6 +17,10 @@ import { } from '../common/entities/indexing-anomaly.entity'; import { CanonicalEvent } from '../events/entities/canonical-event.entity'; import { CanonicalEventQueryService } from '../events/canonical-event-query.service'; +import { EventCheckpoint } from '../events/entities/event-checkpoint.entity'; +import { ContractArtifact } from '../events/entities/contract-artifact.entity'; +import { EventQuarantine } from '../events/entities/event-quarantine.entity'; +import { ProjectionReadinessService } from '../common/projection-readiness/projection-readiness.service'; describe('VerificationProjectorService (integration)', () => { let moduleRef: TestingModule; @@ -55,6 +61,9 @@ describe('VerificationProjectorService (integration)', () => { ProjectParticipantPosition, ProjectorCursor, IndexingAnomaly, + EventCheckpoint, + ContractArtifact, + EventQuarantine, ], synchronize: true, }), @@ -64,12 +73,17 @@ describe('VerificationProjectorService (integration)', () => { ProjectorCursor, IndexingAnomaly, CanonicalEvent, + EventCheckpoint, + ContractArtifact, + EventQuarantine, ]), ], providers: [ VerificationProjectorService, VerificationQueryService, CanonicalEventQueryService, + ProjectionReadinessService, + { provide: ConfigService, useValue: { get: () => undefined } }, ], }).compile(); @@ -100,7 +114,10 @@ describe('VerificationProjectorService (integration)', () => { await projector.processNewEvents(); - const { first, appeal } = await queryService.listRounds(claimId); + const { firstInstanceRounds, appealRounds } = + await queryService.listRounds(claimId); + const first = firstInstanceRounds.items; + const appeal = appealRounds.items; expect(first).toHaveLength(1); expect(appeal).toHaveLength(1); expect(first[0].roundId).toBe(firstRoundId); @@ -132,11 +149,11 @@ describe('VerificationProjectorService (integration)', () => { await projector.processNewEvents(); - const positions = await queryService.listPositions(firstRoundId); - expect(positions).toHaveLength(1); - expect(positions[0].stake).toBe('1000000000000000000'); - expect(positions[0].effectiveWeight).toBe('850'); - expect(positions[0].position).toBe('support'); + const positionPage = await queryService.listPositions(firstRoundId); + expect(positionPage.items).toHaveLength(1); + expect(positionPage.items[0].stake).toBe('1000000000000000000'); + expect(positionPage.items[0].effectiveWeight).toBe('850'); + expect(positionPage.items[0].position).toBe('support'); }); it('detects and records a duplicate position for the same participant/round instead of overwriting it', async () => { @@ -167,9 +184,9 @@ describe('VerificationProjectorService (integration)', () => { const summary = await projector.processNewEvents(); expect(summary.anomalies).toBe(1); - const positions = await queryService.listPositions(firstRoundId); - expect(positions).toHaveLength(1); - expect(positions[0].stake).toBe('100'); // first-committed position wins, not overwritten + const positionPage = await queryService.listPositions(firstRoundId); + expect(positionPage.items).toHaveLength(1); + expect(positionPage.items[0].stake).toBe('100'); // first-committed position wins, not overwritten const anomalies = await dataSource.getRepository(IndexingAnomaly).find(); expect(anomalies).toHaveLength(1); @@ -216,4 +233,30 @@ describe('VerificationProjectorService (integration)', () => { .find(); expect(rounds).toHaveLength(1); }); + + it('fails closed instead of serving a stale round projection when canonical events are unprojected', async () => { + await seedEvent({ + eventName: 'VerificationRoundOpened', + txHash: '0x' + '01'.repeat(32), + blockNumber: '100', + roundId: firstRoundId, + payload: { roundType: 'first', roundNumber: '1' }, + }); + await seedEvent({ + eventName: 'VerificationRoundOpened', + txHash: '0x' + '02'.repeat(32), + blockNumber: '300', + roundId: appealRoundId, + payload: { roundType: 'appeal', roundNumber: '1' }, + }); + + await projector.processNewEvents(); + await dataSource + .getRepository(ProjectorCursor) + .update({ projectorName: 'v2-verification' }, { lastBlockNumber: '100' }); + + await expect(queryService.listRounds(claimId)).rejects.toBeInstanceOf( + ServiceUnavailableException, + ); + }); }); diff --git a/src/v2/verification/verification-projector.service.ts b/src/v2/verification/verification-projector.service.ts index ecada6e3..adf6f2f2 100644 --- a/src/v2/verification/verification-projector.service.ts +++ b/src/v2/verification/verification-projector.service.ts @@ -14,10 +14,20 @@ import { IndexingAnomaly, IndexingAnomalyKind, } from '../common/entities/indexing-anomaly.entity'; +import { + PROJECTOR_HANDLED_EVENTS, + V2_PROJECTORS, + V2ProjectorName, +} from '../common/projection-readiness/projector-registry'; -const PROJECTOR_NAME = 'v2-verification'; +const PROJECTOR_NAME: V2ProjectorName = V2_PROJECTORS.VERIFICATION; const PG_UNIQUE_VIOLATION = '23505'; -const HANDLED_EVENT_NAMES = ['VerificationRoundOpened', 'PositionCommitted']; +// Name and handled-event list come from the projector registry so the +// readiness gate (V2-BE-100) can never disagree with this projector about +// which canonical events it is responsible for consuming. +const HANDLED_EVENT_NAMES: string[] = [ + ...PROJECTOR_HANDLED_EVENTS[PROJECTOR_NAME], +]; export interface ProjectorRunSummary { processed: number; diff --git a/src/v2/verification/verification-query.service.ts b/src/v2/verification/verification-query.service.ts index ead2262b..7f5cf75e 100644 --- a/src/v2/verification/verification-query.service.ts +++ b/src/v2/verification/verification-query.service.ts @@ -1,4 +1,8 @@ -import { Injectable, NotFoundException, BadRequestException } from '@nestjs/common'; +import { + Injectable, + NotFoundException, + BadRequestException, +} from '@nestjs/common'; import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import { @@ -8,8 +12,13 @@ import { import { ProjectParticipantPosition } from './entities/project-participant-position.entity'; import { EventCheckpoint } from '../events/entities/event-checkpoint.entity'; import { DataState } from '../common/data-state.enum'; -import { CursorPage, encodeCursor, decodeCursor } from '../common/cursor-pagination'; - +import { + CursorPage, + encodeCursor, + decodeCursor, +} from '../common/cursor-pagination'; +import { ProjectionReadinessService } from '../common/projection-readiness/projection-readiness.service'; +import { V2_PROJECTORS } from '../common/projection-readiness/projector-registry'; @Injectable() export class VerificationQueryService { @@ -20,15 +29,19 @@ export class VerificationQueryService { private readonly positionRepo: Repository, @InjectRepository(EventCheckpoint) private readonly checkpointRepo: Repository, + private readonly readiness: ProjectionReadinessService, ) {} /** * Calculate the data state for a block number based on chain's safe and finalized blocks */ private async calculateDataState(blockNumber: string): Promise { - // Get the latest checkpoint (assuming single chain for simplicity) - const checkpoint = await this.checkpointRepo.findOne({ + // Get the latest checkpoint (assuming single chain for simplicity). + // `findOne` requires a selection condition, so the newest row is taken + // with an ordered, limited `find` instead of a bare `findOne`. + const [checkpoint] = await this.checkpointRepo.find({ order: { updatedAt: 'DESC' }, + take: 1, }); if (!checkpoint) { @@ -53,8 +66,12 @@ export class VerificationQueryService { limit = 20, cursor?: string, ): Promise<{ - firstInstanceRounds: CursorPage; - appealRounds: CursorPage; + firstInstanceRounds: CursorPage< + ProjectVerificationRound & { computedDataState: DataState } + >; + appealRounds: CursorPage< + ProjectVerificationRound & { computedDataState: DataState } + >; }> { if (!claimId) { throw new BadRequestException('claimId is required'); @@ -63,75 +80,91 @@ export class VerificationQueryService { throw new BadRequestException('limit must be between 1 and 100'); } + // Fail closed: round state is protocol state, so it is only served while + // the projection provably reproduces canonical events. + await this.readiness.assertReady(V2_PROJECTORS.VERIFICATION); + const decoded = cursor ? decodeCursor(cursor) : null; - + // Get first instance rounds with pagination - const firstRoundQuery = this.roundRepo.createQueryBuilder('round') + const firstRoundQuery = this.roundRepo + .createQueryBuilder('round') .where('round.claimId = :claimId', { claimId }) .andWhere('round.roundType = :type', { type: RoundType.FIRST }) .orderBy('round.openedAtBlock', 'ASC') .addOrderBy('round.eventLogIndex', 'ASC'); - + if (decoded) { firstRoundQuery.andWhere( '(round.openedAtBlock > :blockNumber OR ' + - '(round.openedAtBlock = :blockNumber AND round.eventLogIndex > :logIndex))', - { blockNumber: decoded.blockNumber, logIndex: decoded.logIndex } + '(round.openedAtBlock = :blockNumber AND round.eventLogIndex > :logIndex))', + { blockNumber: decoded.blockNumber, logIndex: decoded.logIndex }, ); } - + const firstRounds = await firstRoundQuery.limit(limit).getMany(); - + // Calculate data states for first rounds const firstRoundsWithState = await Promise.all( firstRounds.map(async (round) => ({ ...round, computedDataState: await this.calculateDataState(round.openedAtBlock), - })) + })), ); - + // Get appeal rounds - const appealRoundQuery = this.roundRepo.createQueryBuilder('round') + const appealRoundQuery = this.roundRepo + .createQueryBuilder('round') .where('round.claimId = :claimId', { claimId }) .andWhere('round.roundType = :type', { type: RoundType.APPEAL }) .orderBy('round.openedAtBlock', 'ASC') .addOrderBy('round.eventLogIndex', 'ASC'); - + if (decoded) { appealRoundQuery.andWhere( '(round.openedAtBlock > :blockNumber OR ' + - '(round.openedAtBlock = :blockNumber AND round.eventLogIndex > :logIndex))', - { blockNumber: decoded.blockNumber, logIndex: decoded.logIndex } + '(round.openedAtBlock = :blockNumber AND round.eventLogIndex > :logIndex))', + { blockNumber: decoded.blockNumber, logIndex: decoded.logIndex }, ); } - + const appealRounds = await appealRoundQuery.limit(limit).getMany(); - + // Calculate data states for appeal rounds const appealRoundsWithState = await Promise.all( appealRounds.map(async (round) => ({ ...round, computedDataState: await this.calculateDataState(round.openedAtBlock), - })) + })), ); - + // Generate next cursors - const firstNextCursor = firstRoundsWithState.length === limit - ? encodeCursor({ - blockNumber: firstRoundsWithState[firstRoundsWithState.length - 1].openedAtBlock, - logIndex: firstRoundsWithState[firstRoundsWithState.length - 1].eventLogIndex, - id: firstRoundsWithState[firstRoundsWithState.length - 1].roundId, - }) - : null; - - const appealNextCursor = appealRoundsWithState.length === limit - ? encodeCursor({ - blockNumber: appealRoundsWithState[appealRoundsWithState.length - 1].openedAtBlock, - logIndex: appealRoundsWithState[appealRoundsWithState.length - 1].eventLogIndex, - id: appealRoundsWithState[appealRoundsWithState.length - 1].roundId, - }) - : null; - + const firstNextCursor = + firstRoundsWithState.length === limit + ? encodeCursor({ + blockNumber: + firstRoundsWithState[firstRoundsWithState.length - 1] + .openedAtBlock, + logIndex: + firstRoundsWithState[firstRoundsWithState.length - 1] + .eventLogIndex, + id: firstRoundsWithState[firstRoundsWithState.length - 1].roundId, + }) + : null; + + const appealNextCursor = + appealRoundsWithState.length === limit + ? encodeCursor({ + blockNumber: + appealRoundsWithState[appealRoundsWithState.length - 1] + .openedAtBlock, + logIndex: + appealRoundsWithState[appealRoundsWithState.length - 1] + .eventLogIndex, + id: appealRoundsWithState[appealRoundsWithState.length - 1].roundId, + }) + : null; + return { firstInstanceRounds: { items: firstRoundsWithState, @@ -144,17 +177,26 @@ export class VerificationQueryService { }; } - async getRound(roundId: string): Promise { + async getRound( + roundId: string, + ): Promise { if (!roundId) { throw new BadRequestException('roundId is required'); } - + + // Fail closed: see listRounds. + await this.readiness.assertReady(V2_PROJECTORS.VERIFICATION); + const round = await this.roundRepo.findOne({ where: { roundId } }); if (!round) { - throw new NotFoundException(`No verification round projected for id ${roundId}`); + throw new NotFoundException( + `No verification round projected for id ${roundId}`, + ); } - - const computedDataState = await this.calculateDataState(round.openedAtBlock); + + const computedDataState = await this.calculateDataState( + round.openedAtBlock, + ); return { ...round, computedDataState, @@ -166,7 +208,9 @@ export class VerificationQueryService { roundId: string, limit = 20, cursor?: string, - ): Promise> { + ): Promise< + CursorPage + > { if (!roundId) { throw new BadRequestException('roundId is required'); } @@ -174,43 +218,50 @@ export class VerificationQueryService { throw new BadRequestException('limit must be between 1 and 100'); } + // Fail closed: see listRounds. + await this.readiness.assertReady(V2_PROJECTORS.VERIFICATION); + const decoded = cursor ? decodeCursor(cursor) : null; - - const query = this.positionRepo.createQueryBuilder('position') + + const query = this.positionRepo + .createQueryBuilder('position') .where('position.roundId = :roundId', { roundId }) .orderBy('position.blockNumber', 'ASC') .addOrderBy('position.eventLogIndex', 'ASC'); - + if (decoded) { query.andWhere( '(position.blockNumber > :blockNumber OR ' + - '(position.blockNumber = :blockNumber AND position.eventLogIndex > :logIndex))', - { blockNumber: decoded.blockNumber, logIndex: decoded.logIndex } + '(position.blockNumber = :blockNumber AND position.eventLogIndex > :logIndex))', + { blockNumber: decoded.blockNumber, logIndex: decoded.logIndex }, ); } - + const positions = await query.limit(limit).getMany(); - + // Add computed data states const positionsWithState = await Promise.all( positions.map(async (pos) => ({ ...pos, computedDataState: await this.calculateDataState(pos.blockNumber), - })) + })), ); - + // Generate next cursor - const nextCursor = positionsWithState.length === limit - ? encodeCursor({ - blockNumber: positionsWithState[positionsWithState.length - 1].blockNumber, - logIndex: positionsWithState[positionsWithState.length - 1].eventLogIndex, - id: positionsWithState[positionsWithState.length - 1].id, - }) - : null; - + const nextCursor = + positionsWithState.length === limit + ? encodeCursor({ + blockNumber: + positionsWithState[positionsWithState.length - 1].blockNumber, + logIndex: + positionsWithState[positionsWithState.length - 1].eventLogIndex, + id: positionsWithState[positionsWithState.length - 1].id, + }) + : null; + return { items: positionsWithState, nextCursor, }; } -} \ No newline at end of file +} diff --git a/src/v2/verification/verification.controller.ts b/src/v2/verification/verification.controller.ts index 3d719f82..0fa36810 100644 --- a/src/v2/verification/verification.controller.ts +++ b/src/v2/verification/verification.controller.ts @@ -1,4 +1,4 @@ -import { Controller, Get, Param, Query, ParseUUIDPipe } from '@nestjs/common'; +import { Controller, Get, Param, Query } from '@nestjs/common'; import { VerificationQueryService } from './verification-query.service'; /** @@ -16,7 +16,11 @@ export class VerificationController { @Query('limit') limit?: string, @Query('cursor') cursor?: string, ) { - return this.queryService.listRounds(claimId, limit ? parseInt(limit, 10) : 20, cursor); + return this.queryService.listRounds( + claimId, + limit ? parseInt(limit, 10) : 20, + cursor, + ); } @Get('verification-rounds/:roundId') @@ -30,6 +34,10 @@ export class VerificationController { @Query('limit') limit?: string, @Query('cursor') cursor?: string, ) { - return this.queryService.listPositions(roundId, limit ? parseInt(limit, 10) : 20, cursor); + return this.queryService.listPositions( + roundId, + limit ? parseInt(limit, 10) : 20, + cursor, + ); } -} \ No newline at end of file +}