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
49 changes: 29 additions & 20 deletions src/poller.ts
Original file line number Diff line number Diff line change
Expand Up @@ -161,10 +161,11 @@ async function processIssues(
repo: string,
env: Env,
storeStub: DurableObjectStub,
): Promise<{ processed: number; embedded: number; skipped: number }> {
): Promise<{ processed: number; embedded: number; skipped: number; failed: number }> {
let processed = 0;
let embedded = 0;
let skipped = 0;
let failed = 0;

// Process in batches to manage memory and rate limits
for (const issue of issues) {
Expand All @@ -176,20 +177,6 @@ async function processIssues(
? "pull_request"
: "issue";

const record: IssueRecord = {
repo,
number: issue.number,
type,
state: issue.state,
title,
labels: issue.labels.map((l) => l.name),
milestone: issue.milestone?.title ?? "",
assignees: issue.assignees.map((a) => a.login),
bodyHash,
createdAt: issue.created_at,
updatedAt: issue.updated_at,
};

// Check if body has changed by comparing hash with stored value
const existingResp = await storeStub.fetch(
new Request(
Expand All @@ -206,6 +193,9 @@ async function processIssues(
}
}

// Track whether embedding succeeded — determines whether bodyHash is saved
let embeddingSucceeded = false;

// Generate embedding if content changed
if (needsEmbedding) {
try {
Expand All @@ -217,9 +207,9 @@ async function processIssues(
number: issue.number,
type,
state: issue.state,
labels: record.labels.join(","),
milestone: record.milestone,
assignees: record.assignees.join(","),
labels: issue.labels.map((l) => l.name).join(","),
milestone: issue.milestone?.title ?? "",
assignees: issue.assignees.map((a) => a.login).join(","),
updated_at: issue.updated_at,
};

Expand All @@ -232,16 +222,35 @@ async function processIssues(
},
]);

embeddingSucceeded = true;
embedded++;
} catch (err) {
console.error(
`Failed to embed ${repo}#${issue.number}:`,
err instanceof Error ? err.message : String(err),
);
failed++;
// Continue processing other issues even if one fails
}
}

// Build record — save bodyHash only when embedding succeeded (or was skipped
// because it already exists). When embedding fails, store empty bodyHash so
// the next poll will detect a mismatch and retry embedding.
const record: IssueRecord = {
repo,
number: issue.number,
type,
state: issue.state,
title,
labels: issue.labels.map((l) => l.name),
milestone: issue.milestone?.title ?? "",
assignees: issue.assignees.map((a) => a.login),
bodyHash: needsEmbedding && !embeddingSucceeded ? "" : bodyHash,
createdAt: issue.created_at,
updatedAt: issue.updated_at,
};

// Upsert structured data into IssueStore (always, even if embedding skipped)
await storeStub.fetch(
new Request("http://store/upsert", {
Expand All @@ -254,7 +263,7 @@ async function processIssues(
processed++;
}

return { processed, embedded, skipped };
return { processed, embedded, skipped, failed };
}

/**
Expand Down Expand Up @@ -314,7 +323,7 @@ async function pollRepo(
);

console.log(
`${repo}: ${stats.processed} processed, ${stats.embedded} embedded, ${stats.skipped} unchanged`,
`${repo}: ${stats.processed} processed, ${stats.embedded} embedded, ${stats.skipped} unchanged, ${stats.failed} failed`,
);
}

Expand Down
23 changes: 23 additions & 0 deletions src/store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -179,6 +179,21 @@ export class IssueStore implements DurableObject {
return [...cursor].map(rowToIssueRecord);
}

// ---- Hash reset for re-sync ----

/**
* Reset all bodyHashes for a given repo so that the next poll
* will regenerate embeddings for every issue.
* Returns the number of rows affected.
*/
resetBodyHashes(repo: string): number {
const cursor = this.sql.exec(
`UPDATE issues SET body_hash = '' WHERE repo = ? AND body_hash != ''`,
repo,
);
return cursor.rowsWritten;
}

// ---- Watermark management ----

getWatermark(repo: string): PollWatermark | null {
Expand Down Expand Up @@ -258,6 +273,14 @@ export class IssueStore implements DurableObject {
return Response.json(items);
}

// POST /reset-hashes?repo=... — reset all bodyHashes for a repo to force re-embedding
if (request.method === "POST" && path === "/reset-hashes") {
const repo = url.searchParams.get("repo");
if (!repo) return new Response("missing repo", { status: 400 });
const count = this.resetBodyHashes(repo);
return Response.json({ repo, reset: count });
}

// GET /watermark?repo=... — get poll watermark
if (request.method === "GET" && path === "/watermark") {
const repo = url.searchParams.get("repo");
Expand Down
Loading