Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 8 additions & 2 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -5,8 +5,11 @@ on:
branches: ["main", "master"]
pull_request:

permissions:
contents: read

jobs:
build:
verify:
runs-on: ubuntu-latest
steps:
- name: Checkout
Expand All @@ -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
4 changes: 2 additions & 2 deletions src/dashboard/dashboard.service.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -124,7 +124,8 @@ describe("DashboardService", () => {
source: { code: "eis" }
}
])
}
},
$transaction: vi.fn(async (operations: Array<Promise<unknown>>) => Promise.all(operations))
};

const configService = {
Expand All @@ -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);

Expand Down
101 changes: 100 additions & 1 deletion src/sources/sources.controller.ts
Original file line number Diff line number Diff line change
@@ -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";
Expand Down Expand Up @@ -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<string, unknown>;
}

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<string, unknown>;

@IsArray()
@ValidateNested({ each: true })
@Type(() => QuarantineArtifactDto)
artifacts!: QuarantineArtifactDto[];
}

@Controller("internal/scraper")
export class SourcesController {
constructor(
Expand Down Expand Up @@ -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 };
}
}
166 changes: 163 additions & 3 deletions src/sources/sources.service.ts
Original file line number Diff line number Diff line change
@@ -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";

Expand All @@ -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<string, unknown>;
artifacts: Array<{
kind: string;
bucket: string;
objectKey: string;
mimeType?: string;
checksum?: string;
sizeBytes?: number;
metadata?: Record<string, unknown>;
}>;
};

@Injectable()
export class SourcesService {
constructor(
Expand Down Expand Up @@ -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,
Expand All @@ -163,12 +195,140 @@ 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,
itemsDiscovered: input.itemsDiscovered ?? 0
}
});
}

async quarantineRawEvent(input: QuarantineRawEventInput) {
await syncEnabledSourcesCatalog(
this.prisma,
this.configService.get<string[]>("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;
});
}
}
Loading