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
2 changes: 1 addition & 1 deletion proto
Submodule proto updated 1 files
+9 −0 event/v1/event.proto
153 changes: 152 additions & 1 deletion src/gen/event/v1/event.ts

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

77 changes: 72 additions & 5 deletions src/routes/gRPC/events/streamEvents.ts
Original file line number Diff line number Diff line change
@@ -1,22 +1,57 @@
import type { ServerReadableStream, sendUnaryData } from "@grpc/grpc-js";
import * as Sentry from "@sentry/bun";
import { ZodError } from "zod";
import {
StreamEventRequest,
StreamEventResponse,
EventFailure,
} from "../../../gen/event/v1/event";
import { EventError } from "../../../errors/event";
import { AuthError } from "../../../errors/auth";
import { StorageError } from "../../../errors/storage";
import { streamEventSchema } from "../../../zod/event";
import { createEventInstance, storeEvent } from "../../../utils/eventHelpers";
import { apiKeyContextKey } from "../../../context/auth";
import { wideEventContextKey } from "../../../context/requestContext";
import type { ContextStreamCall } from "../../../interface/types/context";

function getFailureCode(err: unknown): string {
if (err instanceof StorageError) {
if (err.type === "CONSTRAINT_VIOLATION") return "DUPLICATE_IDEMPOTENCY_KEY";
if (err.type === "INVALID_DATA") return "INVALID_DATA";
if (err.type === "INVALID_TIMESTAMP") return "INVALID_TIMESTAMP";
if (err.type === "PRICE_CALCULATION_FAILED") return "PRICE_CALCULATION_FAILED";
return "STORAGE_FAILURE";
}
if (err instanceof EventError) {
if (err.type === "UNSUPPORTED_EVENT_TYPE") return "UNSUPPORTED_EVENT_TYPE";
return "VALIDATION_FAILED";
}
if (err instanceof ZodError) {
return "VALIDATION_FAILED";
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
return "INTERNAL_ERROR";
}

function publicMessageForCode(code: string): string {
switch (code) {
case "DUPLICATE_IDEMPOTENCY_KEY": return "Duplicate idempotency key";
case "VALIDATION_FAILED": return "Event validation failed";
case "INVALID_DATA": return "Invalid event data";
case "INVALID_TIMESTAMP": return "Invalid event timestamp";
case "PRICE_CALCULATION_FAILED": return "Price calculation failed";
case "UNSUPPORTED_EVENT_TYPE": return "Unsupported event type";
case "STORAGE_FAILURE": return "Storage error";
default: return "Internal server error";
}
}

export async function streamEvents(
call: ContextStreamCall,
callback: sendUnaryData<StreamEventResponse>
): Promise<void> {
let eventsProcessed = 0;
const failures: EventFailure[] = [];

const wideEventBuilder = call[wideEventContextKey];
const auth = call[apiKeyContextKey];
Expand All @@ -36,6 +71,8 @@ export async function streamEvents(
);
}

let eventIndex = 0;

for await (const req of call) {
try {
const eventSkeleton = await streamEventSchema.parseAsync({ ...req });
Expand All @@ -52,20 +89,50 @@ export async function streamEvents(
await storeEvent(event, auth);
eventsProcessed++;
} catch (innerError) {
const errorCode = getFailureCode(innerError);

Sentry.addBreadcrumb({
category: "streamEvents",
message: `Event processing failed: ${innerError instanceof Error ? innerError.message : String(innerError)}`,
message: `Event [${eventIndex}] processing failed: ${innerError instanceof Error ? innerError.message : String(innerError)}`,
data: { eventIndex, idempotencyKey: req.idempotencyKey || "<unknown>" },
level: "error",
});
Sentry.captureException(innerError);
callback(innerError as Error, null);
return;
Sentry.captureException(innerError, {
extra: { eventIndex, idempotencyKey: req.idempotencyKey || "<unknown>", errorCode },
});

const failure = EventFailure.create();
failure.eventIndex = eventIndex;
failure.idempotencyKey = req.idempotencyKey || "<unknown>";
failure.errorCode = errorCode;
failure.message = publicMessageForCode(errorCode);
failures.push(failure);
}

eventIndex++;
}

const response = StreamEventResponse.create();
response.eventsProcessed = eventsProcessed;
response.message = `Successfully processed ${eventsProcessed} events`;
response.eventsFailed = failures.length;
response.failures = failures;
const total = eventsProcessed + failures.length;
response.message = `Processed ${total} events (${eventsProcessed} succeeded, ${failures.length} failed)`;

wideEventBuilder?.setEventContext({
eventType: "AI_TOKEN_USAGE",
eventCount: eventsProcessed,
});
if (failures.length > 0) {
wideEventBuilder?.addContext({
eventsFailed: failures.length,
eventFailures: failures.map((f) => ({
eventIndex: f.eventIndex,
errorCode: f.errorCode,
idempotencyKey: f.idempotencyKey,
})),
});
}

callback(null, response);
} catch (error) {
Expand Down
27 changes: 27 additions & 0 deletions src/storage/adapter/postgres/handlers/addEventUtils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,24 @@ export type TransactionFn<T> = (
txn: PgTransaction<any, any, any>
) => Promise<T>;

function getPostgresErrorCode(e: unknown): string | null {
if (e && typeof e === "object" && "code" in e) {
const code = (e as { code: unknown }).code;
if (typeof code === "string") return code;
}
if (e instanceof StorageError && e.originalError) {
return getPostgresErrorCode(e.originalError);
}
if (e && typeof e === "object" && "cause" in e) {
return getPostgresErrorCode((e as { cause: unknown }).cause);
}
return null;
}

function hasPostgresErrorCode(e: unknown, code: string): boolean {
return getPostgresErrorCode(e) === code;
}

export async function executeInTransaction<T>(
connectionObject: PgDatabase<any, any>,
operationName: string,
Expand All @@ -20,6 +38,15 @@ export async function executeInTransaction<T>(
try {
return await connectionObject.transaction(fn);
} catch (e) {
if (hasPostgresErrorCode(e, "23505")) {
throw StorageError.constraintViolation(
"Duplicate key violation",
e instanceof Error ? e : new Error(String(e))
);
}
if (e instanceof StorageError) {
throw e;
}
throw StorageError.transactionFailed(
`Transaction failed while ${operationName}`,
e instanceof Error ? e : new Error(String(e))
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Expand Down
Loading