From 18c6e29f51b73698bdfa0a32da0cfaaf4a0326c8 Mon Sep 17 00:00:00 2001 From: Aleksandr Minkin Date: Tue, 7 Apr 2026 21:43:45 +0300 Subject: [PATCH 1/2] Record quarantined scraper events --- src/sources/sources.controller.ts | 101 +++++++++++++++++- src/sources/sources.service.ts | 166 +++++++++++++++++++++++++++++- 2 files changed, 263 insertions(+), 4 deletions(-) diff --git a/src/sources/sources.controller.ts b/src/sources/sources.controller.ts index 72539c0..fae7873 100644 --- a/src/sources/sources.controller.ts +++ b/src/sources/sources.controller.ts @@ -1,7 +1,18 @@ import { Body, Controller, HttpCode, Post, Req } from "@nestjs/common"; import { SourceRunStatus } from "@prisma/client"; import type { Request } from "express"; -import { IsDateString, IsEnum, IsInt, IsOptional, IsString, Min } from "class-validator"; +import { + IsArray, + IsDateString, + IsEnum, + IsInt, + IsObject, + IsOptional, + IsString, + Min, + ValidateNested +} from "class-validator"; +import { Type } from "class-transformer"; import { AuthService } from "../auth/auth.service"; import type { RequestLike } from "../common/request-context"; import { SourcesService } from "./sources.service"; @@ -33,6 +44,69 @@ class UpsertSourceRunDto { itemsDiscovered?: number; } +class QuarantineArtifactDto { + @IsString() + kind!: string; + + @IsString() + bucket!: string; + + @IsString() + objectKey!: string; + + @IsOptional() + @IsString() + mimeType?: string; + + @IsOptional() + @IsString() + checksum?: string; + + @IsOptional() + @IsInt() + @Min(0) + sizeBytes?: number; + + @IsOptional() + @IsObject() + metadata?: Record; +} + +class QuarantineRawEventDto { + @IsString() + source!: string; + + @IsString() + eventId!: string; + + @IsString() + runKey!: string; + + @IsDateString() + collectedAt!: string; + + @IsString() + sourceUrl!: string; + + @IsString() + payloadVersion!: string; + + @IsOptional() + @IsString() + externalId?: string; + + @IsString() + quarantineReason!: string; + + @IsObject() + rawPayload!: Record; + + @IsArray() + @ValidateNested({ each: true }) + @Type(() => QuarantineArtifactDto) + artifacts!: QuarantineArtifactDto[]; +} + @Controller("internal/scraper") export class SourcesController { constructor( @@ -61,4 +135,29 @@ export class SourcesController { return { ok: true }; } + + @Post("quarantine-events") + @HttpCode(200) + async quarantineRawEvent(@Body() body: QuarantineRawEventDto, @Req() request: Request) { + this.authService.assertIngestToken({ + headers: request.headers as RequestLike["headers"], + id: typeof request.id === "string" ? request.id : undefined, + ip: request.ip + }); + + await this.sourcesService.quarantineRawEvent({ + sourceCode: body.source, + eventId: body.eventId, + runKey: body.runKey, + collectedAt: new Date(body.collectedAt), + sourceUrl: body.sourceUrl, + payloadVersion: body.payloadVersion, + externalId: body.externalId, + quarantineReason: body.quarantineReason, + rawPayload: body.rawPayload, + artifacts: body.artifacts + }); + + return { ok: true }; + } } diff --git a/src/sources/sources.service.ts b/src/sources/sources.service.ts index e94ffcf..2035e3d 100644 --- a/src/sources/sources.service.ts +++ b/src/sources/sources.service.ts @@ -1,6 +1,7 @@ import { BadGatewayException, Injectable } from "@nestjs/common"; import { ConfigService } from "@nestjs/config"; -import { SourceRunStatus } from "@prisma/client"; +import { ArtifactKind, RawEventStatus, SourceRunStatus } from "@prisma/client"; +import { toJson, toNullableJson } from "../prisma/json"; import { PrismaService } from "../prisma/prisma.service"; import { syncEnabledSourcesCatalog } from "./source-catalog"; @@ -14,6 +15,27 @@ type UpsertSourceRunInput = { itemsDiscovered?: number; }; +type QuarantineRawEventInput = { + sourceCode: string; + eventId: string; + runKey: string; + collectedAt: Date; + sourceUrl: string; + payloadVersion: string; + externalId?: string; + quarantineReason: string; + rawPayload: Record; + artifacts: Array<{ + kind: string; + bucket: string; + objectKey: string; + mimeType?: string; + checksum?: string; + sizeBytes?: number; + metadata?: Record; + }>; +}; + @Injectable() export class SourcesService { constructor( @@ -150,10 +172,20 @@ export class SourcesService { return null; } + const existingRun = await this.prisma.sourceRun.findUnique({ + where: { runKey: input.runKey }, + select: { status: true } + }); + const nextStatus = + existingRun && input.status === SourceRunStatus.SUCCESS && + (existingRun.status === SourceRunStatus.PARTIAL || existingRun.status === SourceRunStatus.FAILED) + ? existingRun.status + : input.status; + return this.prisma.sourceRun.upsert({ where: { runKey: input.runKey }, update: { - status: input.status, + status: nextStatus, startedAt: input.startedAt, finishedAt: input.finishedAt ?? null, errorMessage: input.errorMessage ?? null, @@ -163,7 +195,7 @@ export class SourcesService { create: { sourceId: source.id, runKey: input.runKey, - status: input.status, + status: nextStatus, startedAt: input.startedAt, finishedAt: input.finishedAt ?? null, errorMessage: input.errorMessage ?? null, @@ -171,4 +203,132 @@ export class SourcesService { } }); } + + async quarantineRawEvent(input: QuarantineRawEventInput) { + await syncEnabledSourcesCatalog( + this.prisma, + this.configService.get("ENABLED_SOURCES") ?? [] + ); + + const source = await this.prisma.source.findFirst({ + where: { + code: input.sourceCode, + deletedAt: null + }, + select: { id: true } + }); + + if (!source) { + return null; + } + + return this.prisma.$transaction(async (tx) => { + const existingRawEvent = await tx.rawEvent.findUnique({ + where: { eventId: input.eventId }, + select: { id: true, status: true, sourceRunId: true } + }); + + if (existingRawEvent?.status === RawEventStatus.NORMALIZED) { + return existingRawEvent; + } + + let sourceRun = await tx.sourceRun.findUnique({ + where: { runKey: input.runKey }, + select: { id: true, status: true } + }); + + if (!sourceRun) { + sourceRun = await tx.sourceRun.create({ + data: { + sourceId: source.id, + runKey: input.runKey, + status: SourceRunStatus.PARTIAL, + startedAt: input.collectedAt, + finishedAt: input.collectedAt, + itemsFailed: 1, + errorMessage: input.quarantineReason + }, + select: { id: true, status: true } + }); + } else if (!existingRawEvent) { + const shouldStayFailed = sourceRun.status === SourceRunStatus.FAILED; + sourceRun = await tx.sourceRun.update({ + where: { runKey: input.runKey }, + data: { + status: shouldStayFailed ? SourceRunStatus.FAILED : SourceRunStatus.PARTIAL, + finishedAt: new Date(), + itemsFailed: { increment: 1 }, + errorMessage: input.quarantineReason + }, + select: { id: true, status: true } + }); + } + + const rawEvent = await tx.rawEvent.upsert({ + where: { eventId: input.eventId }, + update: { + sourceRunId: sourceRun.id, + externalId: input.externalId, + payloadVersion: input.payloadVersion, + collectedAt: input.collectedAt, + sourceUrl: input.sourceUrl, + rawPayload: toJson(input.rawPayload) ?? {}, + status: RawEventStatus.QUARANTINED, + quarantineReason: input.quarantineReason + }, + create: { + sourceId: source.id, + sourceRunId: sourceRun.id, + eventId: input.eventId, + externalId: input.externalId, + payloadVersion: input.payloadVersion, + collectedAt: input.collectedAt, + sourceUrl: input.sourceUrl, + rawPayload: toJson(input.rawPayload) ?? {}, + status: RawEventStatus.QUARANTINED, + quarantineReason: input.quarantineReason, + checksum: input.eventId + } + }); + + for (const artifact of input.artifacts) { + await tx.artifact.upsert({ + where: { + bucket_objectKey: { + bucket: artifact.bucket, + objectKey: artifact.objectKey + } + }, + update: { + rawEventId: rawEvent.id, + sourceRunId: sourceRun.id, + mimeType: artifact.mimeType, + checksum: artifact.checksum, + sizeBytes: artifact.sizeBytes, + metadata: toNullableJson(artifact.metadata) + }, + create: { + rawEventId: rawEvent.id, + sourceRunId: sourceRun.id, + bucket: artifact.bucket, + objectKey: artifact.objectKey, + kind: + artifact.kind === "RAW_HTML" + ? ArtifactKind.RAW_HTML + : artifact.kind === "REPORT_FILE" + ? ArtifactKind.REPORT_FILE + : artifact.kind === "RAW_JSON" + ? ArtifactKind.RAW_JSON + : ArtifactKind.OTHER, + mimeType: artifact.mimeType, + checksum: artifact.checksum, + sizeBytes: artifact.sizeBytes, + metadata: toNullableJson(artifact.metadata) + } + }); + } + + return rawEvent; + }); + } } From 601d1bb9327ce6332fd104077ead0ae9e774a2c8 Mon Sep 17 00:00:00 2001 From: Aleksandr Minkin Date: Wed, 8 Apr 2026 10:42:32 +0300 Subject: [PATCH 2/2] Run backend tests in CI --- .github/workflows/ci.yml | 10 ++++++++-- src/dashboard/dashboard.service.spec.ts | 4 ++-- 2 files changed, 10 insertions(+), 4 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 31dbd1b..ed13605 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -5,8 +5,11 @@ on: branches: ["main", "master"] pull_request: +permissions: + contents: read + jobs: - build: + verify: runs-on: ubuntu-latest steps: - name: Checkout @@ -19,10 +22,13 @@ jobs: cache: npm - name: Install dependencies - run: npm install + run: npm ci - name: Type check run: npm run check + - name: Run tests + run: npm test + - name: Build run: npm run build diff --git a/src/dashboard/dashboard.service.spec.ts b/src/dashboard/dashboard.service.spec.ts index a8f81e4..8424f49 100644 --- a/src/dashboard/dashboard.service.spec.ts +++ b/src/dashboard/dashboard.service.spec.ts @@ -124,7 +124,8 @@ describe("DashboardService", () => { source: { code: "eis" } } ]) - } + }, + $transaction: vi.fn(async (operations: Array>) => Promise.all(operations)) }; const configService = { @@ -133,7 +134,6 @@ describe("DashboardService", () => { prisma.source.upsert.mockResolvedValue(undefined); prisma.source.updateMany.mockResolvedValue({ count: 0 }); - prisma.$transaction = vi.fn().mockImplementation(async (operations) => Promise.all(operations)); const service = new DashboardService(prisma as never, configService as never);