feat: batch answer submissions with auto-flush on threshold and inactivity (#408)

This commit is contained in:
Zack Pollard
2026-03-23 15:04:53 -04:00
committed by GitHub
parent 73099a9037
commit 6438b9c75e
4 changed files with 333 additions and 181 deletions
@@ -51,46 +51,36 @@ router.post('/api/verify', async (request, env) => {
return Response.json({ success: true });
});
router.post('/api/answers', async (request, env) => {
router.post('/api/answers/batch', async (request, env) => {
const db = env.DB;
const { questionId, value, otherText } = (await request.json()) as {
questionId: string;
value: string;
otherText?: string;
const { answers } = (await request.json()) as {
answers: Array<{ questionId: string; value: string; otherText?: string }>;
};
const ip =
request.headers.get('CF-Connecting-IP') ??
request.headers.get('x-forwarded-for') ??
'unknown';
let respondentId = getRespondentId(request);
const headers = new Headers();
if (!respondentId) {
respondentId = crypto.randomUUID();
setRespondentCookie(headers, respondentId);
if (!answers || !Array.isArray(answers) || answers.length === 0 || answers.length > 20) {
return new Response('Invalid answers payload', { status: 400 });
}
await db
.prepare(
`INSERT INTO respondents (id, ip_address, created_at) VALUES (?, ?, ?)
ON CONFLICT (id) DO UPDATE SET ip_address = excluded.ip_address`,
)
.bind(respondentId, ip, new Date().toISOString())
.run();
const respondentId = getRespondentId(request);
if (!respondentId) {
return new Response('No respondent cookie', { status: 400 });
}
await db
.prepare(
`INSERT INTO answers (respondent_id, question_id, answer, other_text, answered_at)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT (respondent_id, question_id)
DO UPDATE SET answer = excluded.answer, other_text = excluded.other_text, answered_at = excluded.answered_at`,
)
.bind(respondentId, questionId, value, otherText ?? null, new Date().toISOString())
.run();
const now = new Date().toISOString();
const statements = answers.map((a) =>
db
.prepare(
`INSERT INTO answers (respondent_id, question_id, answer, other_text, answered_at)
VALUES (?, ?, ?, ?, ?)
ON CONFLICT (respondent_id, question_id)
DO UPDATE SET answer = excluded.answer, other_text = excluded.other_text, answered_at = excluded.answered_at`,
)
.bind(respondentId, a.questionId, a.value, a.otherText ?? null, now),
);
return new Response(null, { status: 204, headers });
await db.batch(statements);
return new Response(null, { status: 204 });
});
router.get('/api/resume', async (request, env) => {
@@ -1,23 +1,21 @@
import { describe, it, expect, vi, beforeEach, afterEach } from 'vitest';
// Module will be imported after it exists
let saveAnswer: typeof import('./api-client').saveAnswer;
let fireAndForgetSave: typeof import('./api-client').fireAndForgetSave;
let flushPendingQueue: typeof import('./api-client').flushPendingQueue;
let getPendingQueue: typeof import('./api-client').getPendingQueue;
let bufferAnswer: typeof import('./api-client').bufferAnswer;
let flushBuffer: typeof import('./api-client').flushBuffer;
let flushBufferSync: typeof import('./api-client').flushBufferSync;
let getBufferSize: typeof import('./api-client').getBufferSize;
let fetchResume: typeof import('./api-client').fetchResume;
let postComplete: typeof import('./api-client').postComplete;
beforeEach(async () => {
vi.useFakeTimers();
vi.stubGlobal('fetch', vi.fn());
// Re-import module fresh each test to reset pendingQueue
vi.resetModules();
const mod = await import('./api-client');
saveAnswer = mod.saveAnswer;
fireAndForgetSave = mod.fireAndForgetSave;
flushPendingQueue = mod.flushPendingQueue;
getPendingQueue = mod.getPendingQueue;
bufferAnswer = mod.bufferAnswer;
flushBuffer = mod.flushBuffer;
flushBufferSync = mod.flushBufferSync;
getBufferSize = mod.getBufferSize;
fetchResume = mod.fetchResume;
postComplete = mod.postComplete;
});
@@ -27,129 +25,231 @@ afterEach(() => {
vi.restoreAllMocks();
});
describe('saveAnswer', () => {
it('calls fetch with correct body and headers', async () => {
describe('bufferAnswer', () => {
it('adds an answer to the buffer', () => {
bufferAnswer({ questionId: 'q1', value: 'test' });
expect(getBufferSize()).toBe(1);
});
it('overwrites duplicate questionIds', () => {
bufferAnswer({ questionId: 'q1', value: 'first' });
bufferAnswer({ questionId: 'q1', value: 'second' });
expect(getBufferSize()).toBe(1);
});
it('buffers multiple different questions', () => {
bufferAnswer({ questionId: 'q1', value: 'a' });
bufferAnswer({ questionId: 'q2', value: 'b' });
bufferAnswer({ questionId: 'q3', value: 'c' });
expect(getBufferSize()).toBe(3);
});
it('auto-flushes after 4 answers', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: true, status: 204 });
vi.stubGlobal('fetch', mockFetch);
const data = { questionId: 'q1', value: 'test' };
const promise = saveAnswer(data);
await vi.runAllTimersAsync();
await promise;
bufferAnswer({ questionId: 'q1', value: 'a' });
bufferAnswer({ questionId: 'q2', value: 'b' });
bufferAnswer({ questionId: 'q3', value: 'c' });
expect(mockFetch).not.toHaveBeenCalled();
expect(mockFetch).toHaveBeenCalledWith('/api/answers', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(data),
credentials: 'same-origin',
});
bufferAnswer({ questionId: 'q4', value: 'd' });
await vi.runAllTimersAsync();
expect(mockFetch).toHaveBeenCalledTimes(1);
const body = JSON.parse(mockFetch.mock.calls[0][1].body);
expect(body.answers).toHaveLength(4);
expect(getBufferSize()).toBe(0);
});
it('returns true on 204 response', async () => {
vi.stubGlobal('fetch', vi.fn().mockResolvedValue({ ok: true, status: 204 }));
it('auto-flushes after 10 seconds of inactivity', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: true, status: 204 });
vi.stubGlobal('fetch', mockFetch);
const promise = saveAnswer({ questionId: 'q1', value: 'test' });
bufferAnswer({ questionId: 'q1', value: 'a' });
bufferAnswer({ questionId: 'q2', value: 'b' });
expect(mockFetch).not.toHaveBeenCalled();
await vi.advanceTimersByTimeAsync(10_000);
expect(mockFetch).toHaveBeenCalledTimes(1);
const body = JSON.parse(mockFetch.mock.calls[0][1].body);
expect(body.answers).toHaveLength(2);
expect(getBufferSize()).toBe(0);
});
it('resets inactivity timer on each new answer', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: true, status: 204 });
vi.stubGlobal('fetch', mockFetch);
bufferAnswer({ questionId: 'q1', value: 'a' });
// Wait 8 seconds, then add another answer
await vi.advanceTimersByTimeAsync(8_000);
expect(mockFetch).not.toHaveBeenCalled();
bufferAnswer({ questionId: 'q2', value: 'b' });
// 8 more seconds — still within the new 10s window
await vi.advanceTimersByTimeAsync(8_000);
expect(mockFetch).not.toHaveBeenCalled();
// 2 more seconds — now 10s since last answer
await vi.advanceTimersByTimeAsync(2_000);
expect(mockFetch).toHaveBeenCalledTimes(1);
});
it('resets counter after threshold flush', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: true, status: 204 });
vi.stubGlobal('fetch', mockFetch);
// First batch of 4
bufferAnswer({ questionId: 'q1', value: 'a' });
bufferAnswer({ questionId: 'q2', value: 'b' });
bufferAnswer({ questionId: 'q3', value: 'c' });
bufferAnswer({ questionId: 'q4', value: 'd' });
await vi.runAllTimersAsync();
expect(mockFetch).toHaveBeenCalledTimes(1);
// Next 3 answers should NOT trigger a flush
bufferAnswer({ questionId: 'q5', value: 'e' });
bufferAnswer({ questionId: 'q6', value: 'f' });
bufferAnswer({ questionId: 'q7', value: 'g' });
expect(mockFetch).toHaveBeenCalledTimes(1);
// 4th answer triggers second flush
bufferAnswer({ questionId: 'q8', value: 'h' });
await vi.runAllTimersAsync();
expect(mockFetch).toHaveBeenCalledTimes(2);
});
it('duplicate answers do not inflate the flush counter', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: true, status: 204 });
vi.stubGlobal('fetch', mockFetch);
// Answer same question 4 times — should count as 4 towards threshold
// (user changed their mind rapidly)
bufferAnswer({ questionId: 'q1', value: 'a' });
bufferAnswer({ questionId: 'q1', value: 'b' });
bufferAnswer({ questionId: 'q1', value: 'c' });
bufferAnswer({ questionId: 'q1', value: 'd' });
await vi.runAllTimersAsync();
// Threshold is based on answer events, not unique questions
expect(mockFetch).toHaveBeenCalledTimes(1);
const body = JSON.parse(mockFetch.mock.calls[0][1].body);
expect(body.answers).toHaveLength(1); // deduped in buffer
});
});
describe('flushBuffer', () => {
it('returns true immediately when buffer is empty', async () => {
const result = await flushBuffer();
expect(result).toBe(true);
expect(fetch).not.toHaveBeenCalled();
});
it('sends buffered answers to /api/answers/batch', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: true, status: 204 });
vi.stubGlobal('fetch', mockFetch);
bufferAnswer({ questionId: 'q1', value: 'a' });
bufferAnswer({ questionId: 'q2', value: 'b' });
const promise = flushBuffer();
await vi.runAllTimersAsync();
const result = await promise;
expect(result).toBe(true);
expect(mockFetch).toHaveBeenCalledWith('/api/answers/batch', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: expect.any(String),
credentials: 'same-origin',
});
const body = JSON.parse(mockFetch.mock.calls[0][1].body);
expect(body.answers).toHaveLength(2);
expect(getBufferSize()).toBe(0);
});
it('retries on 500 response - fetch called 4 times total', async () => {
it('re-adds items to buffer on failure', async () => {
vi.stubGlobal('fetch', vi.fn().mockResolvedValue({ ok: false, status: 500 }));
bufferAnswer({ questionId: 'q1', value: 'a' });
const promise = flushBuffer();
await vi.runAllTimersAsync();
const result = await promise;
expect(result).toBe(false);
expect(getBufferSize()).toBe(1);
});
it('retries on 500 response - 4 attempts total', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: false, status: 500 });
vi.stubGlobal('fetch', mockFetch);
const promise = saveAnswer({ questionId: 'q1', value: 'test' });
bufferAnswer({ questionId: 'q1', value: 'a' });
const promise = flushBuffer();
await vi.runAllTimersAsync();
await promise;
expect(mockFetch).toHaveBeenCalledTimes(4); // 1 initial + 3 retries
expect(mockFetch).toHaveBeenCalledTimes(4);
});
it('does NOT retry on 400 response - fetch called once', async () => {
it('does NOT retry on 400 response', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: false, status: 400 });
vi.stubGlobal('fetch', mockFetch);
const promise = saveAnswer({ questionId: 'q1', value: 'test' });
bufferAnswer({ questionId: 'q1', value: 'a' });
const promise = flushBuffer();
await vi.runAllTimersAsync();
await promise;
expect(mockFetch).toHaveBeenCalledTimes(1);
});
it('retries on network error (fetch throws)', async () => {
const mockFetch = vi.fn().mockRejectedValue(new TypeError('Failed to fetch'));
vi.stubGlobal('fetch', mockFetch);
const promise = saveAnswer({ questionId: 'q1', value: 'test' });
await vi.runAllTimersAsync();
await promise;
expect(mockFetch).toHaveBeenCalledTimes(4); // 1 initial + 3 retries
});
it('uses delays of 1000, 2000, 4000 ms between retries', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: false, status: 500 });
vi.stubGlobal('fetch', mockFetch);
const promise = saveAnswer({ questionId: 'q1', value: 'test' });
// Initial call happens immediately
await vi.advanceTimersByTimeAsync(0);
expect(mockFetch).toHaveBeenCalledTimes(1);
// After 1000ms: first retry
await vi.advanceTimersByTimeAsync(1000);
expect(mockFetch).toHaveBeenCalledTimes(2);
// After 2000ms more: second retry
await vi.advanceTimersByTimeAsync(2000);
expect(mockFetch).toHaveBeenCalledTimes(3);
// After 4000ms more: third retry
await vi.advanceTimersByTimeAsync(4000);
expect(mockFetch).toHaveBeenCalledTimes(4);
await promise;
});
});
describe('fireAndForgetSave', () => {
it('adds to pendingQueue on total failure', async () => {
it('does not overwrite newer buffer entries on failure', async () => {
vi.stubGlobal('fetch', vi.fn().mockResolvedValue({ ok: false, status: 500 }));
fireAndForgetSave({ questionId: 'q1', value: 'test' });
await vi.runAllTimersAsync();
// Allow microtasks to settle
await vi.advanceTimersByTimeAsync(0);
bufferAnswer({ questionId: 'q1', value: 'old' });
expect(getPendingQueue()).toHaveLength(1);
expect(getPendingQueue()[0]).toEqual({ questionId: 'q1', value: 'test' });
const promise = flushBuffer();
// Buffer a newer answer while flush is in progress
bufferAnswer({ questionId: 'q1', value: 'new' });
await vi.runAllTimersAsync();
await promise;
// The buffer should still have the newer value, not the failed old one
expect(getBufferSize()).toBe(1);
});
});
describe('flushPendingQueue', () => {
it('re-fires queued items', async () => {
const mockFetch = vi.fn().mockResolvedValue({ ok: false, status: 500 });
vi.stubGlobal('fetch', mockFetch);
describe('flushBufferSync', () => {
it('uses sendBeacon when buffer has items', () => {
const mockSendBeacon = vi.fn().mockReturnValue(true);
vi.stubGlobal('navigator', { sendBeacon: mockSendBeacon });
// First: make a save fail to populate the queue
fireAndForgetSave({ questionId: 'q1', value: 'test' });
await vi.runAllTimersAsync();
await vi.advanceTimersByTimeAsync(0);
expect(getPendingQueue()).toHaveLength(1);
bufferAnswer({ questionId: 'q1', value: 'a' });
flushBufferSync();
// Now make fetch succeed and flush
mockFetch.mockResolvedValue({ ok: true, status: 204 });
const callsBefore = mockFetch.mock.calls.length;
flushPendingQueue();
await vi.runAllTimersAsync();
await vi.advanceTimersByTimeAsync(0);
expect(mockSendBeacon).toHaveBeenCalledWith('/api/answers/batch', expect.any(Blob));
expect(getBufferSize()).toBe(0);
});
// Queue should be empty now (re-fired successfully)
expect(getPendingQueue()).toHaveLength(0);
// fetch should have been called at least once more
expect(mockFetch.mock.calls.length).toBeGreaterThan(callsBefore);
it('does nothing when buffer is empty', () => {
const mockSendBeacon = vi.fn();
vi.stubGlobal('navigator', { sendBeacon: mockSendBeacon });
flushBufferSync();
expect(mockSendBeacon).not.toHaveBeenCalled();
});
});
@@ -21,26 +21,60 @@ interface PendingSave {
otherText?: string;
}
let pendingQueue: PendingSave[] = [];
let answerBuffer: Map<string, PendingSave> = new Map();
let inactivityTimer: ReturnType<typeof setTimeout> | null = null;
let unflushedCount = 0;
const BACKOFF_DELAYS = [1000, 2000, 4000];
const INACTIVITY_MS = 10_000;
const FLUSH_THRESHOLD = 4;
let onSaveErrorCallback: ((message: string) => void) | null = null;
/**
* Saves an answer to the server with exponential backoff retry.
* Retries up to 3 times on server errors (5xx) and network errors.
* Does NOT retry on client errors (4xx).
* Registers a callback invoked when a batch save fails after all retries.
*/
export async function saveAnswer(data: PendingSave): Promise<boolean> {
export function onSaveError(cb: (message: string) => void): void {
onSaveErrorCallback = cb;
}
function resetInactivityTimer() {
if (inactivityTimer !== null) {
clearTimeout(inactivityTimer);
}
inactivityTimer = setTimeout(() => {
flushBuffer();
}, INACTIVITY_MS);
}
/**
* Buffers an answer for later batch submission.
* Auto-flushes after 4 answers or 10 seconds of inactivity.
*/
export function bufferAnswer(data: PendingSave): void {
answerBuffer.set(data.questionId, data);
unflushedCount++;
resetInactivityTimer();
if (unflushedCount >= FLUSH_THRESHOLD) {
flushBuffer();
}
}
/**
* Sends a batch of answers to the server with exponential backoff retry.
*/
async function saveBatch(answers: PendingSave[]): Promise<boolean> {
for (let attempt = 0; attempt <= BACKOFF_DELAYS.length; attempt++) {
try {
const res = await fetch('/api/answers', {
const res = await fetch('/api/answers/batch', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify(data),
body: JSON.stringify({ answers }),
credentials: 'same-origin',
});
if (res.ok) return true;
if (res.status < 500) return false; // client error, don't retry
if (res.status < 500) return false;
} catch {
// network error, retry
}
@@ -51,44 +85,60 @@ export async function saveAnswer(data: PendingSave): Promise<boolean> {
return false;
}
let onSaveErrorCallback: ((message: string) => void) | null = null;
/**
* Registers a callback invoked when a save fails after all retries.
* Flushes the answer buffer to the server as a batch.
* Returns true on success, false on failure.
* On failure, items are re-added to the buffer for the next flush attempt.
*/
export function onSaveError(cb: (message: string) => void): void {
onSaveErrorCallback = cb;
}
export async function flushBuffer(): Promise<boolean> {
if (answerBuffer.size === 0) return true;
/**
* Fire-and-forget: saves an answer without blocking.
* On total failure (all retries exhausted), queues for later retry and notifies via callback.
*/
export function fireAndForgetSave(data: PendingSave): void {
saveAnswer(data).then((success) => {
if (!success) {
pendingQueue.push(data);
onSaveErrorCallback?.('Failed to save answer. Your response will be retried automatically.');
}
});
}
/**
* Re-attempts all queued saves that previously failed.
*/
export function flushPendingQueue(): void {
const queue = [...pendingQueue];
pendingQueue = [];
for (const item of queue) {
fireAndForgetSave(item);
if (inactivityTimer !== null) {
clearTimeout(inactivityTimer);
inactivityTimer = null;
}
unflushedCount = 0;
const batch = [...answerBuffer.values()];
answerBuffer.clear();
const success = await saveBatch(batch);
if (!success) {
// Re-add failed items, but don't overwrite newer entries
for (const item of batch) {
if (!answerBuffer.has(item.questionId)) {
answerBuffer.set(item.questionId, item);
}
}
onSaveErrorCallback?.('Failed to save answers. Your responses will be retried automatically.');
}
return success;
}
/**
* Returns the current pending queue (for testing).
* Synchronous flush for beforeunload — uses sendBeacon.
*/
export function getPendingQueue(): PendingSave[] {
return pendingQueue;
export function flushBufferSync(): void {
if (answerBuffer.size === 0) return;
if (inactivityTimer !== null) {
clearTimeout(inactivityTimer);
inactivityTimer = null;
}
unflushedCount = 0;
const batch = [...answerBuffer.values()];
answerBuffer.clear();
const blob = new Blob([JSON.stringify({ answers: batch })], { type: 'application/json' });
navigator.sendBeacon('/api/answers/batch', blob);
}
/**
* Returns the current buffer size (for testing).
*/
export function getBufferSize(): number {
return answerBuffer.size;
}
/**
@@ -3,7 +3,7 @@
import { onMount } from 'svelte';
import { Turnstile } from 'svelte-turnstile';
import { createSurveyEngine } from '$lib/survey-engine.svelte';
import { fetchResume, fireAndForgetSave, postComplete, onSaveError, verifyTurnstile } from '$lib/api-client';
import { fetchResume, bufferAnswer, flushBuffer, flushBufferSync, postComplete, onSaveError, verifyTurnstile } from '$lib/api-client';
import SurveyShell from '$lib/components/SurveyShell.svelte';
import ThankYouScreen from '$lib/components/ThankYouScreen.svelte';
import AlreadyCompleted from '$lib/components/AlreadyCompleted.svelte';
@@ -26,30 +26,42 @@
verifyTurnstile(token).catch(() => {});
}
onMount(async () => {
try {
const resume = await fetchResume();
if (resume.isComplete) {
alreadyCompleted = true;
} else if (resume.answers && resume.nextQuestionIndex !== undefined && resume.nextQuestionIndex > 0) {
needsVerification = !resume.isVerified;
engine.initialize(resume.answers, resume.nextQuestionIndex);
} else {
showWelcome = true;
const handleUnload = () => flushBufferSync();
onMount(() => {
(async () => {
try {
const resume = await fetchResume();
if (resume.isComplete) {
alreadyCompleted = true;
} else if (resume.answers && resume.nextQuestionIndex !== undefined && resume.nextQuestionIndex > 0) {
needsVerification = !resume.isVerified;
engine.initialize(resume.answers, resume.nextQuestionIndex);
} else {
showWelcome = true;
}
} catch (e) {
error = e instanceof Error ? e.message : 'Something went wrong. Please try again later.';
}
} catch (e) {
error = e instanceof Error ? e.message : 'Something went wrong. Please try again later.';
}
loading = false;
loading = false;
})();
window.addEventListener('beforeunload', handleUnload);
return () => window.removeEventListener('beforeunload', handleUnload);
});
function handleAnswer(questionId: string, value: string, otherText?: string) {
engine.setAnswer(questionId, value, otherText);
fireAndForgetSave({ questionId, value, otherText });
bufferAnswer({ questionId, value, otherText });
}
async function handleComplete() {
try {
const flushed = await flushBuffer();
if (!flushed) {
error = 'Failed to save your answers. Please try again.';
return;
}
await postComplete();
surveyFinished = true;
} catch (e) {