Skip to content
Open
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
1 change: 1 addition & 0 deletions CONTEXT.md
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@
## Glossary

- **Queue**: A named, ordered sequence of items (FIFO data structure).
- **Queue Name**: The identifier of a Queue. A valid Queue Name is non-empty and at most 128 Unicode code points long; over HTTP it arrives percent-encoded and is decoded exactly once. Rules live in `src/queue_name.ts`.
- **Payload**: The arbitrary data object placed onto a Queue.
- **Queue Manager**: The central coordinator that tracks the lifecycles of all Queues and persists their state.
- **Persist Engine**: The storage mechanism for Queue state. Currently modeled as an append-only log.
Expand Down
8 changes: 4 additions & 4 deletions openapi.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -84,7 +84,7 @@ paths:
schema:
type: string
'400':
description: Bad Request - Invalid JSON (including invalid UTF-8), unsupported number, or Queue name too long
description: Bad Request - Invalid JSON (including invalid UTF-8), unsupported number, Invalid queue name (malformed percent-encoding), or Queue name too long
'401':
description: Unauthorized
'413':
Expand Down Expand Up @@ -115,7 +115,7 @@ paths:
'204':
description: No Content - Queue is empty
'400':
description: Bad Request - Queue name too long
description: Bad Request - Invalid queue name (malformed percent-encoding) or Queue name too long
'401':
description: Unauthorized
'429':
Expand All @@ -142,7 +142,7 @@ paths:
'204':
description: No Content - Queue is empty
'400':
description: Bad Request - Queue name too long
description: Bad Request - Invalid queue name (malformed percent-encoding) or Queue name too long
'401':
description: Unauthorized
'429':
Expand All @@ -167,7 +167,7 @@ paths:
schema:
type: integer
'400':
description: Bad Request - Queue name too long
description: Bad Request - Invalid queue name (malformed percent-encoding) or Queue name too long
'401':
description: Unauthorized
'429':
Expand Down
1 change: 1 addition & 0 deletions scripts/messcript.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ const productionUnits = new Map([
["queue-manager", ["src/manager.ts"]],
["persist-engine", ["src/persist.ts"]],
["payload", ["src/payload.ts"]],
["queue-name", ["src/queue_name.ts"]],
["http-handler", ["src/handler.ts"]],
["entrypoint", ["main.ts"]],
]);
Expand Down
107 changes: 36 additions & 71 deletions src/handler.ts
Original file line number Diff line number Diff line change
@@ -1,37 +1,41 @@
import QueueManager, { QueueNameTooLongError } from "./manager.ts";
import QueueManager from "./manager.ts";
import { RateLimiter } from "./rate_limiter.ts";
import { withAuth, withRateLimit } from "./middleware.ts";
import { Router } from "./router.ts";
import * as Payload from "./payload.ts";
import * as QueueName from "./queue_name.ts";

type JsonPayload = Payload.Payload;
type RouteHandler = Parameters<Router["get"]>[1];
type RouteMatch = Parameters<RouteHandler>[1];
type QueueRouteHandler = (queueName: string, request: Request) => Response | Promise<Response>;

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

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 }) };
function queueNameErrorResponse(error: unknown): Response {
if (error instanceof QueueName.InvalidQueueNameError) {
return new Response("Invalid queue name", { status: 400 });
}
try {
return { name: decodeURIComponent(raw) };
} catch (error) {
if (error instanceof URIError) {
return { error: new Response("Invalid queue name", { status: 400 }) };
}
throw error;
if (error instanceof QueueName.QueueNameTooLongError) {
return new Response("Queue name too long", { status: 400 });
}
throw error;
}

function enqueueHandler(mgr: QueueManager<JsonPayload>): RouteHandler {
function queueRoute(
handle: QueueRouteHandler,
parseName: (raw: string | undefined) => string = QueueName.parseQueueName,
): RouteHandler {
return async (request, match) => {
const queueResult = extractQueueName(match);
if ("error" in queueResult) {
return queueResult.error;
try {
return await handle(parseName(match.pathname.groups.queue), request);
} catch (error) {
return queueNameErrorResponse(error);
}
const queueName = queueResult.name;
};
}

function enqueueHandler(mgr: QueueManager<JsonPayload>): QueueRouteHandler {
return async (queueName, request) => {
try {
const contentLength = request.headers.get("content-length");
if (contentLength && parseInt(contentLength) > Payload.DEFAULT_MAX_PAYLOAD_SIZE) {
Expand Down Expand Up @@ -65,13 +69,6 @@ function enqueueErrorResponse(error: unknown): Response {
if (error instanceof Payload.InvalidPayloadError) {
return new Response(error.message, { status: 400 });
}
return queueNameErrorResponse(error);
}

function queueNameErrorResponse(error: unknown): Response {
if (error instanceof QueueNameTooLongError) {
return new Response("Queue name too long", { status: 400 });
}
throw error;
}

Expand All @@ -84,53 +81,19 @@ function itemResponse(item: JsonPayload | undefined): Response {
});
}

function dequeueHandler(mgr: QueueManager<JsonPayload>): RouteHandler {
return (request, match) => {
void request;
const queueResult = extractQueueName(match);
if ("error" in queueResult) {
return queueResult.error;
}
try {
const item = request.method === "HEAD"
? mgr.peek(queueResult.name)
: mgr.dequeue(queueResult.name);
return itemResponse(item);
} catch (error) {
return queueNameErrorResponse(error);
}
function dequeueHandler(mgr: QueueManager<JsonPayload>): QueueRouteHandler {
return (queueName, request) => {
const item = request.method === "HEAD" ? mgr.peek(queueName) : mgr.dequeue(queueName);
return itemResponse(item);
};
}

function peekHandler(mgr: QueueManager<JsonPayload>): RouteHandler {
return (request, match) => {
void request;
const queueResult = extractQueueName(match);
if ("error" in queueResult) {
return queueResult.error;
}
try {
return itemResponse(mgr.peek(queueResult.name));
} catch (error) {
return queueNameErrorResponse(error);
}
};
function peekHandler(mgr: QueueManager<JsonPayload>): QueueRouteHandler {
return (queueName) => itemResponse(mgr.peek(queueName));
}

function lengthHandler(mgr: QueueManager<JsonPayload>): RouteHandler {
return (request, match) => {
void request;
const queueResult = extractQueueName(match);
if ("error" in queueResult) {
return queueResult.error;
}
try {
const length = mgr.length(queueResult.name);
return new Response(`${length}`);
} catch (error) {
return queueNameErrorResponse(error);
}
};
function lengthHandler(mgr: QueueManager<JsonPayload>): QueueRouteHandler {
return (queueName) => new Response(`${mgr.length(queueName)}`);
}

function registerRoutes(router: Router, mgr: QueueManager<JsonPayload>): void {
Expand All @@ -146,10 +109,12 @@ function registerRoutes(router: Router, mgr: QueueManager<JsonPayload>): void {
headers: { "Content-Type": "application/json" },
});
});
router.post("/enqueue/:queue", enqueueHandler(mgr));
router.get("/dequeue/:queue", dequeueHandler(mgr));
router.get("/peek/:queue", peekHandler(mgr));
router.get("/length/:queue", lengthHandler(mgr));
// Enqueue only decodes here; QueueManager applies the length rule after the
// payload is validated, so payload errors keep precedence over it.
router.post("/enqueue/:queue", queueRoute(enqueueHandler(mgr), QueueName.decodeQueueName));
router.get("/dequeue/:queue", queueRoute(dequeueHandler(mgr)));
router.get("/peek/:queue", queueRoute(peekHandler(mgr)));
router.get("/length/:queue", queueRoute(lengthHandler(mgr)));
}

function writeLog(destination: { writeSync(data: Uint8Array): number }, message: string): void {
Expand Down
26 changes: 7 additions & 19 deletions src/manager.ts
Original file line number Diff line number Diff line change
@@ -1,12 +1,6 @@
import { QueueEvent, QueueStore } from "./persist.ts"
export const MAX_QUEUE_NAME_LENGTH = 128;

export class QueueNameTooLongError extends Error {
constructor() {
super("Queue name too long");
this.name = "QueueNameTooLongError";
}
}
import { validateQueueName } from "./queue_name.ts";
export { MAX_QUEUE_NAME_LENGTH, QueueNameTooLongError } from "./queue_name.ts";

/**
* FIFO queue with O(1) amortized enqueue and dequeue.
Expand Down Expand Up @@ -77,18 +71,12 @@ export default class Manager<T = string> {
return this;
}

private validateName(name: string): void {
if (Array.from(name).length > MAX_QUEUE_NAME_LENGTH) {
throw new QueueNameTooLongError();
}
}

public canCreateQueue(): boolean {
return this.queues.size < this.queueCountLimit;
}

public canEnqueue(name: string): boolean {
this.validateName(name);
validateQueueName(name);
const queue = this.find(name);
if (!queue) {
return this.canCreateQueue() && 0 < this.queueDepthLimit;
Expand All @@ -101,7 +89,7 @@ export default class Manager<T = string> {
}

public enqueue(name: string, payload: T): Manager<T> {
this.validateName(name);
validateQueueName(name);
const existing = this.find(name);
if (!existing && !this.canCreateQueue()) {
throw new Error("Queue count limit reached");
Expand All @@ -123,7 +111,7 @@ export default class Manager<T = string> {
}

public dequeue(name: string): T | undefined {
this.validateName(name);
validateQueueName(name);
const queue = this.find(name);
if (!queue) {
return undefined;
Expand All @@ -144,7 +132,7 @@ export default class Manager<T = string> {
}

public peek(name: string): T | undefined {
this.validateName(name);
validateQueueName(name);
const queue = this.find(name);

if (queue === undefined) {
Expand All @@ -155,7 +143,7 @@ export default class Manager<T = string> {
}

public length(name: string): number {
this.validateName(name);
validateQueueName(name);
const queue = this.find(name);
return queue ? queue.length : 0;
}
Expand Down
55 changes: 55 additions & 0 deletions src/queue_name.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
export const MAX_QUEUE_NAME_LENGTH = 128;

export class InvalidQueueNameError extends Error {
constructor(message: string = "Invalid queue name") {
super(message);
this.name = "InvalidQueueNameError";
}
}

export class QueueNameTooLongError extends Error {
constructor(message: string = "Queue name too long") {
super(message);
this.name = "QueueNameTooLongError";
}
}

/**
* Checks an already-decoded Queue name against the domain rules: it must be
* non-empty and at most MAX_QUEUE_NAME_LENGTH Unicode code points long.
*/
export function validateQueueName(name: string): string {
if (name === "") {
throw new InvalidQueueNameError();
}
if (Array.from(name).length > MAX_QUEUE_NAME_LENGTH) {
throw new QueueNameTooLongError();
}
return name;
}

/**
* Percent-decodes a raw Queue name without applying the length rule.
* Callers that decode with this must validate the name before use.
*/
export function decodeQueueName(raw: string | undefined): string {
if (raw === undefined) {
throw new InvalidQueueNameError();
}
try {
return decodeURIComponent(raw);
} catch (error) {
if (error instanceof URIError) {
throw new InvalidQueueNameError();
}
throw error;
}
}

/**
* Parses a raw, percent-encoded Queue name (e.g. a URL path segment) into a
* validated Queue name.
*/
export function parseQueueName(raw: string | undefined): string {
return validateQueueName(decodeQueueName(raw));
}
Loading
Loading