mirror of
https://github.com/immich-app/data.immich.app.git
synced 2026-09-30 13:23:26 +08:00
feat: github data pipeline (#20)
This commit is contained in:
Generated
+14
@@ -12,11 +12,13 @@
|
||||
"@influxdata/influxdb-client": "^1.34.0",
|
||||
"fetch-retry": "^6.0.0",
|
||||
"fflate": "^0.8.2",
|
||||
"itty-router": "^5.0.18",
|
||||
"p-limit": "^6.1.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@cloudflare/vitest-pool-workers": "^0.5.8",
|
||||
"@cloudflare/workers-types": "^4.20240729.0",
|
||||
"@octokit/webhooks-types": "^7.6.1",
|
||||
"@types/express": "^4.17.21",
|
||||
"@typescript-eslint/eslint-plugin": "^8.0.0",
|
||||
"@typescript-eslint/parser": "^8.0.0",
|
||||
@@ -2161,6 +2163,12 @@
|
||||
"node": ">= 8"
|
||||
}
|
||||
},
|
||||
"node_modules/@octokit/webhooks-types": {
|
||||
"version": "7.6.1",
|
||||
"resolved": "https://registry.npmjs.org/@octokit/webhooks-types/-/webhooks-types-7.6.1.tgz",
|
||||
"integrity": "sha512-S8u2cJzklBC0FgTwWVLaM8tMrDuDMVE4xiTK4EYXM9GntyvrdbSoxqDQa+Fh57CCNApyIpyeqPhhFEmHPfrXgw==",
|
||||
"dev": true
|
||||
},
|
||||
"node_modules/@pkgjs/parseargs": {
|
||||
"version": "0.11.0",
|
||||
"resolved": "https://registry.npmjs.org/@pkgjs/parseargs/-/parseargs-0.11.0.tgz",
|
||||
@@ -6543,6 +6551,12 @@
|
||||
"node": ">=8"
|
||||
}
|
||||
},
|
||||
"node_modules/itty-router": {
|
||||
"version": "5.0.18",
|
||||
"resolved": "https://registry.npmjs.org/itty-router/-/itty-router-5.0.18.tgz",
|
||||
"integrity": "sha512-mK3ReOt4ARAGy0V0J7uHmArG2USN2x0zprZ+u+YgmeRjXTDbaowDy3kPcsmQY6tH+uHhDgpWit9Vqmv/4rTXwA==",
|
||||
"license": "MIT"
|
||||
},
|
||||
"node_modules/jackspeak": {
|
||||
"version": "3.4.3",
|
||||
"resolved": "https://registry.npmjs.org/jackspeak/-/jackspeak-3.4.3.tgz",
|
||||
|
||||
@@ -19,6 +19,7 @@
|
||||
"devDependencies": {
|
||||
"@cloudflare/vitest-pool-workers": "^0.5.8",
|
||||
"@cloudflare/workers-types": "^4.20240729.0",
|
||||
"@octokit/webhooks-types": "^7.6.1",
|
||||
"@types/express": "^4.17.21",
|
||||
"@typescript-eslint/eslint-plugin": "^8.0.0",
|
||||
"@typescript-eslint/parser": "^8.0.0",
|
||||
@@ -38,6 +39,7 @@
|
||||
"@influxdata/influxdb-client": "^1.34.0",
|
||||
"fetch-retry": "^6.0.0",
|
||||
"fflate": "^0.8.2",
|
||||
"itty-router": "^5.0.18",
|
||||
"p-limit": "^6.1.0"
|
||||
},
|
||||
"volta": {
|
||||
|
||||
Vendored
+1
@@ -1,4 +1,5 @@
|
||||
interface WorkerEnv extends Omit<Env, 'ENVIRONMENT' | 'DEPLOYMENT_KEY'> {
|
||||
SLUG: string;
|
||||
ENVIRONMENT: string;
|
||||
DEPLOYMENT_KEY: string;
|
||||
VMETRICS_API_TOKEN: string;
|
||||
|
||||
+57
-52
@@ -1,59 +1,64 @@
|
||||
/* eslint-disable @typescript-eslint/no-unused-vars */
|
||||
import { IMetricsRepository } from './interface';
|
||||
import {
|
||||
CloudflareDeferredRepository,
|
||||
CloudflareMetricsRepository,
|
||||
HeaderMetricsProvider,
|
||||
InfluxMetricsProvider,
|
||||
} from './repository';
|
||||
import { AutoRouter, error, IRequest } from 'itty-router';
|
||||
import { QueueItem } from 'src/interfaces/queue.interface';
|
||||
import { CloudflareDeferredRepository } from 'src/repositories/cloudflare-deterred.repository';
|
||||
import { CloudflareMetricsRepository } from 'src/repositories/cloudflare-metrics.repository';
|
||||
import { CloudflareQueueRepository } from 'src/repositories/cloudflare-queue.repository';
|
||||
import { InfluxMetricsProvider } from 'src/repositories/influx-metrics-provider.repository';
|
||||
import { ApiWorker } from 'src/workers/api.worker';
|
||||
import { IngestApiWorker } from 'src/workers/ingest-api.worker';
|
||||
import { IngestProcessorWorker } from 'src/workers/ingest-processor.worker';
|
||||
|
||||
enum Header {
|
||||
PMTILES_DEPLOYMENT_KEY = 'PMTiles-Deployment-Key',
|
||||
CACHE_CONTROL = 'Cache-Control',
|
||||
ACCESS_CONTROL_ALLOW_ORIGIN = 'Access-Control-Allow-Origin',
|
||||
VARY = 'Vary',
|
||||
CONTENT_TYPE = 'Content-Type',
|
||||
CONTENT_ENCODING = 'Content-Encoding',
|
||||
SERVER_TIMING = 'Server-Timing',
|
||||
}
|
||||
type FetchRequest = IRequest & Parameters<ExportedHandlerFetchHandler>[0];
|
||||
|
||||
async function handleRequest(
|
||||
request: Request<unknown, IncomingRequestCfProperties>,
|
||||
env: WorkerEnv,
|
||||
deferredRepository: CloudflareDeferredRepository,
|
||||
metrics: IMetricsRepository,
|
||||
) {
|
||||
deferredRepository.defer(() => env.DATA_QUEUE.send({ messages: [{ hello: 'world' }] }));
|
||||
return new Response('Hello world!');
|
||||
}
|
||||
const asTags = (request: FetchRequest) => ({
|
||||
continent: request.cf?.continent ?? '',
|
||||
colo: request.cf?.colo ?? '',
|
||||
asOrg: request.cf?.asOrganization ?? '',
|
||||
});
|
||||
|
||||
const newApiWorker = (request: FetchRequest, env: WorkerEnv, ctx: ExecutionContext) => {
|
||||
const deferredRepository = new CloudflareDeferredRepository(ctx);
|
||||
const influxProvider = new InfluxMetricsProvider(env.VMETRICS_API_TOKEN, env.ENVIRONMENT);
|
||||
deferredRepository.defer(() => influxProvider.flush());
|
||||
const metrics = new CloudflareMetricsRepository('data', asTags(request), [influxProvider]);
|
||||
|
||||
return new ApiWorker(metrics);
|
||||
};
|
||||
|
||||
const newIngestApiWorker = (request: FetchRequest, env: WorkerEnv) => {
|
||||
const queue = new CloudflareQueueRepository(env.DATA_QUEUE);
|
||||
return new IngestApiWorker(queue);
|
||||
};
|
||||
|
||||
const withSlug = (req: FetchRequest, env: WorkerEnv) => {
|
||||
const slug = env.SLUG;
|
||||
if (slug && req.params.slug !== slug) {
|
||||
return error(401, 'Unauthorized');
|
||||
}
|
||||
};
|
||||
|
||||
const handleError = (err: Error) => {
|
||||
return error(500, err.message);
|
||||
};
|
||||
|
||||
const router = AutoRouter<FetchRequest, [WorkerEnv, ExecutionContext]>()
|
||||
.get('/api', (...args) => newApiWorker(...args).getGithubReports())
|
||||
.post('/ingest/github/:slug', withSlug, async (req, env) =>
|
||||
newIngestApiWorker(req, env)
|
||||
.onGithubEvent(await req.json())
|
||||
.catch(handleError)
|
||||
.then(() => new Response(null, { status: 204 })),
|
||||
);
|
||||
|
||||
export default {
|
||||
async queue(batch, env): Promise<void> {
|
||||
const messages = JSON.stringify(batch.messages);
|
||||
console.log(`consumed from our queue: ${messages}`);
|
||||
},
|
||||
fetch: router.fetch,
|
||||
queue: (batch, env) => {
|
||||
const influxProvider = new InfluxMetricsProvider(env.VMETRICS_API_TOKEN, env.ENVIRONMENT);
|
||||
const metricsRepository = new CloudflareMetricsRepository('data', {}, [influxProvider]);
|
||||
const ingestProcessWorker = new IngestProcessorWorker(metricsRepository);
|
||||
|
||||
async fetch(request, env, ctx): Promise<Response> {
|
||||
const workerEnv = env as WorkerEnv;
|
||||
const deferredRepository = new CloudflareDeferredRepository(ctx);
|
||||
const headerProvider = new HeaderMetricsProvider();
|
||||
const influxProvider = new InfluxMetricsProvider(workerEnv.VMETRICS_API_TOKEN, env.ENVIRONMENT);
|
||||
deferredRepository.defer(() => influxProvider.flush());
|
||||
const metrics = new CloudflareMetricsRepository('data', request, [influxProvider, headerProvider]);
|
||||
|
||||
try {
|
||||
const response = await metrics.monitorAsyncFunction({ name: 'handle_request' }, handleRequest)(
|
||||
request,
|
||||
workerEnv,
|
||||
deferredRepository,
|
||||
metrics,
|
||||
);
|
||||
response.headers.set(Header.SERVER_TIMING, headerProvider.getTimingHeader());
|
||||
deferredRepository.runDeferred();
|
||||
return response;
|
||||
} catch (e) {
|
||||
console.error(e);
|
||||
return new Response('Internal Server Error', { status: 500 });
|
||||
if (batch.queue.startsWith('data-ingest')) {
|
||||
return ingestProcessWorker.process(batch);
|
||||
}
|
||||
},
|
||||
} satisfies ExportedHandler<Env>;
|
||||
} satisfies ExportedHandler<WorkerEnv, QueueItem>;
|
||||
|
||||
@@ -1,25 +0,0 @@
|
||||
import { Metric } from './repository';
|
||||
|
||||
export interface IDeferredRepository {
|
||||
defer(promise: AsyncFn): void;
|
||||
runDeferred(): void;
|
||||
}
|
||||
|
||||
export type AsyncFn = (...args: any[]) => Promise<any>;
|
||||
export type Class = { new (...args: any[]): any };
|
||||
export type Operation = { name: string; tags?: { [key: string]: string } };
|
||||
export type Options = { monitorInvocations?: boolean; acceptedErrors?: Class[] };
|
||||
|
||||
export interface IMetricsRepository {
|
||||
monitorAsyncFunction<T extends AsyncFn>(
|
||||
operation: Operation,
|
||||
call: T,
|
||||
options?: Options,
|
||||
): (...args: Parameters<T>) => Promise<Awaited<ReturnType<T>>>;
|
||||
push(metric: Metric): void;
|
||||
}
|
||||
|
||||
export interface IMetricsProviderRepository {
|
||||
pushMetric(metric: Metric): void;
|
||||
flush(): void;
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
import { AsyncFn } from 'src/types';
|
||||
|
||||
export interface IDeferredRepository {
|
||||
defer(promise: AsyncFn): void;
|
||||
runDeferred(): void;
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
import { Metric } from 'src/interfaces/metrics.interface';
|
||||
|
||||
export interface IMetricsProviderRepository {
|
||||
pushMetric(metric: Metric): void;
|
||||
flush(): void;
|
||||
}
|
||||
@@ -0,0 +1,59 @@
|
||||
import { AsyncFn, Operation, Options } from 'src/types';
|
||||
|
||||
export interface IMetricsRepository {
|
||||
monitorAsyncFunction<T extends AsyncFn>(
|
||||
operation: Operation,
|
||||
call: T,
|
||||
options?: Options,
|
||||
): (...args: Parameters<T>) => Promise<Awaited<ReturnType<T>>>;
|
||||
push(metric: Metric): void;
|
||||
}
|
||||
|
||||
export class Metric {
|
||||
private _tags: Map<string, string> = new Map();
|
||||
private _timestamp = performance.now();
|
||||
private _fields = new Map<string, { value: any; type: 'duration' | 'int' }>();
|
||||
private constructor(private _name: string) {}
|
||||
|
||||
static create(name: string) {
|
||||
return new Metric(name);
|
||||
}
|
||||
|
||||
get tags() {
|
||||
return this._tags;
|
||||
}
|
||||
|
||||
get timestamp() {
|
||||
return this._timestamp;
|
||||
}
|
||||
|
||||
get fields() {
|
||||
return this._fields;
|
||||
}
|
||||
|
||||
get name() {
|
||||
return this._name;
|
||||
}
|
||||
|
||||
addTag(key: string, value: string) {
|
||||
this._tags.set(key, value);
|
||||
return this;
|
||||
}
|
||||
|
||||
addTags(tags: { [key: string]: string }) {
|
||||
for (const [key, value] of Object.entries(tags)) {
|
||||
this._tags.set(key, value);
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
durationField(key: string, duration?: number) {
|
||||
this._fields.set(key, { value: duration ?? performance.now() - this._timestamp, type: 'duration' });
|
||||
return this;
|
||||
}
|
||||
|
||||
intField(key: string, value: number) {
|
||||
this._fields.set(key, { value, type: 'int' });
|
||||
return this;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,43 @@
|
||||
import {
|
||||
DiscussionCreatedEvent,
|
||||
IssuesClosedEvent,
|
||||
IssuesOpenedEvent,
|
||||
IssuesReopenedEvent,
|
||||
PullRequestClosedEvent,
|
||||
PullRequestOpenedEvent,
|
||||
PullRequestReopenedEvent,
|
||||
StarCreatedEvent,
|
||||
StarDeletedEvent,
|
||||
} from '@octokit/webhooks-types';
|
||||
import { DiscussionClosedEvent, DiscussionReopenedEvent } from 'src/types';
|
||||
|
||||
export enum EventType {
|
||||
GITHUB_STAR_CREATED = 'github.starCreated',
|
||||
}
|
||||
|
||||
type QueueEventMap = {
|
||||
RepositoryStarCreatedV1: StarCreatedEvent;
|
||||
RepositoryStarDeletedV1: StarDeletedEvent;
|
||||
RepositoryIssueOpenedV1: IssuesOpenedEvent;
|
||||
RepositoryIssueClosedV1: IssuesClosedEvent;
|
||||
RepositoryIssueReopenedV1: IssuesReopenedEvent;
|
||||
RepositoryPullRequestOpenedV1: PullRequestOpenedEvent;
|
||||
RepositoryPullRequestClosedV1: PullRequestClosedEvent;
|
||||
RepositoryPullRequestReopenedV1: PullRequestReopenedEvent;
|
||||
RepositoryDiscussionCreatedV1: DiscussionCreatedEvent;
|
||||
RepositoryDiscussionClosedV1: DiscussionClosedEvent;
|
||||
RepositoryDiscussionReopenedV1: DiscussionReopenedEvent;
|
||||
};
|
||||
|
||||
export type QueueEvent = keyof QueueEventMap;
|
||||
|
||||
export type QueueItem<T extends QueueEvent = QueueEvent> = {
|
||||
id: string;
|
||||
source: string;
|
||||
type: T;
|
||||
data: QueueEventMap[T];
|
||||
};
|
||||
|
||||
export interface IQueueRepository {
|
||||
push<T extends QueueEvent>(item: QueueItem<T>): Promise<void>;
|
||||
}
|
||||
@@ -1,5 +1,5 @@
|
||||
import { AsyncFn, Operation, Options } from './interface';
|
||||
import { Metric } from './repository';
|
||||
import { Metric } from 'src/interfaces/metrics.interface';
|
||||
import { AsyncFn, Operation, Options } from 'src/types';
|
||||
|
||||
export function monitorAsyncFunction<T extends AsyncFn>(
|
||||
operationPrefix: string,
|
||||
|
||||
@@ -0,0 +1,17 @@
|
||||
import { IDeferredRepository } from 'src/interfaces/deferred.interface';
|
||||
import { AsyncFn } from 'src/types';
|
||||
|
||||
export class CloudflareDeferredRepository implements IDeferredRepository {
|
||||
deferred: AsyncFn[] = [];
|
||||
constructor(private ctx: ExecutionContext) {}
|
||||
|
||||
defer(call: AsyncFn): void {
|
||||
this.deferred.push(call);
|
||||
}
|
||||
|
||||
runDeferred() {
|
||||
for (const call of this.deferred) {
|
||||
this.ctx.waitUntil(call());
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,38 @@
|
||||
import { IMetricsProviderRepository } from 'src/interfaces/metrics-provider.interface';
|
||||
import { IMetricsRepository, Metric } from 'src/interfaces/metrics.interface';
|
||||
import { monitorAsyncFunction } from 'src/monitor';
|
||||
import { AsyncFn, Operation, Options } from 'src/types';
|
||||
|
||||
export class CloudflareMetricsRepository implements IMetricsRepository {
|
||||
private readonly defaultTags: { [key: string]: string };
|
||||
|
||||
constructor(
|
||||
private operationPrefix: string,
|
||||
tags: Record<string, string>,
|
||||
private metricsProviders: IMetricsProviderRepository[],
|
||||
) {
|
||||
this.defaultTags = { ...tags };
|
||||
}
|
||||
|
||||
monitorAsyncFunction<T extends AsyncFn>(
|
||||
operation: Operation,
|
||||
call: T,
|
||||
options: Options = {},
|
||||
): (...args: Parameters<T>) => Promise<Awaited<ReturnType<T>>> {
|
||||
operation = { ...operation, tags: { ...operation.tags, ...this.defaultTags } };
|
||||
const callback = (metric: Metric) => {
|
||||
for (const provider of this.metricsProviders) {
|
||||
provider.pushMetric(metric);
|
||||
}
|
||||
};
|
||||
|
||||
return monitorAsyncFunction(this.operationPrefix, operation, call, callback, options);
|
||||
}
|
||||
|
||||
push(metric: Metric) {
|
||||
metric.addTags(this.defaultTags);
|
||||
for (const provider of this.metricsProviders) {
|
||||
provider.pushMetric(metric);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
import { IQueueRepository, QueueItem } from 'src/interfaces/queue.interface';
|
||||
|
||||
export class CloudflareQueueRepository implements IQueueRepository {
|
||||
constructor(private queue: Queue) {}
|
||||
|
||||
async push(item: QueueItem) {
|
||||
await this.queue.send(item);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,48 @@
|
||||
import { Point } from '@influxdata/influxdb-client';
|
||||
import { IMetricsProviderRepository } from 'src/interfaces/metrics-provider.interface';
|
||||
import { Metric } from 'src/interfaces/metrics.interface';
|
||||
|
||||
export class InfluxMetricsProvider implements IMetricsProviderRepository {
|
||||
private metrics: string[] = [];
|
||||
constructor(
|
||||
private influxApiToken: string,
|
||||
private environment: string,
|
||||
) {}
|
||||
|
||||
pushMetric(metric: Metric) {
|
||||
const point = new Point(metric.name);
|
||||
for (const [key, value] of metric.tags) {
|
||||
point.tag(key, value);
|
||||
}
|
||||
for (const [key, { value, type }] of metric.fields) {
|
||||
if (type === 'duration') {
|
||||
point.intField(key, value);
|
||||
} else if (type === 'int') {
|
||||
point.intField(key, value);
|
||||
}
|
||||
}
|
||||
const influxLineProtocol = point.toLineProtocol()?.toString();
|
||||
if (influxLineProtocol) {
|
||||
this.metrics.push(influxLineProtocol);
|
||||
}
|
||||
}
|
||||
|
||||
async flush() {
|
||||
if (this.metrics.length === 0) {
|
||||
return;
|
||||
}
|
||||
const metrics = this.metrics.join('\n');
|
||||
if (this.environment === 'prod') {
|
||||
const response = await fetch('https://cf-workers.monitoring.immich.cloud/write', {
|
||||
method: 'POST',
|
||||
body: this.metrics.join('\n'),
|
||||
headers: {
|
||||
Authorization: `Token ${this.influxApiToken}`,
|
||||
},
|
||||
});
|
||||
await response.body?.cancel();
|
||||
} else {
|
||||
console.log(metrics);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,179 +0,0 @@
|
||||
import { Point } from '@influxdata/influxdb-client';
|
||||
import {
|
||||
AsyncFn,
|
||||
IDeferredRepository,
|
||||
IMetricsProviderRepository,
|
||||
IMetricsRepository,
|
||||
Operation,
|
||||
Options,
|
||||
} from './interface';
|
||||
import { monitorAsyncFunction } from './monitor';
|
||||
|
||||
export class CloudflareDeferredRepository implements IDeferredRepository {
|
||||
deferred: AsyncFn[] = [];
|
||||
constructor(private ctx: ExecutionContext) {}
|
||||
|
||||
defer(call: AsyncFn): void {
|
||||
this.deferred.push(call);
|
||||
}
|
||||
|
||||
runDeferred() {
|
||||
for (const call of this.deferred) {
|
||||
this.ctx.waitUntil(call());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export class HeaderMetricsProvider implements IMetricsProviderRepository {
|
||||
private _metrics: string[] = [];
|
||||
constructor() {}
|
||||
|
||||
pushMetric(metric: Metric) {
|
||||
for (const [label, { value, type }] of metric.fields) {
|
||||
if (type === 'duration') {
|
||||
const suffix = label === 'duration' ? '' : `_${label.replace('_duration', '')}`;
|
||||
this._metrics.push(`${metric.name}${suffix};dur=${value}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
getTimingHeader() {
|
||||
return this._metrics.join(', ');
|
||||
}
|
||||
|
||||
flush() {
|
||||
console.log(this._metrics.join(', '));
|
||||
}
|
||||
}
|
||||
|
||||
export class InfluxMetricsProvider implements IMetricsProviderRepository {
|
||||
private metrics: string[] = [];
|
||||
constructor(
|
||||
private influxApiToken: string,
|
||||
private environment: string,
|
||||
) {}
|
||||
|
||||
pushMetric(metric: Metric) {
|
||||
const point = new Point(metric.name);
|
||||
for (const [key, value] of metric.tags) {
|
||||
point.tag(key, value);
|
||||
}
|
||||
for (const [key, { value, type }] of metric.fields) {
|
||||
if (type === 'duration') {
|
||||
point.intField(key, value);
|
||||
} else if (type === 'int') {
|
||||
point.intField(key, value);
|
||||
}
|
||||
}
|
||||
const influxLineProtocol = point.toLineProtocol()?.toString();
|
||||
if (influxLineProtocol) {
|
||||
this.metrics.push(influxLineProtocol);
|
||||
}
|
||||
}
|
||||
|
||||
async flush() {
|
||||
if (this.metrics.length === 0) {
|
||||
return;
|
||||
}
|
||||
const metrics = this.metrics.join('\n');
|
||||
if (this.environment === 'prod') {
|
||||
const response = await fetch('https://cf-workers.monitoring.immich.cloud/write', {
|
||||
method: 'POST',
|
||||
body: this.metrics.join('\n'),
|
||||
headers: {
|
||||
Authorization: `Token ${this.influxApiToken}`,
|
||||
},
|
||||
});
|
||||
await response.body?.cancel();
|
||||
} else {
|
||||
console.log(metrics);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export class Metric {
|
||||
private _tags: Map<string, string> = new Map();
|
||||
private _timestamp = performance.now();
|
||||
private _fields = new Map<string, { value: any; type: 'duration' | 'int' }>();
|
||||
private constructor(private _name: string) {}
|
||||
|
||||
static create(name: string) {
|
||||
return new Metric(name);
|
||||
}
|
||||
|
||||
get tags() {
|
||||
return this._tags;
|
||||
}
|
||||
|
||||
get timestamp() {
|
||||
return this._timestamp;
|
||||
}
|
||||
|
||||
get fields() {
|
||||
return this._fields;
|
||||
}
|
||||
|
||||
get name() {
|
||||
return this._name;
|
||||
}
|
||||
|
||||
addTag(key: string, value: string) {
|
||||
this._tags.set(key, value);
|
||||
return this;
|
||||
}
|
||||
|
||||
addTags(tags: { [key: string]: string }) {
|
||||
for (const [key, value] of Object.entries(tags)) {
|
||||
this._tags.set(key, value);
|
||||
}
|
||||
return this;
|
||||
}
|
||||
|
||||
durationField(key: string, duration?: number) {
|
||||
this._fields.set(key, { value: duration ?? performance.now() - this._timestamp, type: 'duration' });
|
||||
return this;
|
||||
}
|
||||
|
||||
intField(key: string, value: number) {
|
||||
this._fields.set(key, { value, type: 'int' });
|
||||
return this;
|
||||
}
|
||||
}
|
||||
|
||||
export class CloudflareMetricsRepository implements IMetricsRepository {
|
||||
private readonly defaultTags: { [key: string]: string };
|
||||
|
||||
constructor(
|
||||
private operationPrefix: string,
|
||||
request: Request<unknown, IncomingRequestCfProperties>,
|
||||
private metricsProviders: IMetricsProviderRepository[],
|
||||
) {
|
||||
this.defaultTags = {
|
||||
continent: request.cf?.continent ?? '',
|
||||
colo: request.cf?.colo ?? '',
|
||||
asOrg: request.cf?.asOrganization ?? '',
|
||||
};
|
||||
}
|
||||
|
||||
monitorAsyncFunction<T extends AsyncFn>(
|
||||
operation: Operation,
|
||||
call: T,
|
||||
options: Options = {},
|
||||
): (...args: Parameters<T>) => Promise<Awaited<ReturnType<T>>> {
|
||||
operation = { ...operation, tags: { ...operation.tags, ...this.defaultTags } };
|
||||
const callback = (metric: Metric) => {
|
||||
for (const provider of this.metricsProviders) {
|
||||
provider.pushMetric(metric);
|
||||
}
|
||||
};
|
||||
|
||||
return monitorAsyncFunction(this.operationPrefix, operation, call, callback, options);
|
||||
}
|
||||
|
||||
push(metric: Metric) {
|
||||
metric.addTags(this.defaultTags);
|
||||
for (const provider of this.metricsProviders) {
|
||||
provider.pushMetric(metric);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,9 @@
|
||||
import { DiscussionCreatedEvent, DiscussionDeletedEvent, WebhookEvent } from '@octokit/webhooks-types';
|
||||
|
||||
export type AsyncFn = (...args: any[]) => Promise<any>;
|
||||
export type Class = { new (...args: any[]): any };
|
||||
export type Operation = { name: string; tags?: { [key: string]: string } };
|
||||
export type Options = { monitorInvocations?: boolean; acceptedErrors?: Class[] };
|
||||
export type GithubWebhookEvent = WebhookEvent | DiscussionClosedEvent | DiscussionReopenedEvent;
|
||||
export type DiscussionClosedEvent = Omit<DiscussionDeletedEvent, 'action'> & { action: 'closed' };
|
||||
export type DiscussionReopenedEvent = Omit<DiscussionCreatedEvent, 'action'> & { action: 'reopened' };
|
||||
@@ -0,0 +1,9 @@
|
||||
import { IMetricsRepository } from 'src/interfaces/metrics.interface';
|
||||
|
||||
export class ApiWorker {
|
||||
constructor(private metricsRepository: IMetricsRepository) {}
|
||||
|
||||
getGithubReports() {
|
||||
return [];
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,100 @@
|
||||
import {
|
||||
DiscussionCreatedEvent,
|
||||
DiscussionEvent,
|
||||
IssuesClosedEvent,
|
||||
IssuesEvent,
|
||||
IssuesOpenedEvent,
|
||||
IssuesReopenedEvent,
|
||||
PullRequestClosedEvent,
|
||||
PullRequestEvent,
|
||||
PullRequestOpenedEvent,
|
||||
PullRequestReopenedEvent,
|
||||
StarCreatedEvent,
|
||||
StarDeletedEvent,
|
||||
StarEvent,
|
||||
} from '@octokit/webhooks-types';
|
||||
import { IQueueRepository, QueueEvent, QueueItem } from 'src/interfaces/queue.interface';
|
||||
import { DiscussionClosedEvent, DiscussionReopenedEvent, GithubWebhookEvent as Event } from 'src/types';
|
||||
|
||||
const isStarCreated = (event: Event): event is StarCreatedEvent => 'starred_at' in event && event.action === 'created';
|
||||
const isStarDeleted = (event: Event): event is StarDeletedEvent => 'starred_at' in event && event.action === 'deleted';
|
||||
const asStarId = (event: StarEvent) =>
|
||||
`${event.repository.full_name}-${event.sender.login}-${event.action}-${event.starred_at ?? Date.now()}`;
|
||||
|
||||
const isIssueOpened = (event: Event): event is IssuesOpenedEvent => 'issue' in event && event.action === 'opened';
|
||||
const isIssueClosed = (event: Event): event is IssuesClosedEvent => 'issue' in event && event.action === 'closed';
|
||||
const isIssueReopened = (event: Event): event is IssuesReopenedEvent => 'issue' in event && event.action === 'reopened';
|
||||
const asIssueId = (event: IssuesEvent) =>
|
||||
`${event.repository.full_name}-${event.issue.number}-${event.action}-${event.issue.updated_at}`;
|
||||
|
||||
const isPullRequestOpened = (event: Event): event is PullRequestOpenedEvent =>
|
||||
'pull_request' in event && event.action === 'opened';
|
||||
const isPullRequestClosed = (event: Event): event is PullRequestClosedEvent =>
|
||||
'pull_request' in event && event.action === 'closed';
|
||||
const isPullRequestReopened = (event: Event): event is PullRequestReopenedEvent =>
|
||||
'pull_request' in event && event.action === 'reopened';
|
||||
const asPullRequestId = (event: PullRequestEvent) =>
|
||||
`${event.repository.full_name}-${event.pull_request.number}-${event.action}-${event.pull_request.updated_at}`;
|
||||
|
||||
const isDiscussionCreated = (event: Event): event is DiscussionCreatedEvent =>
|
||||
'discussion' in event && event.action === 'created';
|
||||
const isDiscussionClosed = (event: Event): event is DiscussionClosedEvent =>
|
||||
'discussion' in event && event.action === 'closed';
|
||||
const isDiscussionReopened = (event: Event): event is DiscussionReopenedEvent =>
|
||||
'discussion' in event && event.action === 'reopened';
|
||||
const asDiscussionId = (event: DiscussionEvent | DiscussionClosedEvent | DiscussionReopenedEvent) =>
|
||||
`${event.repository.full_name}-${event.discussion.number}-${event.action}-${event.discussion.updated_at}`;
|
||||
|
||||
export class IngestApiWorker {
|
||||
constructor(private queueRepository: IQueueRepository) {}
|
||||
|
||||
async onGithubEvent(event: Event) {
|
||||
if (isStarCreated(event)) {
|
||||
return this.push(asStarId(event), 'RepositoryStarCreatedV1', event);
|
||||
}
|
||||
|
||||
if (isStarDeleted(event)) {
|
||||
return this.push(asStarId(event), 'RepositoryStarDeletedV1', event);
|
||||
}
|
||||
|
||||
if (isIssueOpened(event)) {
|
||||
return this.push(asIssueId(event), 'RepositoryIssueOpenedV1', event);
|
||||
}
|
||||
|
||||
if (isIssueClosed(event)) {
|
||||
return this.push(asIssueId(event), 'RepositoryIssueClosedV1', event);
|
||||
}
|
||||
|
||||
if (isIssueReopened(event)) {
|
||||
return this.push(asIssueId(event), 'RepositoryIssueReopenedV1', event);
|
||||
}
|
||||
|
||||
if (isPullRequestOpened(event)) {
|
||||
return this.push(asPullRequestId(event), 'RepositoryPullRequestOpenedV1', event);
|
||||
}
|
||||
|
||||
if (isPullRequestClosed(event)) {
|
||||
return this.push(asPullRequestId(event), 'RepositoryPullRequestClosedV1', event);
|
||||
}
|
||||
|
||||
if (isPullRequestReopened(event)) {
|
||||
return this.push(asPullRequestId(event), 'RepositoryPullRequestReopenedV1', event);
|
||||
}
|
||||
|
||||
if (isDiscussionCreated(event)) {
|
||||
return this.push(asDiscussionId(event), 'RepositoryDiscussionCreatedV1', event);
|
||||
}
|
||||
|
||||
if (isDiscussionClosed(event)) {
|
||||
return this.push(asDiscussionId(event), 'RepositoryDiscussionClosedV1', event);
|
||||
}
|
||||
|
||||
if (isDiscussionReopened(event)) {
|
||||
return this.push(asDiscussionId(event), 'RepositoryDiscussionReopenedV1', event);
|
||||
}
|
||||
}
|
||||
|
||||
private async push<T extends QueueEvent>(id: string, type: T, data: QueueItem<T>['data']) {
|
||||
await this.queueRepository.push({ id, type, source: 'ingest-api-worker', data });
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
import { IMetricsRepository } from 'src/interfaces/metrics.interface';
|
||||
import { QueueItem } from 'src/interfaces/queue.interface';
|
||||
|
||||
export class IngestProcessorWorker {
|
||||
constructor(private metricsRepository: IMetricsRepository) {}
|
||||
|
||||
async process(batch: MessageBatch<QueueItem>) {
|
||||
console.log(`Received ${batch.messages.length} messages`);
|
||||
// await this.metricsRepository.push({} as any);
|
||||
}
|
||||
}
|
||||
+22
-21
@@ -1,23 +1,24 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "es2021" /* Set the JavaScript language version for emitted JavaScript and include compatible library declarations. */,
|
||||
"lib": ["es2021"] /* Specify a set of bundled library declaration files that describe the target runtime environment. */,
|
||||
"jsx": "react-jsx" /* Specify what JSX code is generated. */,
|
||||
"module": "es2022" /* Specify what module code is generated. */,
|
||||
"moduleResolution": "Bundler" /* Specify how TypeScript looks up a file from a given module specifier. */,
|
||||
"types": [
|
||||
"@cloudflare/workers-types"
|
||||
],
|
||||
"resolveJsonModule": true /* Enable importing .json files */,
|
||||
"allowJs": true /* Allow JavaScript files to be a part of your program. Use the `checkJS` option to get errors from these files. */,
|
||||
"checkJs": false /* Enable error reporting in type-checked JavaScript files. */,
|
||||
"noEmit": true /* Disable emitting files from a compilation. */,
|
||||
"isolatedModules": true /* Ensure that each file can be safely transpiled without relying on other imports. */,
|
||||
"allowSyntheticDefaultImports": true /* Allow 'import x from y' when a module doesn't have a default export. */,
|
||||
"forceConsistentCasingInFileNames": true /* Ensure that casing is correct in imports. */,
|
||||
"strict": true /* Enable all strict type-checking options. */,
|
||||
"skipLibCheck": true /* Skip type checking all .d.ts files. */
|
||||
},
|
||||
"exclude": ["test"],
|
||||
"include": ["worker-configuration.d.ts", "src/**/*.ts", "vitest.config.mts"]
|
||||
"compilerOptions": {
|
||||
"target": "es2021" /* Set the JavaScript language version for emitted JavaScript and include compatible library declarations. */,
|
||||
"lib": [
|
||||
"es2021"
|
||||
] /* Specify a set of bundled library declaration files that describe the target runtime environment. */,
|
||||
"jsx": "react-jsx" /* Specify what JSX code is generated. */,
|
||||
"module": "es2022" /* Specify what module code is generated. */,
|
||||
"baseUrl": "./",
|
||||
"moduleResolution": "Bundler" /* Specify how TypeScript looks up a file from a given module specifier. */,
|
||||
"types": ["@cloudflare/workers-types"],
|
||||
"resolveJsonModule": true /* Enable importing .json files */,
|
||||
"allowJs": true /* Allow JavaScript files to be a part of your program. Use the `checkJS` option to get errors from these files. */,
|
||||
"checkJs": false /* Enable error reporting in type-checked JavaScript files. */,
|
||||
"noEmit": true /* Disable emitting files from a compilation. */,
|
||||
"isolatedModules": true /* Ensure that each file can be safely transpiled without relying on other imports. */,
|
||||
"allowSyntheticDefaultImports": true /* Allow 'import x from y' when a module doesn't have a default export. */,
|
||||
"forceConsistentCasingInFileNames": true /* Ensure that casing is correct in imports. */,
|
||||
"strict": true /* Enable all strict type-checking options. */,
|
||||
"skipLibCheck": true /* Skip type checking all .d.ts files. */
|
||||
},
|
||||
"exclude": ["test"],
|
||||
"include": ["worker-configuration.d.ts", "src/**/*.ts", "vitest.config.mts"]
|
||||
}
|
||||
|
||||
Vendored
+1
@@ -3,5 +3,6 @@
|
||||
interface Env {
|
||||
DEPLOYMENT_KEY: "20240812";
|
||||
ENVIRONMENT: "production";
|
||||
SLUG: "slug";
|
||||
DATA_QUEUE: Queue;
|
||||
}
|
||||
|
||||
@@ -8,6 +8,7 @@ compatibility_flags = ["nodejs_compat"]
|
||||
[vars]
|
||||
DEPLOYMENT_KEY = "20240812"
|
||||
ENVIRONMENT = "production"
|
||||
SLUG = "slug"
|
||||
|
||||
[[queues.producers]]
|
||||
queue = "data-ingest-dev"
|
||||
|
||||
Reference in New Issue
Block a user