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
4 changes: 4 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,3 +25,7 @@ npm test

Demo baselines are intentionally red. Do not fix a defect directly on its
baseline branch; the repair should arrive through a RepoPilot pull request.

The webhook replay repair must make task creation and dispatch one idempotent
operation. Concurrent retries with the same delivery ID share one in-flight
creation, while different delivery IDs remain independent.
29 changes: 29 additions & 0 deletions src/webhooks/processor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -40,4 +40,33 @@ describe("IssueWebhookProcessor", () => {
expect(store.size()).toBe(1);
expect(dispatchTask).toHaveBeenCalledOnce();
});

it("returns the existing task for a sequential retry", async () => {
const dispatchTask = vi.fn(async () => undefined);
const store = new DeliveryTaskStore();
const processor = new IssueWebhookProcessor(store, dispatchTask);

const first = await processor.process(delivery);
const retry = await processor.process(delivery);

expect(retry).toEqual({ task: first.task, newlyCreated: false });
expect(dispatchTask).toHaveBeenCalledOnce();
});

it("processes different delivery IDs independently", async () => {
const dispatchTask = vi.fn(async () => undefined);
const store = new DeliveryTaskStore();
const processor = new IssueWebhookProcessor(store, dispatchTask);

const [first, second] = await Promise.all([
processor.process(delivery),
processor.process({ ...delivery, deliveryId: "delivery-issue-43", issueNumber: 43 })
]);

expect(first.task.taskId).not.toBe(second.task.taskId);
expect(first.newlyCreated).toBe(true);
expect(second.newlyCreated).toBe(true);
expect(store.size()).toBe(2);
expect(dispatchTask).toHaveBeenCalledTimes(2);
});
});
27 changes: 11 additions & 16 deletions src/webhooks/processor.ts
Original file line number Diff line number Diff line change
Expand Up @@ -15,21 +15,16 @@ export class IssueWebhookProcessor {
) {}

async process(delivery: IssueOpenedDelivery): Promise<ProcessIssueResult> {
const existing = await this.store.find(delivery.deliveryId);
if (existing) {
return { task: existing, newlyCreated: false };
}

await Promise.resolve();

const task: MaintenanceTask = {
taskId: randomUUID(),
deliveryId: delivery.deliveryId,
repository: delivery.repository,
issueNumber: delivery.issueNumber
};
await this.store.save(task);
await this.dispatchTask(task);
return { task, newlyCreated: true };
return this.store.getOrCreate(delivery.deliveryId, async () => {
await Promise.resolve();
const task: MaintenanceTask = {
taskId: randomUUID(),
deliveryId: delivery.deliveryId,
repository: delivery.repository,
issueNumber: delivery.issueNumber
};
await this.dispatchTask(task);
return task;
});
}
}
32 changes: 26 additions & 6 deletions src/webhooks/store.ts
Original file line number Diff line number Diff line change
@@ -1,14 +1,34 @@
import type { MaintenanceTask } from "./types.js";
import type { MaintenanceTask, StoredTaskResult } from "./types.js";

export class DeliveryTaskStore {
private readonly tasks = new Map<string, MaintenanceTask>();
private readonly inFlight = new Map<string, Promise<StoredTaskResult>>();

async find(deliveryId: string): Promise<MaintenanceTask | undefined> {
return this.tasks.get(deliveryId);
}
async getOrCreate(
deliveryId: string,
createTask: () => Promise<MaintenanceTask>
): Promise<StoredTaskResult> {
const existing = this.tasks.get(deliveryId);
if (existing) {
return { task: existing, newlyCreated: false };
}

const pending = this.inFlight.get(deliveryId);
if (pending) {
const result = await pending;
return { task: result.task, newlyCreated: false };
}

async save(task: MaintenanceTask): Promise<void> {
this.tasks.set(task.deliveryId, task);
const creation = createTask()
.then((task) => {
this.tasks.set(deliveryId, task);
return { task, newlyCreated: true };
})
.finally(() => {
this.inFlight.delete(deliveryId);
});
this.inFlight.set(deliveryId, creation);
return creation;
}

size(): number {
Expand Down
5 changes: 5 additions & 0 deletions src/webhooks/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -17,4 +17,9 @@ export interface ProcessIssueResult {
newlyCreated: boolean;
}

export interface StoredTaskResult {
task: MaintenanceTask;
newlyCreated: boolean;
}

export type DispatchTask = (task: MaintenanceTask) => Promise<void>;
Loading