perf: drain background outbox concurrently (#297)
All checks were successful
CI / lint (pull_request) Successful in 27s
CI / build-frontend (pull_request) Successful in 4s

This commit is contained in:
timmy 2026-08-08 11:55:58 +00:00
parent d0b8c3dc0d
commit 43c7a85eec
4 changed files with 206 additions and 10 deletions

View File

@ -96,6 +96,32 @@ function createIssueSyncStore({ transaction, indexedDB = globalThis.indexedDB, n
});
}
async function claimBatch(ownerLogin, limit = 70) {
return transact(async records => {
const timestamp = Number(now());
const eligible = (await records.getAll()).filter(candidate =>
candidate.recordType !== 'receipt-preference' &&
candidate.ownerLogin === ownerLogin &&
(candidate.status === 'queued' || candidate.status === 'sending') &&
(candidate.status !== 'sending' || Number(candidate.claimUntil) <= timestamp));
const lanes = {
issue: eligible.filter(item => (item.outboxLane || 'issue') !== 'authored'),
authored: eligible.filter(item => item.outboxLane === 'authored'),
};
const selected = [];
const maximum = Math.max(0, Number(limit) || 0);
while (selected.length < maximum && (lanes.issue.length || lanes.authored.length)) {
if (lanes.issue.length) selected.push(lanes.issue.shift());
if (selected.length < maximum && lanes.authored.length) selected.push(lanes.authored.shift());
}
const claimed = selected.map(item => ({
...item, status: 'sending', claimUntil: timestamp + claimMs,
}));
for (const item of claimed) await records.put(item);
return claimed;
});
}
async function claim(id, ownerLogin) {
return transact(async records => {
const timestamp = Number(now());
@ -153,6 +179,7 @@ function createIssueSyncStore({ transaction, indexedDB = globalThis.indexedDB, n
upsert,
claim,
claimNext,
claimBatch,
complete: id => update(id, item => ({ ...item, status: 'sent', claimUntil: 0 })),
release: id => update(id, item => ({ ...item, status: 'queued', claimUntil: 0 })),
fail: (id, error) => update(id, item => ({ ...item, status: 'attention', claimUntil: 0, error })),
@ -167,7 +194,10 @@ function createIssueSyncStore({ transaction, indexedDB = globalThis.indexedDB, n
};
}
function createBackgroundIssueSync({ store, fetchJson, base = '' }) {
function createBackgroundIssueSync({
store, fetchJson, base = '', maxConcurrency = 3, batchSize = 70,
batch = work => work(),
}) {
let purgeRequested = false;
let activeFlush = null;
function receiptFor(item, status, delivered = {}) {
@ -272,14 +302,41 @@ function createBackgroundIssueSync({ store, fetchJson, base = '' }) {
const receipts = [];
let attention = 0;
if (!login) return { confirmed, blocked: 0, attention, login, receipts };
while (true) {
const item = await store.claimNext(login);
if (!item) break;
const result = await deliver(item);
const collect = result => {
if (result.issue) confirmed.push(result.issue);
if (result.message) confirmed.push(result.message);
if (result.attention) attention += 1;
if (result.receipt) receipts.push(result.receipt);
};
if (store.claimBatch) {
const claimed = await store.claimBatch(login, batchSize);
let next = 0;
let authenticationError = null;
let transientError = null;
const worker = async () => {
while (!authenticationError && next < claimed.length) {
const item = claimed[next++];
try {
collect(await deliver(item));
} catch (error) {
if (Number(error?.status || 0) === 401) authenticationError = error;
else if (!transientError) transientError = error;
}
}
};
const concurrency = Math.max(1, Math.min(Number(maxConcurrency) || 1, claimed.length));
await Promise.all(Array.from({ length: concurrency }, worker));
if (authenticationError) {
await Promise.all(claimed.slice(next).map(item => store.release(item.id)));
throw authenticationError;
}
if (transientError) throw transientError;
} else {
while (true) {
const item = await store.claimNext(login);
if (!item) break;
collect(await deliver(item));
}
}
const blocked = store.countBlocked ? await store.countBlocked(login) : 0;
return { confirmed, blocked, attention, login, receipts };
@ -288,7 +345,7 @@ function createBackgroundIssueSync({ store, fetchJson, base = '' }) {
function flush() {
if (purgeRequested) return Promise.resolve({ confirmed: [], blocked: 0, attention: 0, login: '', receipts: [] });
if (activeFlush) return activeFlush;
activeFlush = runFlush().finally(() => { activeFlush = null; });
activeFlush = Promise.resolve().then(() => batch(runFlush)).finally(() => { activeFlush = null; });
return activeFlush;
}

View File

@ -1,6 +1,6 @@
const BASE = new URL('./', self.location.href).pathname;
importScripts(BASE + 'static/background-issue-sync.js');
const CACHE = 'stackchain-dashboard-shell-v26';
const CACHE = 'stackchain-dashboard-shell-v27';
const OUTAGE_STATUSES = new Set([500, 502, 503, 504]);
const SHELL = [
BASE,
@ -45,12 +45,22 @@ async function sessionCsrf() {
return typeof payload.csrf_token === 'string' ? payload.csrf_token : '';
}
let batchedCsrf = null;
async function withSessionCsrf(work) {
batchedCsrf = sessionCsrf();
try {
return await work();
} finally {
batchedCsrf = null;
}
}
async function fetchJson(url, options = {}) {
const requestOptions = { ...options };
const method = String(options.method || 'GET').toUpperCase();
if (!['GET', 'HEAD', 'OPTIONS'].includes(method)) {
const headers = new Headers(options.headers || {});
const csrf = await sessionCsrf();
const csrf = await (batchedCsrf || sessionCsrf());
if (csrf) headers.set('X-CSRF-Token', csrf);
requestOptions.headers = headers;
}
@ -65,7 +75,7 @@ async function fetchJson(url, options = {}) {
}
const issueSync = self.__issueSync || createBackgroundIssueSync({
store: createIssueSyncStore(), fetchJson, base: BASE,
store: createIssueSyncStore(), fetchJson, base: BASE, batch: withSessionCsrf,
});
async function flushAndNotify() {

View File

@ -359,6 +359,126 @@ const fetchJson = async url => {{
assert output["error"] == "Gitea unavailable"
def test_bounded_flush_delivers_healthy_records_after_transient_failure_without_same_run_retry():
script = f"""
const createBackgroundIssueSync = require({json.dumps(str(SYNC))});
const items = [
{{id:'issue-1',operationId:'issue-1',ownerLogin:'timmy',outboxLane:'issue',repository:'o/r',title:'One',body:'',labelIds:[]}},
{{id:'message-1',operationId:'message-1',ownerLogin:'timmy',outboxLane:'authored',kind:'issue-comment',repository:'o/r',number:1,body:'Reply'}},
{{id:'issue-bad',operationId:'issue-bad',ownerLogin:'timmy',outboxLane:'issue',repository:'o/r',title:'Bad',body:'',labelIds:[]}},
{{id:'message-2',operationId:'message-2',ownerLogin:'timmy',outboxLane:'authored',kind:'issue-comment',repository:'o/r',number:2,body:'Reply'}},
{{id:'issue-2',operationId:'issue-2',ownerLogin:'timmy',outboxLane:'issue',repository:'o/r',title:'Two',body:'',labelIds:[]}},
];
const state = {{active:0,maxActive:0,attempted:[],completed:[],released:[],batches:0}};
const store = {{
claimBatch: async (owner, limit) => items.slice(0, limit),
complete: async id => state.completed.push(id),
release: async id => state.released.push(id),
fail: async () => {{}}, countBlocked: async () => 0,
}};
const fetchJson = async (url, options={{}}) => {{
if (url === 'api/v1/background-identity') return {{login:'timmy'}};
const key = options.headers['Idempotency-Key'];
state.attempted.push(key); state.active += 1;
state.maxActive = Math.max(state.maxActive, state.active);
await new Promise(resolve => setTimeout(resolve, key === 'issue-bad' ? 5 : 20));
state.active -= 1;
if (key === 'issue-bad') {{ const error=new Error('Temporary outage'); error.status=503; throw error; }}
return {{id:key,number:7}};
}};
(async () => {{
let error;
try {{ await createBackgroundIssueSync({{
store,fetchJson,maxConcurrency:3,batchSize:10,
batch: async work => {{ state.batches += 1; return work(); }},
}}).flush(); }}
catch (caught) {{ error=caught.message; }}
process.stdout.write(JSON.stringify({{state,error}}));
}})();
"""
output = run_node(script)
assert output["state"]["maxActive"] == 3
assert output["state"]["batches"] == 1
assert sorted(output["state"]["attempted"]) == [
"issue-1", "issue-2", "issue-bad", "message-1", "message-2"
]
assert len(output["state"]["attempted"]) == len(set(output["state"]["attempted"]))
assert sorted(output["state"]["completed"]) == [
"issue-1", "issue-2", "message-1", "message-2"
]
assert output["state"]["released"] == ["issue-bad"]
assert output["error"] == "Temporary outage"
def test_bounded_flush_stops_admission_and_releases_unstarted_claims_on_auth_loss():
script = f"""
const createBackgroundIssueSync = require({json.dumps(str(SYNC))});
const items = Array.from({{length:5}}, (_, index) => ({{
id:'item-'+index,operationId:'item-'+index,ownerLogin:'timmy',outboxLane:'issue',
repository:'o/r',title:'Item',body:'',labelIds:[],
}}));
const state = {{attempted:[],released:[],completed:[]}};
const store = {{
claimBatch: async () => items,
complete: async id => state.completed.push(id),
release: async id => state.released.push(id), fail:async()=>{{}}, countBlocked:async()=>0,
}};
const fetchJson = async (url, options={{}}) => {{
if (url === 'api/v1/background-identity') return {{login:'timmy'}};
const key=options.headers['Idempotency-Key']; state.attempted.push(key);
if (key === 'item-0') {{ const error=new Error('Authentication required'); error.status=401; throw error; }}
await new Promise(resolve => setTimeout(resolve, 20));
return {{number:1}};
}};
(async()=>{{
let error;
try {{ await createBackgroundIssueSync({{store,fetchJson,maxConcurrency:2}}).flush(); }}
catch (caught) {{ error={{message:caught.message,status:caught.status}}; }}
process.stdout.write(JSON.stringify({{state,error}}));
}})();
"""
output = run_node(script)
assert output["state"]["attempted"] == ["item-0", "item-1"]
assert sorted(output["state"]["released"]) == ["item-0", "item-2", "item-3", "item-4"]
assert output["state"]["completed"] == ["item-1"]
assert output["error"] == {"message": "Authentication required", "status": 401}
def test_issue_store_claims_one_finite_batch_fairly_across_outbox_lanes():
script = f"""
const createBackgroundIssueSync = require({json.dumps(str(SYNC))});
const records=new Map();let tail=Promise.resolve();
const transaction=work=>{{const run=tail.then(()=>work({{
getAll:async()=>[...records.values()].map(value=>({{...value}})),
put:async value=>records.set(value.id,{{...value}}),delete:async id=>records.delete(id),
}}));tail=run.catch(()=>{{}});return run;}};
(async()=>{{
const store=createBackgroundIssueSync.createIssueSyncStore({{transaction,now:()=>100,claimMs:5000}});
await store.reconcile([
{{id:'issue-1',ownerLogin:'timmy',status:'queued'}},
{{id:'issue-2',ownerLogin:'timmy',status:'queued'}},
{{id:'issue-3',ownerLogin:'timmy',status:'queued'}},
], 'issue');
await store.reconcile([
{{id:'message-1',ownerLogin:'timmy',status:'queued'}},
{{id:'message-2',ownerLogin:'timmy',status:'queued'}},
], 'authored');
const claimed=await store.claimBatch('timmy', 4);
process.stdout.write(JSON.stringify({{claimed,snapshot:await store.snapshot()}}));
}})();
"""
output = run_node(script)
assert [item["id"] for item in output["claimed"]] == [
"issue-1", "message-1", "issue-2", "message-2"
]
assert all(item["status"] == "sending" for item in output["claimed"])
assert all(item["claimUntil"] == 5100 for item in output["claimed"])
assert next(item for item in output["snapshot"] if item["id"] == "issue-3")["status"] == "queued"
def test_permanent_delivery_failure_marks_issue_for_foreground_attention():
script = f"""
const createBackgroundIssueSync = require({json.dumps(str(SYNC))});

View File

@ -95,7 +95,7 @@ async function dispatchNotificationClick(route) {{
def test_strict_browser_assets_ship_in_a_new_shell_cache():
source = WORKER.read_text()
assert "stackchain-dashboard-shell-v26" in source
assert "stackchain-dashboard-shell-v27" in source
assert "BASE + 'static/dashboard.css'" in source
assert "BASE + 'static/dashboard.js'" in source
assert "BASE + 'static/install-app.js'" in source
@ -185,6 +185,15 @@ def test_background_mutations_obtain_session_bound_csrf_proof():
assert "headers.set('X-CSRF-Token', csrf)" in source
def test_one_session_bound_csrf_proof_is_reused_for_a_background_drain():
source = WORKER.read_text()
assert "async function withSessionCsrf(work)" in source
assert "batchedCsrf = sessionCsrf()" in source
assert "await (batchedCsrf || sessionCsrf())" in source
assert "batch: withSessionCsrf" in source
def test_install_precaches_complete_subpath_scoped_app_shell():
result = run_worker_scenario(
"""