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
7 changes: 4 additions & 3 deletions main.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ import * as Persistency from "./src/persist.ts";
import QueueManager from "./src/manager.ts";
import { createHandler } from "./src/handler.ts";
import { parseConfig } from "./src/config.ts";
import type { Payload } from "./src/payload.ts";

const LOG_ENCODER = new TextEncoder();

Expand Down Expand Up @@ -35,13 +36,13 @@ const CONFIG = parseConfig(readEnv(), Deno.args);

// Set up our persistency manager
const PERSIST_ENGINE = CONFIG.persistEnabled
? new Persistency.FileStore
: new Persistency.MemoryStore;
? new Persistency.FileStore<Payload>()
: new Persistency.MemoryStore<Payload>();

PERSIST_ENGINE.dir(CONFIG.persistDir);

// Set up the manager, which will handle our queues for us
const MANAGER = new QueueManager(PERSIST_ENGINE, CONFIG.queueDepthLimit, CONFIG.queueCountLimit, CONFIG.persistEnabled);
const MANAGER = new QueueManager<Payload>(PERSIST_ENGINE, CONFIG.queueDepthLimit, CONFIG.queueCountLimit, CONFIG.persistEnabled);

// Load up any existing queue data, if we're persisting
if (PERSIST_ENGINE instanceof Persistency.FileStore) {
Expand Down
24 changes: 14 additions & 10 deletions src/handler.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,16 @@
import QueueManager, { QueueNameTooLongError } from "./manager.ts";
import { RateLimiter } from "./rate_limiter.ts";
import { withAuth, withRateLimit } from "./middleware.ts";
import { RouteHandler, Router } from "./router.ts";
import { Router } from "./router.ts";
import * as Payload from "./payload.ts";

type JsonPayload = Payload.Payload;
type RouteHandler = Parameters<Router["get"]>[1];
type RouteMatch = Parameters<RouteHandler>[1];

const LOG_ENCODER = Reflect.construct(TextEncoder, []);

function extractQueueName(match: Parameters<RouteHandler>[1]): { name: string } | { error: Response } {
function extractQueueName(match: RouteMatch): { name: string } | { error: Response } {
const raw = match.pathname.groups.queue;
if (raw === undefined) {
return { error: new Response("Invalid queue name", { status: 400 }) };
Expand All @@ -21,7 +25,7 @@ function extractQueueName(match: Parameters<RouteHandler>[1]): { name: string }
}
}

function enqueueHandler(mgr: QueueManager<string>): RouteHandler {
function enqueueHandler(mgr: QueueManager<JsonPayload>): RouteHandler {
return async (request, match) => {
const queueResult = extractQueueName(match);
if ("error" in queueResult) {
Expand All @@ -33,7 +37,7 @@ function enqueueHandler(mgr: QueueManager<string>): RouteHandler {
if (contentLength && parseInt(contentLength) > Payload.DEFAULT_MAX_PAYLOAD_SIZE) {
return new Response("Payload too large", { status: 413 });
}
const payload = await Payload.readAndValidatePayload<string>(
const payload = await Payload.readAndValidatePayload(
request.body,
Payload.DEFAULT_MAX_PAYLOAD_SIZE,
);
Expand Down Expand Up @@ -71,7 +75,7 @@ function queueNameErrorResponse(error: unknown): Response {
throw error;
}

function itemResponse(item: unknown): Response {
function itemResponse(item: JsonPayload | undefined): Response {
if (item === undefined) {
return new Response(null, { status: 204 });
}
Expand All @@ -80,7 +84,7 @@ function itemResponse(item: unknown): Response {
});
}

function dequeueHandler(mgr: QueueManager<string>): RouteHandler {
function dequeueHandler(mgr: QueueManager<JsonPayload>): RouteHandler {
return (request, match) => {
void request;
const queueResult = extractQueueName(match);
Expand All @@ -98,7 +102,7 @@ function dequeueHandler(mgr: QueueManager<string>): RouteHandler {
};
}

function peekHandler(mgr: QueueManager<string>): RouteHandler {
function peekHandler(mgr: QueueManager<JsonPayload>): RouteHandler {
return (request, match) => {
void request;
const queueResult = extractQueueName(match);
Expand All @@ -113,7 +117,7 @@ function peekHandler(mgr: QueueManager<string>): RouteHandler {
};
}

function lengthHandler(mgr: QueueManager<string>): RouteHandler {
function lengthHandler(mgr: QueueManager<JsonPayload>): RouteHandler {
return (request, match) => {
void request;
const queueResult = extractQueueName(match);
Expand All @@ -129,7 +133,7 @@ function lengthHandler(mgr: QueueManager<string>): RouteHandler {
};
}

function registerRoutes(router: Router, mgr: QueueManager<string>): void {
function registerRoutes(router: Router, mgr: QueueManager<JsonPayload>): void {
router.get("/health{/}?", () => {
return new Response(JSON.stringify({ status: "ok" }), {
status: 200,
Expand All @@ -153,7 +157,7 @@ function writeLog(destination: { writeSync(data: Uint8Array): number }, message:
}

export function createHandler(
mgr: QueueManager<string>,
mgr: QueueManager<JsonPayload>,
apiToken: string,
rateLimitRequests?: number,
) {
Expand Down
42 changes: 28 additions & 14 deletions src/payload.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,13 @@
export type JsonValue =
| string
| number
| boolean
| null
| JsonValue[]
| { [key: string]: JsonValue };

export type Payload = Exclude<JsonValue, null>;

export class InvalidPayloadError extends Error {
constructor(message: string = "Invalid JSON") {
super(message);
Expand Down Expand Up @@ -101,11 +111,12 @@ function decodePayloadBody(body: string | Uint8Array): string {
}
}

function parseAndValidateJson(text: string): unknown {
function parseAndValidateJson(text: string): JsonValue {
try {
const json = JSON.parse(text);
const json: unknown = JSON.parse(text);
validateJsonSource(text);
return json;
// JSON.parse only creates JSON values; source validation also rules out non-finite numbers.
return json as JsonValue;
} catch (error) {
if (error instanceof RangeError || error instanceof JsonNestingTooDeepError) {
throw new JsonNestingTooDeepError();
Expand All @@ -117,32 +128,36 @@ function parseAndValidateJson(text: string): unknown {
}
}

function extractPayloadValue<T>(json: unknown): T {
if (json === null || typeof json !== "object") {
function isJsonObject(value: JsonValue): value is { [key: string]: JsonValue } {
return typeof value === "object" && value !== null && !Array.isArray(value);
}

function extractPayloadValue(json: JsonValue): Payload {
if (!isJsonObject(json)) {
throw new InvalidPayloadError("Missing payload key");
}
if (!("payload" in json)) {
throw new InvalidPayloadError("Missing payload key");
}
const payload = (json as Record<string, unknown>).payload;
const payload = json.payload;
if (payload === null) {
throw new InvalidPayloadError("Null payload not allowed");
}
return payload as T;
return payload;
}

export function parsePayloadBody<T = unknown>(body: string | Uint8Array): T {
export function parsePayloadBody(body: string | Uint8Array): Payload {
const text = decodePayloadBody(body);
const json = parseAndValidateJson(text);
return extractPayloadValue<T>(json);
return extractPayloadValue(json);
}

export async function readAndValidatePayload<T = unknown>(
export async function readAndValidatePayload(
stream: ReadableStream<Uint8Array> | null,
maxBytes: number = DEFAULT_MAX_PAYLOAD_SIZE,
): Promise<T> {
): Promise<Payload> {
if (stream === null) {
return parsePayloadBody<T>("");
return parsePayloadBody("");
}

const reader = stream.getReader();
Expand Down Expand Up @@ -182,6 +197,5 @@ export async function readAndValidatePayload<T = unknown>(
body.set(chunk, offset);
offset += chunk.byteLength;
}
return parsePayloadBody<T>(body);
return parsePayloadBody(body);
}

3 changes: 2 additions & 1 deletion tests/e2e_test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -4,13 +4,14 @@ import { createHandler } from "../src/handler.ts";
import * as Persistency from "../src/persist.ts";
import { RateLimiter } from "../src/rate_limiter.ts";
import { parseConfig, ConfigError } from "../src/config.ts";
import type { Payload } from "../src/payload.ts";

// Shared helpers
const API_TOKEN = "test-token";
const authHeaders = { "Authorization": `Bearer ${API_TOKEN}` };

function makeHandler(token = API_TOKEN, rateLimit = 100) {
const mgr = new QueueManager(new Persistency.MemoryStore());
const mgr = new QueueManager<Payload>(new Persistency.MemoryStore<Payload>());
return createHandler(mgr, token, rateLimit);
}

Expand Down
Loading
Loading