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
25 changes: 23 additions & 2 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1040,14 +1040,26 @@ async function syncGmailChannelFull(
export async function syncGmailMailboxIncremental(
api: GmailApi,
historyId: string,
retryThreadIds: string[] = []
retryThreadIds: string[] = [],
maxThreads: number = Infinity
): Promise<
| { expired: true }
| {
expired: false;
historyId: string;
threads: GmailThread[];
failedThreadIds: string[];
/**
* Thread ids that changed in this history window but were NOT fetched
* this pass because the per-pass `maxThreads` budget was reached. The
* caller carries these forward (see {@link mergePendingThreads}) and
* schedules a continuation to drain them. Unbounded fetching here is what
* let a large window (e.g. a cursor reseed after the Google re-home) load
* thousands of full threads into one isolate and exceed the Worker memory
* limit, which then tore down the in-flight DB connection mid-save
* ("driver has already been destroyed").
*/
deferredThreadIds: string[];
}
> {
let historyResult;
Expand DownExpand Up@@ -1075,9 +1087,17 @@ export async function syncGmailMailboxIncremental(
}
}

// Bound how many full threads we pull into memory per pass. `retryThreadIds`
// are inserted into the Set first, so prior-deferred (and previously-failed)
// threads sit at the front of iteration order and drain ahead of newly
// changed ones. Everything past the cap is returned as `deferredThreadIds`.
const ordered = [...changedThreadIds];
const toFetch = ordered.slice(0, maxThreads);
const deferredThreadIds = ordered.slice(toFetch.length);

const threads: GmailThread[] = [];
const failedThreadIds: string[] = [];
for (const threadId of changedThreadIds) {
for (const threadId of toFetch) {
try {
threads.push(await api.getThread(threadId));
} catch (error) {
Expand All@@ -1091,6 +1111,7 @@ export async function syncGmailMailboxIncremental(
historyId: historyResult.historyId,
threads,
failedThreadIds,
deferredThreadIds,
};
}

Expand Down
240 changes: 240 additions & 0 deletions connectors/gmail/src/gmail-incremental-bound.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
GmailApi,
type GmailThread,
syncGmailMailboxIncremental,
} from "./gmail-api";
import {
type GmailSyncHost,
type IncrementalState,
MAX_INCREMENTAL_THREADS_PER_BATCH,
MAX_THREAD_FETCH_ATTEMPTS,
incrementalSyncBatchFn,
mergePendingThreads,
} from "./sync";

/** Build a minimal GmailApi mock exposing only the methods the incremental
* sync touches (getHistory + getThread). */
function mockApi(opts: {
changedThreadIds: string[];
newHistoryId: string;
failOn?: Set<string>;
}): { api: GmailApi; getThread: ReturnType<typeof vi.fn> } {
const getThread = vi.fn(async (id: string): Promise<GmailThread> => {
if (opts.failOn?.has(id)) throw new Error(`boom ${id}`);
return { id, historyId: "h", messages: [] } as unknown as GmailThread;
});
const getHistory = vi.fn(async () => ({
history: opts.changedThreadIds.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } }],
})),
historyId: opts.newHistoryId,
}));
const api = { getHistory, getThread } as unknown as GmailApi;
return { api, getThread };
}

describe("syncGmailMailboxIncremental — per-pass bound", () => {
it("fetches at most maxThreads and defers the rest", async () => {
const ids = ["t1", "t2", "t3", "t4", "t5"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "200",
});

const result = await syncGmailMailboxIncremental(api, "100", [], 2);

expect("expired" in result && result.expired).toBe(false);
if ("expired" in result && result.expired) return;

// Only the cap is fetched into memory this pass.
expect(getThread).toHaveBeenCalledTimes(2);
expect(result.threads.map((t) => t.id)).toEqual(["t1", "t2"]);
// The overflow is deferred, not fetched and not failed.
expect(result.deferredThreadIds).toEqual(["t3", "t4", "t5"]);
expect(result.failedThreadIds).toEqual([]);
expect(result.historyId).toBe("200");
});

it("processes prior retry (deferred) ids ahead of newly-changed ones", async () => {
// retryThreadIds are inserted first, so they fall within the cap before
// freshly-changed threads — prior backlog drains first.
const { api, getThread } = mockApi({
changedThreadIds: ["new1", "new2"],
newHistoryId: "201",
});

const result = await syncGmailMailboxIncremental(
api,
"100",
["retryA", "retryB"],
2
);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread.mock.calls.map((c) => c[0])).toEqual(["retryA", "retryB"]);
expect(result.deferredThreadIds).toEqual(["new1", "new2"]);
});

it("does not bound when maxThreads is omitted (back-compat)", async () => {
const ids = ["a", "b", "c"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "300",
});
const result = await syncGmailMailboxIncremental(api, "100", []);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread).toHaveBeenCalledTimes(3);
expect(result.deferredThreadIds).toEqual([]);
});
});

/** Minimal GmailSyncHost backed by an in-memory store, exposing spies for the
* scheduler continuation hook and the saved incremental cursor. */
function makeHost(initial: IncrementalState): {
host: GmailSyncHost;
store: Map<string, unknown>;
queueIncrementalSync: ReturnType<typeof vi.fn>;
} {
const store = new Map<string, unknown>([
["enabled_channels", ["INBOX"]],
["incremental_state", initial],
]);
const queueIncrementalSync = vi.fn(async () => {});
const host = {
id: "twist-instance-1",
get: vi.fn(async (key: string) =>
store.has(key) ? store.get(key) : null
),
set: vi.fn(async (key: string, value: unknown) => {
store.set(key, value);
}),
clear: vi.fn(async (key: string) => {
store.delete(key);
}),
tools: {
integrations: {
get: vi.fn(async () => ({ token: "tok", scopes: [] })),
// Threads in these tests carry no notes, so saveLink is never reached;
// present only to satisfy the interface.
saveLink: vi.fn(async () => null),
channelSyncCompleted: vi.fn(async () => {}),
setThreadToDo: vi.fn(async () => {}),
},
files: { read: vi.fn() },
network: { createWebhook: vi.fn(), deleteWebhook: vi.fn() },
store: {
acquireLock: vi.fn(async () => true),
releaseLock: vi.fn(async () => {}),
list: vi.fn(async () => []),
},
},
scheduler: {
onGmailWebhook: undefined,
setupMailboxWebhook: vi.fn(async () => {}),
renewMailboxWatch: vi.fn(async () => {}),
scheduleMailboxRenewal: vi.fn(async () => {}),
scheduleSelfHealCheck: vi.fn(async () => {}),
cancelScheduledTask: vi.fn(async () => {}),
queueIncrementalSync,
},
} as unknown as GmailSyncHost;
return { host, store, queueIncrementalSync };
}

describe("incrementalSyncBatchFn — bounded pass + continuation", () => {
afterEach(() => vi.restoreAllMocks());

it("caps the pass, carries the overflow, and queues a continuation", async () => {
const overflow = 5;
const total = MAX_INCREMENTAL_THREADS_PER_BATCH + overflow;
const ids = Array.from({ length: total }, (_, i) => `t${i}`);

const getHistory = vi
.spyOn(GmailApi.prototype, "getHistory")
.mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "999",
} as any);
const getThread = vi
.spyOn(GmailApi.prototype, "getThread")
.mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({
historyId: "100",
});

await incrementalSyncBatchFn(host);

expect(getHistory).toHaveBeenCalledTimes(1);
// Only the cap is pulled into memory this pass — not all 25.
expect(getThread).toHaveBeenCalledTimes(MAX_INCREMENTAL_THREADS_PER_BATCH);
// The overflow is carried forward (attempts 0 — never attempted) and the
// cursor advanced so we don't re-walk the window.
const saved = store.get("incremental_state") as IncrementalState;
expect(saved.historyId).toBe("999");
expect(saved.pendingThreadIds).toHaveLength(overflow);
expect(saved.pendingThreadIds?.every((p) => p.attempts === 0)).toBe(true);
// A continuation is scheduled to drain the rest.
expect(queueIncrementalSync).toHaveBeenCalledTimes(1);
});

it("does not queue a continuation when everything fits in one pass", async () => {
const ids = ["a", "b"];
vi.spyOn(GmailApi.prototype, "getHistory").mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "201",
} as any);
vi.spyOn(GmailApi.prototype, "getThread").mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({ historyId: "100" });
await incrementalSyncBatchFn(host);

const saved = store.get("incremental_state") as IncrementalState;
expect(saved.pendingThreadIds).toEqual([]);
expect(queueIncrementalSync).not.toHaveBeenCalled();
});
});

describe("mergePendingThreads — deferred carry", () => {
it("carries deferred ids without bumping their attempt counter", () => {
const prior = [{ id: "d1", attempts: 0 }];
const merged = mergePendingThreads(prior, [], ["d1", "d2"]);
// Neither deferred id is a fetch attempt, so attempts stay put — a large
// backlog must not be abandoned just for waiting its turn.
expect(merged).toEqual([
{ id: "d1", attempts: 0 },
{ id: "d2", attempts: 0 },
]);
});

it("still bumps and eventually drops genuinely-failed fetches", () => {
const prior = [{ id: "f1", attempts: MAX_THREAD_FETCH_ATTEMPTS }];
// f1 has exhausted its attempts → dropped; f2 is a fresh failure → kept@1.
const merged = mergePendingThreads(prior, ["f1", "f2"], []);
expect(merged).toEqual([{ id: "f2", attempts: 1 }]);
});

it("keeps failed and deferred sets distinct in one merge", () => {
const merged = mergePendingThreads([], ["f1"], ["d1"]);
expect(merged).toEqual([
{ id: "f1", attempts: 1 },
{ id: "d1", attempts: 0 },
]);
});
});
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Add copy buttons to all
 blocks\n(function() {\n function addCopyButtons() {\n document.querySelectorAll('pre code').forEach(function(codeBlock) {\n if (codeBlock.parentElement.hasAttribute('data-copy-added')) return;\n codeBlock.parentElement.setAttribute('data-copy-added', 'true');\n \n var btn = document.createElement('button');\n btn.textContent = 'Copy';\n btn.style.cssText = 'position:absolute;top:4px;right:4px;padding:2px 8px;font-size:11px;background:#4ecdc4;border:none;border-radius:4px;color:#1a1a2e;cursor:pointer;opacity:0.7;transition:opacity 0.2s;';\n btn.onmouseover = function() { this.style.opacity = '1'; };\n btn.onmouseout = function() { this.style.opacity = '0.7'; };\n btn.onclick = function() {\n navigator.clipboard.writeText(codeBlock.textContent).then(function() {\n btn.textContent = 'Copied!';\n setTimeout(function() { btn.textContent = 'Copy'; }, 1500);\n });\n };\n codeBlock.parentElement.style.position = 'relative';\n codeBlock.parentElement.appendChild(btn);\n });\n }\n \n addCopyButtons();\n \n // Re-run on dynamic content\n var observer = new MutationObserver(addCopyButtons);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Add Copy Buttons to Code Blocks");
}
} catch(__e) { console.warn('[Userscript:Add Copy Buttons to Code Blocks]', __e); }
})();
(function(){
try {
var __m = "github.com";
var __re = new RegExp('^' + "github\\.com" + '
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
25 changes: 23 additions & 2 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1040,14 +1040,26 @@ async function syncGmailChannelFull(
export async function syncGmailMailboxIncremental(
api: GmailApi,
historyId: string,
retryThreadIds: string[] = []
retryThreadIds: string[] = [],
maxThreads: number = Infinity
): Promise<
| { expired: true }
| {
expired: false;
historyId: string;
threads: GmailThread[];
failedThreadIds: string[];
/**
* Thread ids that changed in this history window but were NOT fetched
* this pass because the per-pass `maxThreads` budget was reached. The
* caller carries these forward (see {@link mergePendingThreads}) and
* schedules a continuation to drain them. Unbounded fetching here is what
* let a large window (e.g. a cursor reseed after the Google re-home) load
* thousands of full threads into one isolate and exceed the Worker memory
* limit, which then tore down the in-flight DB connection mid-save
* ("driver has already been destroyed").
*/
deferredThreadIds: string[];
}
> {
let historyResult;
Expand DownExpand Up@@ -1075,9 +1087,17 @@ export async function syncGmailMailboxIncremental(
}
}

// Bound how many full threads we pull into memory per pass. `retryThreadIds`
// are inserted into the Set first, so prior-deferred (and previously-failed)
// threads sit at the front of iteration order and drain ahead of newly
// changed ones. Everything past the cap is returned as `deferredThreadIds`.
const ordered = [...changedThreadIds];
const toFetch = ordered.slice(0, maxThreads);
const deferredThreadIds = ordered.slice(toFetch.length);

const threads: GmailThread[] = [];
const failedThreadIds: string[] = [];
for (const threadId of changedThreadIds) {
for (const threadId of toFetch) {
try {
threads.push(await api.getThread(threadId));
} catch (error) {
Expand All@@ -1091,6 +1111,7 @@ export async function syncGmailMailboxIncremental(
historyId: historyResult.historyId,
threads,
failedThreadIds,
deferredThreadIds,
};
}

Expand Down
240 changes: 240 additions & 0 deletions connectors/gmail/src/gmail-incremental-bound.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
GmailApi,
type GmailThread,
syncGmailMailboxIncremental,
} from "./gmail-api";
import {
type GmailSyncHost,
type IncrementalState,
MAX_INCREMENTAL_THREADS_PER_BATCH,
MAX_THREAD_FETCH_ATTEMPTS,
incrementalSyncBatchFn,
mergePendingThreads,
} from "./sync";

/** Build a minimal GmailApi mock exposing only the methods the incremental
* sync touches (getHistory + getThread). */
function mockApi(opts: {
changedThreadIds: string[];
newHistoryId: string;
failOn?: Set<string>;
}): { api: GmailApi; getThread: ReturnType<typeof vi.fn> } {
const getThread = vi.fn(async (id: string): Promise<GmailThread> => {
if (opts.failOn?.has(id)) throw new Error(`boom ${id}`);
return { id, historyId: "h", messages: [] } as unknown as GmailThread;
});
const getHistory = vi.fn(async () => ({
history: opts.changedThreadIds.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } }],
})),
historyId: opts.newHistoryId,
}));
const api = { getHistory, getThread } as unknown as GmailApi;
return { api, getThread };
}

describe("syncGmailMailboxIncremental — per-pass bound", () => {
it("fetches at most maxThreads and defers the rest", async () => {
const ids = ["t1", "t2", "t3", "t4", "t5"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "200",
});

const result = await syncGmailMailboxIncremental(api, "100", [], 2);

expect("expired" in result && result.expired).toBe(false);
if ("expired" in result && result.expired) return;

// Only the cap is fetched into memory this pass.
expect(getThread).toHaveBeenCalledTimes(2);
expect(result.threads.map((t) => t.id)).toEqual(["t1", "t2"]);
// The overflow is deferred, not fetched and not failed.
expect(result.deferredThreadIds).toEqual(["t3", "t4", "t5"]);
expect(result.failedThreadIds).toEqual([]);
expect(result.historyId).toBe("200");
});

it("processes prior retry (deferred) ids ahead of newly-changed ones", async () => {
// retryThreadIds are inserted first, so they fall within the cap before
// freshly-changed threads — prior backlog drains first.
const { api, getThread } = mockApi({
changedThreadIds: ["new1", "new2"],
newHistoryId: "201",
});

const result = await syncGmailMailboxIncremental(
api,
"100",
["retryA", "retryB"],
2
);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread.mock.calls.map((c) => c[0])).toEqual(["retryA", "retryB"]);
expect(result.deferredThreadIds).toEqual(["new1", "new2"]);
});

it("does not bound when maxThreads is omitted (back-compat)", async () => {
const ids = ["a", "b", "c"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "300",
});
const result = await syncGmailMailboxIncremental(api, "100", []);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread).toHaveBeenCalledTimes(3);
expect(result.deferredThreadIds).toEqual([]);
});
});

/** Minimal GmailSyncHost backed by an in-memory store, exposing spies for the
* scheduler continuation hook and the saved incremental cursor. */
function makeHost(initial: IncrementalState): {
host: GmailSyncHost;
store: Map<string, unknown>;
queueIncrementalSync: ReturnType<typeof vi.fn>;
} {
const store = new Map<string, unknown>([
["enabled_channels", ["INBOX"]],
["incremental_state", initial],
]);
const queueIncrementalSync = vi.fn(async () => {});
const host = {
id: "twist-instance-1",
get: vi.fn(async (key: string) =>
store.has(key) ? store.get(key) : null
),
set: vi.fn(async (key: string, value: unknown) => {
store.set(key, value);
}),
clear: vi.fn(async (key: string) => {
store.delete(key);
}),
tools: {
integrations: {
get: vi.fn(async () => ({ token: "tok", scopes: [] })),
// Threads in these tests carry no notes, so saveLink is never reached;
// present only to satisfy the interface.
saveLink: vi.fn(async () => null),
channelSyncCompleted: vi.fn(async () => {}),
setThreadToDo: vi.fn(async () => {}),
},
files: { read: vi.fn() },
network: { createWebhook: vi.fn(), deleteWebhook: vi.fn() },
store: {
acquireLock: vi.fn(async () => true),
releaseLock: vi.fn(async () => {}),
list: vi.fn(async () => []),
},
},
scheduler: {
onGmailWebhook: undefined,
setupMailboxWebhook: vi.fn(async () => {}),
renewMailboxWatch: vi.fn(async () => {}),
scheduleMailboxRenewal: vi.fn(async () => {}),
scheduleSelfHealCheck: vi.fn(async () => {}),
cancelScheduledTask: vi.fn(async () => {}),
queueIncrementalSync,
},
} as unknown as GmailSyncHost;
return { host, store, queueIncrementalSync };
}

describe("incrementalSyncBatchFn — bounded pass + continuation", () => {
afterEach(() => vi.restoreAllMocks());

it("caps the pass, carries the overflow, and queues a continuation", async () => {
const overflow = 5;
const total = MAX_INCREMENTAL_THREADS_PER_BATCH + overflow;
const ids = Array.from({ length: total }, (_, i) => `t${i}`);

const getHistory = vi
.spyOn(GmailApi.prototype, "getHistory")
.mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "999",
} as any);
const getThread = vi
.spyOn(GmailApi.prototype, "getThread")
.mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({
historyId: "100",
});

await incrementalSyncBatchFn(host);

expect(getHistory).toHaveBeenCalledTimes(1);
// Only the cap is pulled into memory this pass — not all 25.
expect(getThread).toHaveBeenCalledTimes(MAX_INCREMENTAL_THREADS_PER_BATCH);
// The overflow is carried forward (attempts 0 — never attempted) and the
// cursor advanced so we don't re-walk the window.
const saved = store.get("incremental_state") as IncrementalState;
expect(saved.historyId).toBe("999");
expect(saved.pendingThreadIds).toHaveLength(overflow);
expect(saved.pendingThreadIds?.every((p) => p.attempts === 0)).toBe(true);
// A continuation is scheduled to drain the rest.
expect(queueIncrementalSync).toHaveBeenCalledTimes(1);
});

it("does not queue a continuation when everything fits in one pass", async () => {
const ids = ["a", "b"];
vi.spyOn(GmailApi.prototype, "getHistory").mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "201",
} as any);
vi.spyOn(GmailApi.prototype, "getThread").mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({ historyId: "100" });
await incrementalSyncBatchFn(host);

const saved = store.get("incremental_state") as IncrementalState;
expect(saved.pendingThreadIds).toEqual([]);
expect(queueIncrementalSync).not.toHaveBeenCalled();
});
});

describe("mergePendingThreads — deferred carry", () => {
it("carries deferred ids without bumping their attempt counter", () => {
const prior = [{ id: "d1", attempts: 0 }];
const merged = mergePendingThreads(prior, [], ["d1", "d2"]);
// Neither deferred id is a fetch attempt, so attempts stay put — a large
// backlog must not be abandoned just for waiting its turn.
expect(merged).toEqual([
{ id: "d1", attempts: 0 },
{ id: "d2", attempts: 0 },
]);
});

it("still bumps and eventually drops genuinely-failed fetches", () => {
const prior = [{ id: "f1", attempts: MAX_THREAD_FETCH_ATTEMPTS }];
// f1 has exhausted its attempts → dropped; f2 is a fresh failure → kept@1.
const merged = mergePendingThreads(prior, ["f1", "f2"], []);
expect(merged).toEqual([{ id: "f2", attempts: 1 }]);
});

it("keeps failed and deferred sets distinct in one merge", () => {
const merged = mergePendingThreads([], ["f1"], ["d1"]);
expect(merged).toEqual([
{ id: "f1", attempts: 1 },
{ id: "d1", attempts: 0 },
]);
});
});
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Force GitHub README to respect dark mode\n(function() {\n var style = document.createElement('style');\n style.textContent = '\n .markdown-body {\n color-scheme: dark light;\n }\n .markdown-body pre { background: #161b22 !important; }\n .markdown-body code { background: rgba(110, 118, 129, 0.4) !important; }\n .markdown-body table th, .markdown-body table td { border-color: #30363d !important; }\n .markdown-body img { background: #0d1117; }\n .markdown-body blockquote { border-left-color: #8b949e; }\n .markdown-body hr { border-color: #30363d; }\n ';\n document.head.appendChild(style);\n})();", "GitHub Dark Mode README Fix"); } } catch(__e) { console.warn('[Userscript:GitHub Dark Mode README Fix]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
25 changes: 23 additions & 2 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1040,14 +1040,26 @@ async function syncGmailChannelFull(
export async function syncGmailMailboxIncremental(
api: GmailApi,
historyId: string,
retryThreadIds: string[] = []
retryThreadIds: string[] = [],
maxThreads: number = Infinity
): Promise<
| { expired: true }
| {
expired: false;
historyId: string;
threads: GmailThread[];
failedThreadIds: string[];
/**
* Thread ids that changed in this history window but were NOT fetched
* this pass because the per-pass `maxThreads` budget was reached. The
* caller carries these forward (see {@link mergePendingThreads}) and
* schedules a continuation to drain them. Unbounded fetching here is what
* let a large window (e.g. a cursor reseed after the Google re-home) load
* thousands of full threads into one isolate and exceed the Worker memory
* limit, which then tore down the in-flight DB connection mid-save
* ("driver has already been destroyed").
*/
deferredThreadIds: string[];
}
> {
let historyResult;
Expand DownExpand Up@@ -1075,9 +1087,17 @@ export async function syncGmailMailboxIncremental(
}
}

// Bound how many full threads we pull into memory per pass. `retryThreadIds`
// are inserted into the Set first, so prior-deferred (and previously-failed)
// threads sit at the front of iteration order and drain ahead of newly
// changed ones. Everything past the cap is returned as `deferredThreadIds`.
const ordered = [...changedThreadIds];
const toFetch = ordered.slice(0, maxThreads);
const deferredThreadIds = ordered.slice(toFetch.length);

const threads: GmailThread[] = [];
const failedThreadIds: string[] = [];
for (const threadId of changedThreadIds) {
for (const threadId of toFetch) {
try {
threads.push(await api.getThread(threadId));
} catch (error) {
Expand All@@ -1091,6 +1111,7 @@ export async function syncGmailMailboxIncremental(
historyId: historyResult.historyId,
threads,
failedThreadIds,
deferredThreadIds,
};
}

Expand Down
240 changes: 240 additions & 0 deletions connectors/gmail/src/gmail-incremental-bound.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
GmailApi,
type GmailThread,
syncGmailMailboxIncremental,
} from "./gmail-api";
import {
type GmailSyncHost,
type IncrementalState,
MAX_INCREMENTAL_THREADS_PER_BATCH,
MAX_THREAD_FETCH_ATTEMPTS,
incrementalSyncBatchFn,
mergePendingThreads,
} from "./sync";

/** Build a minimal GmailApi mock exposing only the methods the incremental
* sync touches (getHistory + getThread). */
function mockApi(opts: {
changedThreadIds: string[];
newHistoryId: string;
failOn?: Set<string>;
}): { api: GmailApi; getThread: ReturnType<typeof vi.fn> } {
const getThread = vi.fn(async (id: string): Promise<GmailThread> => {
if (opts.failOn?.has(id)) throw new Error(`boom ${id}`);
return { id, historyId: "h", messages: [] } as unknown as GmailThread;
});
const getHistory = vi.fn(async () => ({
history: opts.changedThreadIds.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } }],
})),
historyId: opts.newHistoryId,
}));
const api = { getHistory, getThread } as unknown as GmailApi;
return { api, getThread };
}

describe("syncGmailMailboxIncremental — per-pass bound", () => {
it("fetches at most maxThreads and defers the rest", async () => {
const ids = ["t1", "t2", "t3", "t4", "t5"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "200",
});

const result = await syncGmailMailboxIncremental(api, "100", [], 2);

expect("expired" in result && result.expired).toBe(false);
if ("expired" in result && result.expired) return;

// Only the cap is fetched into memory this pass.
expect(getThread).toHaveBeenCalledTimes(2);
expect(result.threads.map((t) => t.id)).toEqual(["t1", "t2"]);
// The overflow is deferred, not fetched and not failed.
expect(result.deferredThreadIds).toEqual(["t3", "t4", "t5"]);
expect(result.failedThreadIds).toEqual([]);
expect(result.historyId).toBe("200");
});

it("processes prior retry (deferred) ids ahead of newly-changed ones", async () => {
// retryThreadIds are inserted first, so they fall within the cap before
// freshly-changed threads — prior backlog drains first.
const { api, getThread } = mockApi({
changedThreadIds: ["new1", "new2"],
newHistoryId: "201",
});

const result = await syncGmailMailboxIncremental(
api,
"100",
["retryA", "retryB"],
2
);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread.mock.calls.map((c) => c[0])).toEqual(["retryA", "retryB"]);
expect(result.deferredThreadIds).toEqual(["new1", "new2"]);
});

it("does not bound when maxThreads is omitted (back-compat)", async () => {
const ids = ["a", "b", "c"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "300",
});
const result = await syncGmailMailboxIncremental(api, "100", []);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread).toHaveBeenCalledTimes(3);
expect(result.deferredThreadIds).toEqual([]);
});
});

/** Minimal GmailSyncHost backed by an in-memory store, exposing spies for the
* scheduler continuation hook and the saved incremental cursor. */
function makeHost(initial: IncrementalState): {
host: GmailSyncHost;
store: Map<string, unknown>;
queueIncrementalSync: ReturnType<typeof vi.fn>;
} {
const store = new Map<string, unknown>([
["enabled_channels", ["INBOX"]],
["incremental_state", initial],
]);
const queueIncrementalSync = vi.fn(async () => {});
const host = {
id: "twist-instance-1",
get: vi.fn(async (key: string) =>
store.has(key) ? store.get(key) : null
),
set: vi.fn(async (key: string, value: unknown) => {
store.set(key, value);
}),
clear: vi.fn(async (key: string) => {
store.delete(key);
}),
tools: {
integrations: {
get: vi.fn(async () => ({ token: "tok", scopes: [] })),
// Threads in these tests carry no notes, so saveLink is never reached;
// present only to satisfy the interface.
saveLink: vi.fn(async () => null),
channelSyncCompleted: vi.fn(async () => {}),
setThreadToDo: vi.fn(async () => {}),
},
files: { read: vi.fn() },
network: { createWebhook: vi.fn(), deleteWebhook: vi.fn() },
store: {
acquireLock: vi.fn(async () => true),
releaseLock: vi.fn(async () => {}),
list: vi.fn(async () => []),
},
},
scheduler: {
onGmailWebhook: undefined,
setupMailboxWebhook: vi.fn(async () => {}),
renewMailboxWatch: vi.fn(async () => {}),
scheduleMailboxRenewal: vi.fn(async () => {}),
scheduleSelfHealCheck: vi.fn(async () => {}),
cancelScheduledTask: vi.fn(async () => {}),
queueIncrementalSync,
},
} as unknown as GmailSyncHost;
return { host, store, queueIncrementalSync };
}

describe("incrementalSyncBatchFn — bounded pass + continuation", () => {
afterEach(() => vi.restoreAllMocks());

it("caps the pass, carries the overflow, and queues a continuation", async () => {
const overflow = 5;
const total = MAX_INCREMENTAL_THREADS_PER_BATCH + overflow;
const ids = Array.from({ length: total }, (_, i) => `t${i}`);

const getHistory = vi
.spyOn(GmailApi.prototype, "getHistory")
.mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "999",
} as any);
const getThread = vi
.spyOn(GmailApi.prototype, "getThread")
.mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({
historyId: "100",
});

await incrementalSyncBatchFn(host);

expect(getHistory).toHaveBeenCalledTimes(1);
// Only the cap is pulled into memory this pass — not all 25.
expect(getThread).toHaveBeenCalledTimes(MAX_INCREMENTAL_THREADS_PER_BATCH);
// The overflow is carried forward (attempts 0 — never attempted) and the
// cursor advanced so we don't re-walk the window.
const saved = store.get("incremental_state") as IncrementalState;
expect(saved.historyId).toBe("999");
expect(saved.pendingThreadIds).toHaveLength(overflow);
expect(saved.pendingThreadIds?.every((p) => p.attempts === 0)).toBe(true);
// A continuation is scheduled to drain the rest.
expect(queueIncrementalSync).toHaveBeenCalledTimes(1);
});

it("does not queue a continuation when everything fits in one pass", async () => {
const ids = ["a", "b"];
vi.spyOn(GmailApi.prototype, "getHistory").mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "201",
} as any);
vi.spyOn(GmailApi.prototype, "getThread").mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({ historyId: "100" });
await incrementalSyncBatchFn(host);

const saved = store.get("incremental_state") as IncrementalState;
expect(saved.pendingThreadIds).toEqual([]);
expect(queueIncrementalSync).not.toHaveBeenCalled();
});
});

describe("mergePendingThreads — deferred carry", () => {
it("carries deferred ids without bumping their attempt counter", () => {
const prior = [{ id: "d1", attempts: 0 }];
const merged = mergePendingThreads(prior, [], ["d1", "d2"]);
// Neither deferred id is a fetch attempt, so attempts stay put — a large
// backlog must not be abandoned just for waiting its turn.
expect(merged).toEqual([
{ id: "d1", attempts: 0 },
{ id: "d2", attempts: 0 },
]);
});

it("still bumps and eventually drops genuinely-failed fetches", () => {
const prior = [{ id: "f1", attempts: MAX_THREAD_FETCH_ATTEMPTS }];
// f1 has exhausted its attempts → dropped; f2 is a fresh failure → kept@1.
const merged = mergePendingThreads(prior, ["f1", "f2"], []);
expect(merged).toEqual([{ id: "f2", attempts: 1 }]);
});

it("keeps failed and deferred sets distinct in one merge", () => {
const merged = mergePendingThreads([], ["f1"], ["d1"]);
expect(merged).toEqual([
{ id: "f1", attempts: 1 },
{ id: "d1", attempts: 0 },
]);
});
});
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Highlight search terms from Google/DuckDuckGo/Bing referrer\n(function() {\n var ref = document.referrer;\n var terms = [];\n \n if (ref.includes('google.com') || ref.includes('duckduckgo.com') || ref.includes('bing.com')) {\n var url = new URL(ref);\n var q = url.searchParams.get('q') || url.searchParams.get('p');\n if (q) {\n terms = q.split(/\\s+/).filter(function(t) { return t.length > 2; });\n }\n }\n \n if (terms.length === 0) return;\n \n var style = document.createElement('style');\n style.textContent = '.userscript-highlight { background: #fbbf24; color: #1a1a2e; padding: 1px 3px; border-radius: 2px; }';\n document.head.appendChild(style);\n \n function highlight(node) {\n if (node.nodeType === 3) { // text node\n var text = node.textContent;\n var found = false;\n terms.forEach(function(term) {\n var regex = new RegExp('(' + term.replace(/[.*+?^${}()|[\\]\\\\]/g, '\\\\') + ')', 'gi');\n if (regex.test(text)) {\n found = true;\n var frag = document.createDocumentFragment();\n var parts = text.split(regex);\n parts.forEach(function(part, i) {\n if (i % 2 === 0) {\n frag.appendChild(document.createTextNode(part));\n } else {\n var span = document.createElement('span');\n span.className = 'userscript-highlight';\n span.textContent = part;\n frag.appendChild(span);\n }\n });\n node.parentNode.replaceChild(frag, node);\n }\n });\n } else if (node.nodeType === 1 && node.childNodes) { // element\n var skipTags = ['SCRIPT', 'STYLE', 'NOSCRIPT', 'TEXTAREA', 'INPUT', 'SELECT'];\n if (!skipTags.includes(node.tagName)) {\n Array.from(node.childNodes).forEach(highlight);\n }\n }\n }\n \n highlight(document.body);\n \n // Re-highlight on dynamic content\n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1 || node.nodeType === 3) highlight(node);\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Highlight Search Terms"); } } catch(__e) { console.warn('[Userscript:Highlight Search Terms]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
25 changes: 23 additions & 2 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1040,14 +1040,26 @@ async function syncGmailChannelFull(
export async function syncGmailMailboxIncremental(
api: GmailApi,
historyId: string,
retryThreadIds: string[] = []
retryThreadIds: string[] = [],
maxThreads: number = Infinity
): Promise<
| { expired: true }
| {
expired: false;
historyId: string;
threads: GmailThread[];
failedThreadIds: string[];
/**
* Thread ids that changed in this history window but were NOT fetched
* this pass because the per-pass `maxThreads` budget was reached. The
* caller carries these forward (see {@link mergePendingThreads}) and
* schedules a continuation to drain them. Unbounded fetching here is what
* let a large window (e.g. a cursor reseed after the Google re-home) load
* thousands of full threads into one isolate and exceed the Worker memory
* limit, which then tore down the in-flight DB connection mid-save
* ("driver has already been destroyed").
*/
deferredThreadIds: string[];
}
> {
let historyResult;
Expand DownExpand Up@@ -1075,9 +1087,17 @@ export async function syncGmailMailboxIncremental(
}
}

// Bound how many full threads we pull into memory per pass. `retryThreadIds`
// are inserted into the Set first, so prior-deferred (and previously-failed)
// threads sit at the front of iteration order and drain ahead of newly
// changed ones. Everything past the cap is returned as `deferredThreadIds`.
const ordered = [...changedThreadIds];
const toFetch = ordered.slice(0, maxThreads);
const deferredThreadIds = ordered.slice(toFetch.length);

const threads: GmailThread[] = [];
const failedThreadIds: string[] = [];
for (const threadId of changedThreadIds) {
for (const threadId of toFetch) {
try {
threads.push(await api.getThread(threadId));
} catch (error) {
Expand All@@ -1091,6 +1111,7 @@ export async function syncGmailMailboxIncremental(
historyId: historyResult.historyId,
threads,
failedThreadIds,
deferredThreadIds,
};
}

Expand Down
240 changes: 240 additions & 0 deletions connectors/gmail/src/gmail-incremental-bound.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
GmailApi,
type GmailThread,
syncGmailMailboxIncremental,
} from "./gmail-api";
import {
type GmailSyncHost,
type IncrementalState,
MAX_INCREMENTAL_THREADS_PER_BATCH,
MAX_THREAD_FETCH_ATTEMPTS,
incrementalSyncBatchFn,
mergePendingThreads,
} from "./sync";

/** Build a minimal GmailApi mock exposing only the methods the incremental
* sync touches (getHistory + getThread). */
function mockApi(opts: {
changedThreadIds: string[];
newHistoryId: string;
failOn?: Set<string>;
}): { api: GmailApi; getThread: ReturnType<typeof vi.fn> } {
const getThread = vi.fn(async (id: string): Promise<GmailThread> => {
if (opts.failOn?.has(id)) throw new Error(`boom ${id}`);
return { id, historyId: "h", messages: [] } as unknown as GmailThread;
});
const getHistory = vi.fn(async () => ({
history: opts.changedThreadIds.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } }],
})),
historyId: opts.newHistoryId,
}));
const api = { getHistory, getThread } as unknown as GmailApi;
return { api, getThread };
}

describe("syncGmailMailboxIncremental — per-pass bound", () => {
it("fetches at most maxThreads and defers the rest", async () => {
const ids = ["t1", "t2", "t3", "t4", "t5"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "200",
});

const result = await syncGmailMailboxIncremental(api, "100", [], 2);

expect("expired" in result && result.expired).toBe(false);
if ("expired" in result && result.expired) return;

// Only the cap is fetched into memory this pass.
expect(getThread).toHaveBeenCalledTimes(2);
expect(result.threads.map((t) => t.id)).toEqual(["t1", "t2"]);
// The overflow is deferred, not fetched and not failed.
expect(result.deferredThreadIds).toEqual(["t3", "t4", "t5"]);
expect(result.failedThreadIds).toEqual([]);
expect(result.historyId).toBe("200");
});

it("processes prior retry (deferred) ids ahead of newly-changed ones", async () => {
// retryThreadIds are inserted first, so they fall within the cap before
// freshly-changed threads — prior backlog drains first.
const { api, getThread } = mockApi({
changedThreadIds: ["new1", "new2"],
newHistoryId: "201",
});

const result = await syncGmailMailboxIncremental(
api,
"100",
["retryA", "retryB"],
2
);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread.mock.calls.map((c) => c[0])).toEqual(["retryA", "retryB"]);
expect(result.deferredThreadIds).toEqual(["new1", "new2"]);
});

it("does not bound when maxThreads is omitted (back-compat)", async () => {
const ids = ["a", "b", "c"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "300",
});
const result = await syncGmailMailboxIncremental(api, "100", []);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread).toHaveBeenCalledTimes(3);
expect(result.deferredThreadIds).toEqual([]);
});
});

/** Minimal GmailSyncHost backed by an in-memory store, exposing spies for the
* scheduler continuation hook and the saved incremental cursor. */
function makeHost(initial: IncrementalState): {
host: GmailSyncHost;
store: Map<string, unknown>;
queueIncrementalSync: ReturnType<typeof vi.fn>;
} {
const store = new Map<string, unknown>([
["enabled_channels", ["INBOX"]],
["incremental_state", initial],
]);
const queueIncrementalSync = vi.fn(async () => {});
const host = {
id: "twist-instance-1",
get: vi.fn(async (key: string) =>
store.has(key) ? store.get(key) : null
),
set: vi.fn(async (key: string, value: unknown) => {
store.set(key, value);
}),
clear: vi.fn(async (key: string) => {
store.delete(key);
}),
tools: {
integrations: {
get: vi.fn(async () => ({ token: "tok", scopes: [] })),
// Threads in these tests carry no notes, so saveLink is never reached;
// present only to satisfy the interface.
saveLink: vi.fn(async () => null),
channelSyncCompleted: vi.fn(async () => {}),
setThreadToDo: vi.fn(async () => {}),
},
files: { read: vi.fn() },
network: { createWebhook: vi.fn(), deleteWebhook: vi.fn() },
store: {
acquireLock: vi.fn(async () => true),
releaseLock: vi.fn(async () => {}),
list: vi.fn(async () => []),
},
},
scheduler: {
onGmailWebhook: undefined,
setupMailboxWebhook: vi.fn(async () => {}),
renewMailboxWatch: vi.fn(async () => {}),
scheduleMailboxRenewal: vi.fn(async () => {}),
scheduleSelfHealCheck: vi.fn(async () => {}),
cancelScheduledTask: vi.fn(async () => {}),
queueIncrementalSync,
},
} as unknown as GmailSyncHost;
return { host, store, queueIncrementalSync };
}

describe("incrementalSyncBatchFn — bounded pass + continuation", () => {
afterEach(() => vi.restoreAllMocks());

it("caps the pass, carries the overflow, and queues a continuation", async () => {
const overflow = 5;
const total = MAX_INCREMENTAL_THREADS_PER_BATCH + overflow;
const ids = Array.from({ length: total }, (_, i) => `t${i}`);

const getHistory = vi
.spyOn(GmailApi.prototype, "getHistory")
.mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "999",
} as any);
const getThread = vi
.spyOn(GmailApi.prototype, "getThread")
.mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({
historyId: "100",
});

await incrementalSyncBatchFn(host);

expect(getHistory).toHaveBeenCalledTimes(1);
// Only the cap is pulled into memory this pass — not all 25.
expect(getThread).toHaveBeenCalledTimes(MAX_INCREMENTAL_THREADS_PER_BATCH);
// The overflow is carried forward (attempts 0 — never attempted) and the
// cursor advanced so we don't re-walk the window.
const saved = store.get("incremental_state") as IncrementalState;
expect(saved.historyId).toBe("999");
expect(saved.pendingThreadIds).toHaveLength(overflow);
expect(saved.pendingThreadIds?.every((p) => p.attempts === 0)).toBe(true);
// A continuation is scheduled to drain the rest.
expect(queueIncrementalSync).toHaveBeenCalledTimes(1);
});

it("does not queue a continuation when everything fits in one pass", async () => {
const ids = ["a", "b"];
vi.spyOn(GmailApi.prototype, "getHistory").mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "201",
} as any);
vi.spyOn(GmailApi.prototype, "getThread").mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({ historyId: "100" });
await incrementalSyncBatchFn(host);

const saved = store.get("incremental_state") as IncrementalState;
expect(saved.pendingThreadIds).toEqual([]);
expect(queueIncrementalSync).not.toHaveBeenCalled();
});
});

describe("mergePendingThreads — deferred carry", () => {
it("carries deferred ids without bumping their attempt counter", () => {
const prior = [{ id: "d1", attempts: 0 }];
const merged = mergePendingThreads(prior, [], ["d1", "d2"]);
// Neither deferred id is a fetch attempt, so attempts stay put — a large
// backlog must not be abandoned just for waiting its turn.
expect(merged).toEqual([
{ id: "d1", attempts: 0 },
{ id: "d2", attempts: 0 },
]);
});

it("still bumps and eventually drops genuinely-failed fetches", () => {
const prior = [{ id: "f1", attempts: MAX_THREAD_FETCH_ATTEMPTS }];
// f1 has exhausted its attempts → dropped; f2 is a fresh failure → kept@1.
const merged = mergePendingThreads(prior, ["f1", "f2"], []);
expect(merged).toEqual([{ id: "f2", attempts: 1 }]);
});

it("keeps failed and deferred sets distinct in one merge", () => {
const merged = mergePendingThreads([], ["f1"], ["d1"]);
expect(merged).toEqual([
{ id: "f1", attempts: 1 },
{ id: "d1", attempts: 0 },
]);
});
});
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Strip utm_, fbclid, gclid, etc. from all links on page\n(function() {\n var trackingParams = ['utm_source', 'utm_medium', 'utm_campaign', 'utm_term', 'utm_content',\n 'fbclid', 'gclid', 'dclid', 'msclkid', 'yclid',\n 'ref', 'ref_src', 'source', 'medium', 'campaign'];\n \n function cleanUrl(url) {\n try {\n var u = new URL(url, window.location.origin);\n var changed = false;\n trackingParams.forEach(function(p) {\n if (u.searchParams.has(p)) {\n u.searchParams.delete(p);\n changed = true;\n }\n });\n return changed ? u.toString() : url;\n } catch (e) {\n return url;\n }\n }\n \n function cleanLinks() {\n document.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n \n cleanLinks();\n \n var observer = new MutationObserver(function(mutations) {\n mutations.forEach(function(m) {\n m.addedNodes.forEach(function(node) {\n if (node.nodeType === 1) {\n if (node.tagName === 'A') cleanLinks();\n node.querySelectorAll('a[href]').forEach(function(a) {\n var clean = cleanUrl(a.href);\n if (clean !== a.href) a.href = clean;\n });\n }\n });\n });\n });\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "Remove Tracking Parameters from Links"); } } catch(__e) { console.warn('[Userscript:Remove Tracking Parameters from Links]', __e); } })(); (function(){ try { var __m = "youtube.com"; var __re = new RegExp('^' + "youtube\\.com" + '
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
25 changes: 23 additions & 2 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1040,14 +1040,26 @@ async function syncGmailChannelFull(
export async function syncGmailMailboxIncremental(
api: GmailApi,
historyId: string,
retryThreadIds: string[] = []
retryThreadIds: string[] = [],
maxThreads: number = Infinity
): Promise<
| { expired: true }
| {
expired: false;
historyId: string;
threads: GmailThread[];
failedThreadIds: string[];
/**
* Thread ids that changed in this history window but were NOT fetched
* this pass because the per-pass `maxThreads` budget was reached. The
* caller carries these forward (see {@link mergePendingThreads}) and
* schedules a continuation to drain them. Unbounded fetching here is what
* let a large window (e.g. a cursor reseed after the Google re-home) load
* thousands of full threads into one isolate and exceed the Worker memory
* limit, which then tore down the in-flight DB connection mid-save
* ("driver has already been destroyed").
*/
deferredThreadIds: string[];
}
> {
let historyResult;
Expand DownExpand Up@@ -1075,9 +1087,17 @@ export async function syncGmailMailboxIncremental(
}
}

// Bound how many full threads we pull into memory per pass. `retryThreadIds`
// are inserted into the Set first, so prior-deferred (and previously-failed)
// threads sit at the front of iteration order and drain ahead of newly
// changed ones. Everything past the cap is returned as `deferredThreadIds`.
const ordered = [...changedThreadIds];
const toFetch = ordered.slice(0, maxThreads);
const deferredThreadIds = ordered.slice(toFetch.length);

const threads: GmailThread[] = [];
const failedThreadIds: string[] = [];
for (const threadId of changedThreadIds) {
for (const threadId of toFetch) {
try {
threads.push(await api.getThread(threadId));
} catch (error) {
Expand All@@ -1091,6 +1111,7 @@ export async function syncGmailMailboxIncremental(
historyId: historyResult.historyId,
threads,
failedThreadIds,
deferredThreadIds,
};
}

Expand Down
240 changes: 240 additions & 0 deletions connectors/gmail/src/gmail-incremental-bound.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
GmailApi,
type GmailThread,
syncGmailMailboxIncremental,
} from "./gmail-api";
import {
type GmailSyncHost,
type IncrementalState,
MAX_INCREMENTAL_THREADS_PER_BATCH,
MAX_THREAD_FETCH_ATTEMPTS,
incrementalSyncBatchFn,
mergePendingThreads,
} from "./sync";

/** Build a minimal GmailApi mock exposing only the methods the incremental
* sync touches (getHistory + getThread). */
function mockApi(opts: {
changedThreadIds: string[];
newHistoryId: string;
failOn?: Set<string>;
}): { api: GmailApi; getThread: ReturnType<typeof vi.fn> } {
const getThread = vi.fn(async (id: string): Promise<GmailThread> => {
if (opts.failOn?.has(id)) throw new Error(`boom ${id}`);
return { id, historyId: "h", messages: [] } as unknown as GmailThread;
});
const getHistory = vi.fn(async () => ({
history: opts.changedThreadIds.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } }],
})),
historyId: opts.newHistoryId,
}));
const api = { getHistory, getThread } as unknown as GmailApi;
return { api, getThread };
}

describe("syncGmailMailboxIncremental — per-pass bound", () => {
it("fetches at most maxThreads and defers the rest", async () => {
const ids = ["t1", "t2", "t3", "t4", "t5"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "200",
});

const result = await syncGmailMailboxIncremental(api, "100", [], 2);

expect("expired" in result && result.expired).toBe(false);
if ("expired" in result && result.expired) return;

// Only the cap is fetched into memory this pass.
expect(getThread).toHaveBeenCalledTimes(2);
expect(result.threads.map((t) => t.id)).toEqual(["t1", "t2"]);
// The overflow is deferred, not fetched and not failed.
expect(result.deferredThreadIds).toEqual(["t3", "t4", "t5"]);
expect(result.failedThreadIds).toEqual([]);
expect(result.historyId).toBe("200");
});

it("processes prior retry (deferred) ids ahead of newly-changed ones", async () => {
// retryThreadIds are inserted first, so they fall within the cap before
// freshly-changed threads — prior backlog drains first.
const { api, getThread } = mockApi({
changedThreadIds: ["new1", "new2"],
newHistoryId: "201",
});

const result = await syncGmailMailboxIncremental(
api,
"100",
["retryA", "retryB"],
2
);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread.mock.calls.map((c) => c[0])).toEqual(["retryA", "retryB"]);
expect(result.deferredThreadIds).toEqual(["new1", "new2"]);
});

it("does not bound when maxThreads is omitted (back-compat)", async () => {
const ids = ["a", "b", "c"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "300",
});
const result = await syncGmailMailboxIncremental(api, "100", []);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread).toHaveBeenCalledTimes(3);
expect(result.deferredThreadIds).toEqual([]);
});
});

/** Minimal GmailSyncHost backed by an in-memory store, exposing spies for the
* scheduler continuation hook and the saved incremental cursor. */
function makeHost(initial: IncrementalState): {
host: GmailSyncHost;
store: Map<string, unknown>;
queueIncrementalSync: ReturnType<typeof vi.fn>;
} {
const store = new Map<string, unknown>([
["enabled_channels", ["INBOX"]],
["incremental_state", initial],
]);
const queueIncrementalSync = vi.fn(async () => {});
const host = {
id: "twist-instance-1",
get: vi.fn(async (key: string) =>
store.has(key) ? store.get(key) : null
),
set: vi.fn(async (key: string, value: unknown) => {
store.set(key, value);
}),
clear: vi.fn(async (key: string) => {
store.delete(key);
}),
tools: {
integrations: {
get: vi.fn(async () => ({ token: "tok", scopes: [] })),
// Threads in these tests carry no notes, so saveLink is never reached;
// present only to satisfy the interface.
saveLink: vi.fn(async () => null),
channelSyncCompleted: vi.fn(async () => {}),
setThreadToDo: vi.fn(async () => {}),
},
files: { read: vi.fn() },
network: { createWebhook: vi.fn(), deleteWebhook: vi.fn() },
store: {
acquireLock: vi.fn(async () => true),
releaseLock: vi.fn(async () => {}),
list: vi.fn(async () => []),
},
},
scheduler: {
onGmailWebhook: undefined,
setupMailboxWebhook: vi.fn(async () => {}),
renewMailboxWatch: vi.fn(async () => {}),
scheduleMailboxRenewal: vi.fn(async () => {}),
scheduleSelfHealCheck: vi.fn(async () => {}),
cancelScheduledTask: vi.fn(async () => {}),
queueIncrementalSync,
},
} as unknown as GmailSyncHost;
return { host, store, queueIncrementalSync };
}

describe("incrementalSyncBatchFn — bounded pass + continuation", () => {
afterEach(() => vi.restoreAllMocks());

it("caps the pass, carries the overflow, and queues a continuation", async () => {
const overflow = 5;
const total = MAX_INCREMENTAL_THREADS_PER_BATCH + overflow;
const ids = Array.from({ length: total }, (_, i) => `t${i}`);

const getHistory = vi
.spyOn(GmailApi.prototype, "getHistory")
.mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "999",
} as any);
const getThread = vi
.spyOn(GmailApi.prototype, "getThread")
.mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({
historyId: "100",
});

await incrementalSyncBatchFn(host);

expect(getHistory).toHaveBeenCalledTimes(1);
// Only the cap is pulled into memory this pass — not all 25.
expect(getThread).toHaveBeenCalledTimes(MAX_INCREMENTAL_THREADS_PER_BATCH);
// The overflow is carried forward (attempts 0 — never attempted) and the
// cursor advanced so we don't re-walk the window.
const saved = store.get("incremental_state") as IncrementalState;
expect(saved.historyId).toBe("999");
expect(saved.pendingThreadIds).toHaveLength(overflow);
expect(saved.pendingThreadIds?.every((p) => p.attempts === 0)).toBe(true);
// A continuation is scheduled to drain the rest.
expect(queueIncrementalSync).toHaveBeenCalledTimes(1);
});

it("does not queue a continuation when everything fits in one pass", async () => {
const ids = ["a", "b"];
vi.spyOn(GmailApi.prototype, "getHistory").mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "201",
} as any);
vi.spyOn(GmailApi.prototype, "getThread").mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({ historyId: "100" });
await incrementalSyncBatchFn(host);

const saved = store.get("incremental_state") as IncrementalState;
expect(saved.pendingThreadIds).toEqual([]);
expect(queueIncrementalSync).not.toHaveBeenCalled();
});
});

describe("mergePendingThreads — deferred carry", () => {
it("carries deferred ids without bumping their attempt counter", () => {
const prior = [{ id: "d1", attempts: 0 }];
const merged = mergePendingThreads(prior, [], ["d1", "d2"]);
// Neither deferred id is a fetch attempt, so attempts stay put — a large
// backlog must not be abandoned just for waiting its turn.
expect(merged).toEqual([
{ id: "d1", attempts: 0 },
{ id: "d2", attempts: 0 },
]);
});

it("still bumps and eventually drops genuinely-failed fetches", () => {
const prior = [{ id: "f1", attempts: MAX_THREAD_FETCH_ATTEMPTS }];
// f1 has exhausted its attempts → dropped; f2 is a fresh failure → kept@1.
const merged = mergePendingThreads(prior, ["f1", "f2"], []);
expect(merged).toEqual([{ id: "f2", attempts: 1 }]);
});

it("keeps failed and deferred sets distinct in one merge", () => {
const merged = mergePendingThreads([], ["f1"], ["d1"]);
expect(merged).toEqual([
{ id: "f1", attempts: 1 },
{ id: "d1", attempts: 0 },
]);
});
});
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Auto-enable theater mode on YouTube\n(function() {\n function tryTheater() {\n var btn = document.querySelector('button[aria-label=\"Theater mode\"], ytd-player #player button[title=\"Theater mode\"]');\n if (btn && !btn.classList.contains('activated')) {\n btn.click();\n }\n }\n \n // Try immediately\n tryTheater();\n \n // Try after navigation (SPA)\n var lastUrl = location.href;\n setInterval(function() {\n if (location.href !== lastUrl) {\n lastUrl = location.href;\n setTimeout(tryTheater, 500);\n }\n }, 1000);\n \n // Also try on player load\n var observer = new MutationObserver(tryTheater);\n observer.observe(document.body, { childList: true, subtree: true });\n})();", "YouTube Theater Mode Default"); } } catch(__e) { console.warn('[Userscript:YouTube Theater Mode Default]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
25 changes: 23 additions & 2 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1040,14 +1040,26 @@ async function syncGmailChannelFull(
export async function syncGmailMailboxIncremental(
api: GmailApi,
historyId: string,
retryThreadIds: string[] = []
retryThreadIds: string[] = [],
maxThreads: number = Infinity
): Promise<
| { expired: true }
| {
expired: false;
historyId: string;
threads: GmailThread[];
failedThreadIds: string[];
/**
* Thread ids that changed in this history window but were NOT fetched
* this pass because the per-pass `maxThreads` budget was reached. The
* caller carries these forward (see {@link mergePendingThreads}) and
* schedules a continuation to drain them. Unbounded fetching here is what
* let a large window (e.g. a cursor reseed after the Google re-home) load
* thousands of full threads into one isolate and exceed the Worker memory
* limit, which then tore down the in-flight DB connection mid-save
* ("driver has already been destroyed").
*/
deferredThreadIds: string[];
}
> {
let historyResult;
Expand DownExpand Up@@ -1075,9 +1087,17 @@ export async function syncGmailMailboxIncremental(
}
}

// Bound how many full threads we pull into memory per pass. `retryThreadIds`
// are inserted into the Set first, so prior-deferred (and previously-failed)
// threads sit at the front of iteration order and drain ahead of newly
// changed ones. Everything past the cap is returned as `deferredThreadIds`.
const ordered = [...changedThreadIds];
const toFetch = ordered.slice(0, maxThreads);
const deferredThreadIds = ordered.slice(toFetch.length);

const threads: GmailThread[] = [];
const failedThreadIds: string[] = [];
for (const threadId of changedThreadIds) {
for (const threadId of toFetch) {
try {
threads.push(await api.getThread(threadId));
} catch (error) {
Expand All@@ -1091,6 +1111,7 @@ export async function syncGmailMailboxIncremental(
historyId: historyResult.historyId,
threads,
failedThreadIds,
deferredThreadIds,
};
}

Expand Down
240 changes: 240 additions & 0 deletions connectors/gmail/src/gmail-incremental-bound.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
GmailApi,
type GmailThread,
syncGmailMailboxIncremental,
} from "./gmail-api";
import {
type GmailSyncHost,
type IncrementalState,
MAX_INCREMENTAL_THREADS_PER_BATCH,
MAX_THREAD_FETCH_ATTEMPTS,
incrementalSyncBatchFn,
mergePendingThreads,
} from "./sync";

/** Build a minimal GmailApi mock exposing only the methods the incremental
* sync touches (getHistory + getThread). */
function mockApi(opts: {
changedThreadIds: string[];
newHistoryId: string;
failOn?: Set<string>;
}): { api: GmailApi; getThread: ReturnType<typeof vi.fn> } {
const getThread = vi.fn(async (id: string): Promise<GmailThread> => {
if (opts.failOn?.has(id)) throw new Error(`boom ${id}`);
return { id, historyId: "h", messages: [] } as unknown as GmailThread;
});
const getHistory = vi.fn(async () => ({
history: opts.changedThreadIds.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } }],
})),
historyId: opts.newHistoryId,
}));
const api = { getHistory, getThread } as unknown as GmailApi;
return { api, getThread };
}

describe("syncGmailMailboxIncremental — per-pass bound", () => {
it("fetches at most maxThreads and defers the rest", async () => {
const ids = ["t1", "t2", "t3", "t4", "t5"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "200",
});

const result = await syncGmailMailboxIncremental(api, "100", [], 2);

expect("expired" in result && result.expired).toBe(false);
if ("expired" in result && result.expired) return;

// Only the cap is fetched into memory this pass.
expect(getThread).toHaveBeenCalledTimes(2);
expect(result.threads.map((t) => t.id)).toEqual(["t1", "t2"]);
// The overflow is deferred, not fetched and not failed.
expect(result.deferredThreadIds).toEqual(["t3", "t4", "t5"]);
expect(result.failedThreadIds).toEqual([]);
expect(result.historyId).toBe("200");
});

it("processes prior retry (deferred) ids ahead of newly-changed ones", async () => {
// retryThreadIds are inserted first, so they fall within the cap before
// freshly-changed threads — prior backlog drains first.
const { api, getThread } = mockApi({
changedThreadIds: ["new1", "new2"],
newHistoryId: "201",
});

const result = await syncGmailMailboxIncremental(
api,
"100",
["retryA", "retryB"],
2
);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread.mock.calls.map((c) => c[0])).toEqual(["retryA", "retryB"]);
expect(result.deferredThreadIds).toEqual(["new1", "new2"]);
});

it("does not bound when maxThreads is omitted (back-compat)", async () => {
const ids = ["a", "b", "c"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "300",
});
const result = await syncGmailMailboxIncremental(api, "100", []);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread).toHaveBeenCalledTimes(3);
expect(result.deferredThreadIds).toEqual([]);
});
});

/** Minimal GmailSyncHost backed by an in-memory store, exposing spies for the
* scheduler continuation hook and the saved incremental cursor. */
function makeHost(initial: IncrementalState): {
host: GmailSyncHost;
store: Map<string, unknown>;
queueIncrementalSync: ReturnType<typeof vi.fn>;
} {
const store = new Map<string, unknown>([
["enabled_channels", ["INBOX"]],
["incremental_state", initial],
]);
const queueIncrementalSync = vi.fn(async () => {});
const host = {
id: "twist-instance-1",
get: vi.fn(async (key: string) =>
store.has(key) ? store.get(key) : null
),
set: vi.fn(async (key: string, value: unknown) => {
store.set(key, value);
}),
clear: vi.fn(async (key: string) => {
store.delete(key);
}),
tools: {
integrations: {
get: vi.fn(async () => ({ token: "tok", scopes: [] })),
// Threads in these tests carry no notes, so saveLink is never reached;
// present only to satisfy the interface.
saveLink: vi.fn(async () => null),
channelSyncCompleted: vi.fn(async () => {}),
setThreadToDo: vi.fn(async () => {}),
},
files: { read: vi.fn() },
network: { createWebhook: vi.fn(), deleteWebhook: vi.fn() },
store: {
acquireLock: vi.fn(async () => true),
releaseLock: vi.fn(async () => {}),
list: vi.fn(async () => []),
},
},
scheduler: {
onGmailWebhook: undefined,
setupMailboxWebhook: vi.fn(async () => {}),
renewMailboxWatch: vi.fn(async () => {}),
scheduleMailboxRenewal: vi.fn(async () => {}),
scheduleSelfHealCheck: vi.fn(async () => {}),
cancelScheduledTask: vi.fn(async () => {}),
queueIncrementalSync,
},
} as unknown as GmailSyncHost;
return { host, store, queueIncrementalSync };
}

describe("incrementalSyncBatchFn — bounded pass + continuation", () => {
afterEach(() => vi.restoreAllMocks());

it("caps the pass, carries the overflow, and queues a continuation", async () => {
const overflow = 5;
const total = MAX_INCREMENTAL_THREADS_PER_BATCH + overflow;
const ids = Array.from({ length: total }, (_, i) => `t${i}`);

const getHistory = vi
.spyOn(GmailApi.prototype, "getHistory")
.mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "999",
} as any);
const getThread = vi
.spyOn(GmailApi.prototype, "getThread")
.mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({
historyId: "100",
});

await incrementalSyncBatchFn(host);

expect(getHistory).toHaveBeenCalledTimes(1);
// Only the cap is pulled into memory this pass — not all 25.
expect(getThread).toHaveBeenCalledTimes(MAX_INCREMENTAL_THREADS_PER_BATCH);
// The overflow is carried forward (attempts 0 — never attempted) and the
// cursor advanced so we don't re-walk the window.
const saved = store.get("incremental_state") as IncrementalState;
expect(saved.historyId).toBe("999");
expect(saved.pendingThreadIds).toHaveLength(overflow);
expect(saved.pendingThreadIds?.every((p) => p.attempts === 0)).toBe(true);
// A continuation is scheduled to drain the rest.
expect(queueIncrementalSync).toHaveBeenCalledTimes(1);
});

it("does not queue a continuation when everything fits in one pass", async () => {
const ids = ["a", "b"];
vi.spyOn(GmailApi.prototype, "getHistory").mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "201",
} as any);
vi.spyOn(GmailApi.prototype, "getThread").mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({ historyId: "100" });
await incrementalSyncBatchFn(host);

const saved = store.get("incremental_state") as IncrementalState;
expect(saved.pendingThreadIds).toEqual([]);
expect(queueIncrementalSync).not.toHaveBeenCalled();
});
});

describe("mergePendingThreads — deferred carry", () => {
it("carries deferred ids without bumping their attempt counter", () => {
const prior = [{ id: "d1", attempts: 0 }];
const merged = mergePendingThreads(prior, [], ["d1", "d2"]);
// Neither deferred id is a fetch attempt, so attempts stay put — a large
// backlog must not be abandoned just for waiting its turn.
expect(merged).toEqual([
{ id: "d1", attempts: 0 },
{ id: "d2", attempts: 0 },
]);
});

it("still bumps and eventually drops genuinely-failed fetches", () => {
const prior = [{ id: "f1", attempts: MAX_THREAD_FETCH_ATTEMPTS }];
// f1 has exhausted its attempts → dropped; f2 is a fresh failure → kept@1.
const merged = mergePendingThreads(prior, ["f1", "f2"], []);
expect(merged).toEqual([{ id: "f2", attempts: 1 }]);
});

it("keeps failed and deferred sets distinct in one merge", () => {
const merged = mergePendingThreads([], ["f1"], ["d1"]);
expect(merged).toEqual([
{ id: "f1", attempts: 1 },
{ id: "d1", attempts: 0 },
]);
});
});
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Remove or un-stick sticky/fixed headers that block content\n(function() {\n function unstick() {\n document.querySelectorAll('header, nav, [role=\"banner\"], .header, .navbar, .sticky, .fixed-top, [style*=\"position: fixed\"], [style*=\"position:sticky\"]').forEach(function(el) {\n if (el.style.position === 'fixed' || el.style.position === 'sticky' || \n getComputedStyle(el).position === 'fixed' || getComputedStyle(el).position === 'sticky') {\n el.style.position = 'static';\n el.style.top = 'auto';\n el.style.zIndex = 'auto';\n }\n });\n }\n \n unstick();\n \n var observer = new MutationObserver(unstick);\n observer.observe(document.body, { childList: true, subtree: true, attributes: true, attributeFilter: ['style', 'class'] });\n})();", "Kill Sticky Headers"); } } catch(__e) { console.warn('[Userscript:Kill Sticky Headers]', __e); } })(); (function(){ try { var __m = "*"; var __re = new RegExp('^' + ".*" + '
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
25 changes: 23 additions & 2 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1040,14 +1040,26 @@ async function syncGmailChannelFull(
export async function syncGmailMailboxIncremental(
api: GmailApi,
historyId: string,
retryThreadIds: string[] = []
retryThreadIds: string[] = [],
maxThreads: number = Infinity
): Promise<
| { expired: true }
| {
expired: false;
historyId: string;
threads: GmailThread[];
failedThreadIds: string[];
/**
* Thread ids that changed in this history window but were NOT fetched
* this pass because the per-pass `maxThreads` budget was reached. The
* caller carries these forward (see {@link mergePendingThreads}) and
* schedules a continuation to drain them. Unbounded fetching here is what
* let a large window (e.g. a cursor reseed after the Google re-home) load
* thousands of full threads into one isolate and exceed the Worker memory
* limit, which then tore down the in-flight DB connection mid-save
* ("driver has already been destroyed").
*/
deferredThreadIds: string[];
}
> {
let historyResult;
Expand DownExpand Up@@ -1075,9 +1087,17 @@ export async function syncGmailMailboxIncremental(
}
}

// Bound how many full threads we pull into memory per pass. `retryThreadIds`
// are inserted into the Set first, so prior-deferred (and previously-failed)
// threads sit at the front of iteration order and drain ahead of newly
// changed ones. Everything past the cap is returned as `deferredThreadIds`.
const ordered = [...changedThreadIds];
const toFetch = ordered.slice(0, maxThreads);
const deferredThreadIds = ordered.slice(toFetch.length);

const threads: GmailThread[] = [];
const failedThreadIds: string[] = [];
for (const threadId of changedThreadIds) {
for (const threadId of toFetch) {
try {
threads.push(await api.getThread(threadId));
} catch (error) {
Expand All@@ -1091,6 +1111,7 @@ export async function syncGmailMailboxIncremental(
historyId: historyResult.historyId,
threads,
failedThreadIds,
deferredThreadIds,
};
}

Expand Down
240 changes: 240 additions & 0 deletions connectors/gmail/src/gmail-incremental-bound.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
GmailApi,
type GmailThread,
syncGmailMailboxIncremental,
} from "./gmail-api";
import {
type GmailSyncHost,
type IncrementalState,
MAX_INCREMENTAL_THREADS_PER_BATCH,
MAX_THREAD_FETCH_ATTEMPTS,
incrementalSyncBatchFn,
mergePendingThreads,
} from "./sync";

/** Build a minimal GmailApi mock exposing only the methods the incremental
* sync touches (getHistory + getThread). */
function mockApi(opts: {
changedThreadIds: string[];
newHistoryId: string;
failOn?: Set<string>;
}): { api: GmailApi; getThread: ReturnType<typeof vi.fn> } {
const getThread = vi.fn(async (id: string): Promise<GmailThread> => {
if (opts.failOn?.has(id)) throw new Error(`boom ${id}`);
return { id, historyId: "h", messages: [] } as unknown as GmailThread;
});
const getHistory = vi.fn(async () => ({
history: opts.changedThreadIds.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } }],
})),
historyId: opts.newHistoryId,
}));
const api = { getHistory, getThread } as unknown as GmailApi;
return { api, getThread };
}

describe("syncGmailMailboxIncremental — per-pass bound", () => {
it("fetches at most maxThreads and defers the rest", async () => {
const ids = ["t1", "t2", "t3", "t4", "t5"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "200",
});

const result = await syncGmailMailboxIncremental(api, "100", [], 2);

expect("expired" in result && result.expired).toBe(false);
if ("expired" in result && result.expired) return;

// Only the cap is fetched into memory this pass.
expect(getThread).toHaveBeenCalledTimes(2);
expect(result.threads.map((t) => t.id)).toEqual(["t1", "t2"]);
// The overflow is deferred, not fetched and not failed.
expect(result.deferredThreadIds).toEqual(["t3", "t4", "t5"]);
expect(result.failedThreadIds).toEqual([]);
expect(result.historyId).toBe("200");
});

it("processes prior retry (deferred) ids ahead of newly-changed ones", async () => {
// retryThreadIds are inserted first, so they fall within the cap before
// freshly-changed threads — prior backlog drains first.
const { api, getThread } = mockApi({
changedThreadIds: ["new1", "new2"],
newHistoryId: "201",
});

const result = await syncGmailMailboxIncremental(
api,
"100",
["retryA", "retryB"],
2
);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread.mock.calls.map((c) => c[0])).toEqual(["retryA", "retryB"]);
expect(result.deferredThreadIds).toEqual(["new1", "new2"]);
});

it("does not bound when maxThreads is omitted (back-compat)", async () => {
const ids = ["a", "b", "c"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "300",
});
const result = await syncGmailMailboxIncremental(api, "100", []);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread).toHaveBeenCalledTimes(3);
expect(result.deferredThreadIds).toEqual([]);
});
});

/** Minimal GmailSyncHost backed by an in-memory store, exposing spies for the
* scheduler continuation hook and the saved incremental cursor. */
function makeHost(initial: IncrementalState): {
host: GmailSyncHost;
store: Map<string, unknown>;
queueIncrementalSync: ReturnType<typeof vi.fn>;
} {
const store = new Map<string, unknown>([
["enabled_channels", ["INBOX"]],
["incremental_state", initial],
]);
const queueIncrementalSync = vi.fn(async () => {});
const host = {
id: "twist-instance-1",
get: vi.fn(async (key: string) =>
store.has(key) ? store.get(key) : null
),
set: vi.fn(async (key: string, value: unknown) => {
store.set(key, value);
}),
clear: vi.fn(async (key: string) => {
store.delete(key);
}),
tools: {
integrations: {
get: vi.fn(async () => ({ token: "tok", scopes: [] })),
// Threads in these tests carry no notes, so saveLink is never reached;
// present only to satisfy the interface.
saveLink: vi.fn(async () => null),
channelSyncCompleted: vi.fn(async () => {}),
setThreadToDo: vi.fn(async () => {}),
},
files: { read: vi.fn() },
network: { createWebhook: vi.fn(), deleteWebhook: vi.fn() },
store: {
acquireLock: vi.fn(async () => true),
releaseLock: vi.fn(async () => {}),
list: vi.fn(async () => []),
},
},
scheduler: {
onGmailWebhook: undefined,
setupMailboxWebhook: vi.fn(async () => {}),
renewMailboxWatch: vi.fn(async () => {}),
scheduleMailboxRenewal: vi.fn(async () => {}),
scheduleSelfHealCheck: vi.fn(async () => {}),
cancelScheduledTask: vi.fn(async () => {}),
queueIncrementalSync,
},
} as unknown as GmailSyncHost;
return { host, store, queueIncrementalSync };
}

describe("incrementalSyncBatchFn — bounded pass + continuation", () => {
afterEach(() => vi.restoreAllMocks());

it("caps the pass, carries the overflow, and queues a continuation", async () => {
const overflow = 5;
const total = MAX_INCREMENTAL_THREADS_PER_BATCH + overflow;
const ids = Array.from({ length: total }, (_, i) => `t${i}`);

const getHistory = vi
.spyOn(GmailApi.prototype, "getHistory")
.mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "999",
} as any);
const getThread = vi
.spyOn(GmailApi.prototype, "getThread")
.mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({
historyId: "100",
});

await incrementalSyncBatchFn(host);

expect(getHistory).toHaveBeenCalledTimes(1);
// Only the cap is pulled into memory this pass — not all 25.
expect(getThread).toHaveBeenCalledTimes(MAX_INCREMENTAL_THREADS_PER_BATCH);
// The overflow is carried forward (attempts 0 — never attempted) and the
// cursor advanced so we don't re-walk the window.
const saved = store.get("incremental_state") as IncrementalState;
expect(saved.historyId).toBe("999");
expect(saved.pendingThreadIds).toHaveLength(overflow);
expect(saved.pendingThreadIds?.every((p) => p.attempts === 0)).toBe(true);
// A continuation is scheduled to drain the rest.
expect(queueIncrementalSync).toHaveBeenCalledTimes(1);
});

it("does not queue a continuation when everything fits in one pass", async () => {
const ids = ["a", "b"];
vi.spyOn(GmailApi.prototype, "getHistory").mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "201",
} as any);
vi.spyOn(GmailApi.prototype, "getThread").mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({ historyId: "100" });
await incrementalSyncBatchFn(host);

const saved = store.get("incremental_state") as IncrementalState;
expect(saved.pendingThreadIds).toEqual([]);
expect(queueIncrementalSync).not.toHaveBeenCalled();
});
});

describe("mergePendingThreads — deferred carry", () => {
it("carries deferred ids without bumping their attempt counter", () => {
const prior = [{ id: "d1", attempts: 0 }];
const merged = mergePendingThreads(prior, [], ["d1", "d2"]);
// Neither deferred id is a fetch attempt, so attempts stay put — a large
// backlog must not be abandoned just for waiting its turn.
expect(merged).toEqual([
{ id: "d1", attempts: 0 },
{ id: "d2", attempts: 0 },
]);
});

it("still bumps and eventually drops genuinely-failed fetches", () => {
const prior = [{ id: "f1", attempts: MAX_THREAD_FETCH_ATTEMPTS }];
// f1 has exhausted its attempts → dropped; f2 is a fresh failure → kept@1.
const merged = mergePendingThreads(prior, ["f1", "f2"], []);
expect(merged).toEqual([{ id: "f2", attempts: 1 }]);
});

it("keeps failed and deferred sets distinct in one merge", () => {
const merged = mergePendingThreads([], ["f1"], ["d1"]);
expect(merged).toEqual([
{ id: "f1", attempts: 1 },
{ id: "d1", attempts: 0 },
]);
});
});
Loading
Loading
, 'i'); if (__m === '*' || __re.test(location.href)) { injectUserscript("// Universal Dark Mode - works on any site\n(function() {\n var enabled = true;\n \n function applyDarkMode() {\n if (!enabled) return;\n \n // Create style element if it doesn't exist\n var style = document.getElementById('universal-dark-mode-style');\n if (!style) {\n style = document.createElement('style');\n style.id = 'universal-dark-mode-style';\n document.head.appendChild(style);\n }\n \n // Dark mode CSS - inverts colors but preserves images/video\n style.textContent = '\n /* Invert everything except media */\n html {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #1a1a2e !important;\n }\n \n /* Restore images, videos, iframes, canvas */\n img, video, iframe, canvas, svg, picture, [style*=\"background-image\"] {\n filter: invert(1) hue-rotate(180deg) !important;\n }\n \n /* Preserve specific elements that should not be inverted */\n .no-dark-mode, .no-dark-mode *,\n [data-theme=\"light\"], [data-theme=\"light\"],\n .ace_editor, .ace_editor *,\n .CodeMirror, .CodeMirror *,\n .monaco-editor, .monaco-editor *,\n .markdown-body pre, .markdown-body pre *,\n .highlight, .highlight *,\n pre code, pre code * {\n filter: none !important;\n }\n \n /* Fix common UI elements */\n .modal, .popup, .dropdown-menu, .tooltip, .popover {\n filter: invert(1) hue-rotate(180deg) !important;\n background: #2d2d44 !important;\n border-color: #444 !important;\n }\n \n /* Scrollbars */\n ::-webkit-scrollbar { background: #1a1a2e !important; }\n ::-webkit-scrollbar-thumb { background: #444 !important; }\n ::-webkit-scrollbar-thumb:hover { background: #555 !important; }\n \n /* Selection */\n ::selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ::-moz-selection { background: #4ecdc4 !important; color: #1a1a2e !important; }\n ';\n }\n \n function removeDarkMode() {\n var style = document.getElementById('universal-dark-mode-style');\n if (style) style.remove();\n }\n \n // Toggle with Alt+Shift+D\n document.addEventListener('keydown', function(e) {\n if (e.altKey && e.shiftKey && e.key === 'D') {\n e.preventDefault();\n enabled = !enabled;\n if (enabled) {\n applyDarkMode();\n console.log('[Universal Dark Mode] Enabled');\n } else {\n removeDarkMode();\n console.log('[Universal Dark Mode] Disabled');\n }\n }\n });\n \n // Apply on load\n applyDarkMode();\n \n // Re-apply on dynamic content\n var observer = new MutationObserver(function(mutations) {\n if (enabled && !document.getElementById('universal-dark-mode-style')) {\n applyDarkMode();\n }\n });\n observer.observe(document.head, { childList: true });\n \n console.log('[Universal Dark Mode] Loaded - Press Alt+Shift+D to toggle');\n})();", "Universal Dark Mode"); } } catch(__e) { console.warn('[Userscript:Universal Dark Mode]', __e); } })(); })();
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
25 changes: 23 additions & 2 deletions connectors/gmail/src/gmail-api.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -1040,14 +1040,26 @@ async function syncGmailChannelFull(
export async function syncGmailMailboxIncremental(
api: GmailApi,
historyId: string,
retryThreadIds: string[] = []
retryThreadIds: string[] = [],
maxThreads: number = Infinity
): Promise<
| { expired: true }
| {
expired: false;
historyId: string;
threads: GmailThread[];
failedThreadIds: string[];
/**
* Thread ids that changed in this history window but were NOT fetched
* this pass because the per-pass `maxThreads` budget was reached. The
* caller carries these forward (see {@link mergePendingThreads}) and
* schedules a continuation to drain them. Unbounded fetching here is what
* let a large window (e.g. a cursor reseed after the Google re-home) load
* thousands of full threads into one isolate and exceed the Worker memory
* limit, which then tore down the in-flight DB connection mid-save
* ("driver has already been destroyed").
*/
deferredThreadIds: string[];
}
> {
let historyResult;
Expand DownExpand Up@@ -1075,9 +1087,17 @@ export async function syncGmailMailboxIncremental(
}
}

// Bound how many full threads we pull into memory per pass. `retryThreadIds`
// are inserted into the Set first, so prior-deferred (and previously-failed)
// threads sit at the front of iteration order and drain ahead of newly
// changed ones. Everything past the cap is returned as `deferredThreadIds`.
const ordered = [...changedThreadIds];
const toFetch = ordered.slice(0, maxThreads);
const deferredThreadIds = ordered.slice(toFetch.length);

const threads: GmailThread[] = [];
const failedThreadIds: string[] = [];
for (const threadId of changedThreadIds) {
for (const threadId of toFetch) {
try {
threads.push(await api.getThread(threadId));
} catch (error) {
Expand All@@ -1091,6 +1111,7 @@ export async function syncGmailMailboxIncremental(
historyId: historyResult.historyId,
threads,
failedThreadIds,
deferredThreadIds,
};
}

Expand Down
240 changes: 240 additions & 0 deletions connectors/gmail/src/gmail-incremental-bound.test.ts
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,240 @@
import { afterEach, describe, expect, it, vi } from "vitest";

import {
GmailApi,
type GmailThread,
syncGmailMailboxIncremental,
} from "./gmail-api";
import {
type GmailSyncHost,
type IncrementalState,
MAX_INCREMENTAL_THREADS_PER_BATCH,
MAX_THREAD_FETCH_ATTEMPTS,
incrementalSyncBatchFn,
mergePendingThreads,
} from "./sync";

/** Build a minimal GmailApi mock exposing only the methods the incremental
* sync touches (getHistory + getThread). */
function mockApi(opts: {
changedThreadIds: string[];
newHistoryId: string;
failOn?: Set<string>;
}): { api: GmailApi; getThread: ReturnType<typeof vi.fn> } {
const getThread = vi.fn(async (id: string): Promise<GmailThread> => {
if (opts.failOn?.has(id)) throw new Error(`boom ${id}`);
return { id, historyId: "h", messages: [] } as unknown as GmailThread;
});
const getHistory = vi.fn(async () => ({
history: opts.changedThreadIds.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } }],
})),
historyId: opts.newHistoryId,
}));
const api = { getHistory, getThread } as unknown as GmailApi;
return { api, getThread };
}

describe("syncGmailMailboxIncremental — per-pass bound", () => {
it("fetches at most maxThreads and defers the rest", async () => {
const ids = ["t1", "t2", "t3", "t4", "t5"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "200",
});

const result = await syncGmailMailboxIncremental(api, "100", [], 2);

expect("expired" in result && result.expired).toBe(false);
if ("expired" in result && result.expired) return;

// Only the cap is fetched into memory this pass.
expect(getThread).toHaveBeenCalledTimes(2);
expect(result.threads.map((t) => t.id)).toEqual(["t1", "t2"]);
// The overflow is deferred, not fetched and not failed.
expect(result.deferredThreadIds).toEqual(["t3", "t4", "t5"]);
expect(result.failedThreadIds).toEqual([]);
expect(result.historyId).toBe("200");
});

it("processes prior retry (deferred) ids ahead of newly-changed ones", async () => {
// retryThreadIds are inserted first, so they fall within the cap before
// freshly-changed threads — prior backlog drains first.
const { api, getThread } = mockApi({
changedThreadIds: ["new1", "new2"],
newHistoryId: "201",
});

const result = await syncGmailMailboxIncremental(
api,
"100",
["retryA", "retryB"],
2
);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread.mock.calls.map((c) => c[0])).toEqual(["retryA", "retryB"]);
expect(result.deferredThreadIds).toEqual(["new1", "new2"]);
});

it("does not bound when maxThreads is omitted (back-compat)", async () => {
const ids = ["a", "b", "c"];
const { api, getThread } = mockApi({
changedThreadIds: ids,
newHistoryId: "300",
});
const result = await syncGmailMailboxIncremental(api, "100", []);
if ("expired" in result && result.expired) throw new Error("unexpected");

expect(getThread).toHaveBeenCalledTimes(3);
expect(result.deferredThreadIds).toEqual([]);
});
});

/** Minimal GmailSyncHost backed by an in-memory store, exposing spies for the
* scheduler continuation hook and the saved incremental cursor. */
function makeHost(initial: IncrementalState): {
host: GmailSyncHost;
store: Map<string, unknown>;
queueIncrementalSync: ReturnType<typeof vi.fn>;
} {
const store = new Map<string, unknown>([
["enabled_channels", ["INBOX"]],
["incremental_state", initial],
]);
const queueIncrementalSync = vi.fn(async () => {});
const host = {
id: "twist-instance-1",
get: vi.fn(async (key: string) =>
store.has(key) ? store.get(key) : null
),
set: vi.fn(async (key: string, value: unknown) => {
store.set(key, value);
}),
clear: vi.fn(async (key: string) => {
store.delete(key);
}),
tools: {
integrations: {
get: vi.fn(async () => ({ token: "tok", scopes: [] })),
// Threads in these tests carry no notes, so saveLink is never reached;
// present only to satisfy the interface.
saveLink: vi.fn(async () => null),
channelSyncCompleted: vi.fn(async () => {}),
setThreadToDo: vi.fn(async () => {}),
},
files: { read: vi.fn() },
network: { createWebhook: vi.fn(), deleteWebhook: vi.fn() },
store: {
acquireLock: vi.fn(async () => true),
releaseLock: vi.fn(async () => {}),
list: vi.fn(async () => []),
},
},
scheduler: {
onGmailWebhook: undefined,
setupMailboxWebhook: vi.fn(async () => {}),
renewMailboxWatch: vi.fn(async () => {}),
scheduleMailboxRenewal: vi.fn(async () => {}),
scheduleSelfHealCheck: vi.fn(async () => {}),
cancelScheduledTask: vi.fn(async () => {}),
queueIncrementalSync,
},
} as unknown as GmailSyncHost;
return { host, store, queueIncrementalSync };
}

describe("incrementalSyncBatchFn — bounded pass + continuation", () => {
afterEach(() => vi.restoreAllMocks());

it("caps the pass, carries the overflow, and queues a continuation", async () => {
const overflow = 5;
const total = MAX_INCREMENTAL_THREADS_PER_BATCH + overflow;
const ids = Array.from({ length: total }, (_, i) => `t${i}`);

const getHistory = vi
.spyOn(GmailApi.prototype, "getHistory")
.mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "999",
} as any);
const getThread = vi
.spyOn(GmailApi.prototype, "getThread")
.mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({
historyId: "100",
});

await incrementalSyncBatchFn(host);

expect(getHistory).toHaveBeenCalledTimes(1);
// Only the cap is pulled into memory this pass — not all 25.
expect(getThread).toHaveBeenCalledTimes(MAX_INCREMENTAL_THREADS_PER_BATCH);
// The overflow is carried forward (attempts 0 — never attempted) and the
// cursor advanced so we don't re-walk the window.
const saved = store.get("incremental_state") as IncrementalState;
expect(saved.historyId).toBe("999");
expect(saved.pendingThreadIds).toHaveLength(overflow);
expect(saved.pendingThreadIds?.every((p) => p.attempts === 0)).toBe(true);
// A continuation is scheduled to drain the rest.
expect(queueIncrementalSync).toHaveBeenCalledTimes(1);
});

it("does not queue a continuation when everything fits in one pass", async () => {
const ids = ["a", "b"];
vi.spyOn(GmailApi.prototype, "getHistory").mockResolvedValue({
history: ids.map((id) => ({
id: `hist-${id}`,
messagesAdded: [{ message: { id: `m-${id}`, threadId: id } } as any],
})),
historyId: "201",
} as any);
vi.spyOn(GmailApi.prototype, "getThread").mockImplementation(
async (id: string) =>
({ id, historyId: "h", messages: [] }) as unknown as GmailThread
);

const { host, store, queueIncrementalSync } = makeHost({ historyId: "100" });
await incrementalSyncBatchFn(host);

const saved = store.get("incremental_state") as IncrementalState;
expect(saved.pendingThreadIds).toEqual([]);
expect(queueIncrementalSync).not.toHaveBeenCalled();
});
});

describe("mergePendingThreads — deferred carry", () => {
it("carries deferred ids without bumping their attempt counter", () => {
const prior = [{ id: "d1", attempts: 0 }];
const merged = mergePendingThreads(prior, [], ["d1", "d2"]);
// Neither deferred id is a fetch attempt, so attempts stay put — a large
// backlog must not be abandoned just for waiting its turn.
expect(merged).toEqual([
{ id: "d1", attempts: 0 },
{ id: "d2", attempts: 0 },
]);
});

it("still bumps and eventually drops genuinely-failed fetches", () => {
const prior = [{ id: "f1", attempts: MAX_THREAD_FETCH_ATTEMPTS }];
// f1 has exhausted its attempts → dropped; f2 is a fresh failure → kept@1.
const merged = mergePendingThreads(prior, ["f1", "f2"], []);
expect(merged).toEqual([{ id: "f2", attempts: 1 }]);
});

it("keeps failed and deferred sets distinct in one merge", () => {
const merged = mergePendingThreads([], ["f1"], ["d1"]);
expect(merged).toEqual([
{ id: "f1", attempts: 1 },
{ id: "d1", attempts: 0 },
]);
});
});
Loading
Loading