Repository navigation
Expand file tree
/
Copy pathreadCache.js
More file actions
367 lines (332 loc) · 15.5 KB
/
Copy pathreadCache.js
File metadata and controls
367 lines (332 loc) · 15.5 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
/**
* The console's read cache: upstream GET answers, per node, shared by every worker thread on this
* host and by every operator signed in through it.
*
* WHY THE CONSOLE AND NOT THE PRERENDER NODES. The plugin caches its analytics window per WORKER,
* and Harper hands each incoming connection to a worker (SO_REUSEPORT, else the most idle one), so
* a refresh is answered from that cache only when it happens to land on the worker that scanned. The
* console is the one place every operator's read passes through, and it already owns the fan-out:
* a repeat read answered here never reaches a node that is also serving crawlers.
*
* WHY A TABLE AND NOT A MAP. The same distribution applies to this component — a module-level Map
* would be per worker here too, and miss the same way. `ProxyRead` is a node-local table (not
* replicated, not audited, not exported) that every worker on this host reads. A crash loses it,
* which for a cache is the right trade: audited, every cached body would also be written to the
* transaction log and kept there for its retention.
*
* WHO MAY READ A HIT. Every data route on the plugin requires a super_user and no answer depends on
* which one asked — that is what makes one answer shareable between operators. But the console
* cannot validate a session token itself: the token is opaque, and the console cookie is not signed,
* so "sent a cookie" proves nothing. A hit for node X is served only to a requester whose token node X
* itself confirmed as a super_user within `VERIFY_MS` — by answering a data route 200, by the shell's
* own session check, or by a session check this cache makes when it holds neither. That is the
* uncached authorization, at most `VERIFY_MS` stale, and never weaker.
*
* WHAT A WRITE DOES. Every console POST that can change state stamps a write generation, and an entry
* whose fetch began before it is a miss: the reload after "Pause" or "Apply" reads the node, not the
* answer from before the click. Writes made elsewhere — another console, the plugin's own schedulers
* — are bounded by the TTL alone.
*/
import { createHash } from 'node:crypto';
/**
* How long an answer may be served, per route, measured from when its DATA was produced — see
* `bornAt` below. Only the routes that cost the node something, or cannot change faster, are listed.
*/
export const READ_TTL = Object.freeze({
// Harper aggregates analytics once a minute: re-reading inside that window cannot surface a new
// row, only repeat the scan. (The plugin's `management.analytics.cacheTtl` defaults to it too.)
'analytics': 60_000,
// Day-bucketed crawl sketches: a minute of staleness is invisible at that grain.
'crawl-breadth': 60_000,
// The metric catalog — static for a plugin release.
'metrics': 600_000,
// A capped page-cache scan on one of the node's two heavy slots per worker.
'pages': 15_000,
});
/**
* Everything else — overview, queue-state, the probe and purge status, config, invalidations,
* sitemaps, unrouted — is cheap on the node by design, and is what an operator watches change. A
* short TTL collapses a burst (a view switch, several operators, several tabs) without making a
* Refresh look stuck.
*/
export const DEFAULT_READ_TTL = 5_000;
/**
* Never cached. `session` IS the authorization check. `page-content` is a byte-exact download of one
* stored page: unique per request, and up to megabytes.
*/
export const UNCACHED_READS = Object.freeze(new Set(['session', 'page-content']));
export const ttlFor = (route) => (UNCACHED_READS.has(route) ? 0 : (READ_TTL[route] ?? DEFAULT_READ_TTL));
/**
* The POSTs that only read, and so leave the cache alone. Every other POST stamps the write
* generation — including dry runs, because a dry-run purge or sweep still STARTS a pass, and the
* status routes then read differently.
*/
export const READ_ONLY_POST = Object.freeze(new Set(['explain', 'schedule', 'sitemap']));
export const invalidatesReads = (route) => !READ_ONLY_POST.has(route);
/** How long a node's confirmation of one operator's token is trusted before it is asked again. */
export const VERIFY_MS = 60_000;
/**
* Larger answers (in bytes, checked before decoding) are served but not stored. Today's largest (a 24h analytics window) is tens of KB per
* node; anything near this is not a dashboard read.
*/
export const MAX_CACHED_BODY = 4 * 1024 * 1024;
/**
* The write-generation rows: ONE PER WORKER, `write-generation:<worker>`, each only ever moved forward
* by its own worker, and read as the max of all of them. A single shared row would be last-writer-wins:
* two POSTs on two workers can commit in the opposite order to their timestamps, leave the row at the
* earlier one, and let a third worker serve an entry fetched between them. None of these is a cache key
* — those are URLs, or `sha256:` digests.
*/
export const GENERATION_PREFIX = 'write-generation:';
/** Bound on remembered confirmations and revocations, per worker. One per operator per node in practice. */
const MAX_REMEMBERED = 1000;
/**
* Harper refuses a primary key past 1,978 bytes. A key is the URL the console would ask, which a page
* browse's prefix or cursor can lengthen without limit, so a long one is stored by its digest instead.
*/
const MAX_PLAIN_KEY_BYTES = 1024;
/** The cache key is the URL the console would ask: one origin, one route, one query. */
export const cacheKeyFor = (origin, path) => {
const url = `${origin}/prerender_admin/${path}`;
return Buffer.byteLength(url) <= MAX_PLAIN_KEY_BYTES
? url
: `sha256:${createHash('sha256').update(url).digest('base64url')}`;
};
const isJson = (contentType) => /\bjson\b/i.test(String(contentType ?? ''));
const parse = (text) => {
try {
return JSON.parse(text);
} catch {
return null;
}
};
const ageOf = (payload) =>
payload && typeof payload === 'object' && Number.isFinite(payload.cacheAgeMs) ? payload.cacheAgeMs : null;
/**
* Whether a stored entry may answer now.
*
* Freshness runs from `bornAt`, when the DATA was produced, not from when the console fetched it. The
* plugin caches analytics itself and says how old its answer was (`cacheAgeMs`); counting from the
* fetch would stack the two TTLs, and a "60s" window could be two minutes old. Data is never younger
* than its fetch, whatever a node's clock claimed. An entry the clock has gone backwards past is not
* fresh either: it would otherwise stay fresh for however far the clock moved.
*/
export function isFresh(entry, { ttl, now, generation }) {
if (!entry || typeof entry.body !== 'string') return false;
if (!Number.isFinite(entry.startedAt) || !Number.isFinite(entry.bornAt)) return false;
if (entry.startedAt < generation) return false;
if (now < entry.startedAt) return false;
return now - Math.min(entry.bornAt, entry.startedAt) < ttl;
}
/** A Map that keeps its newest `MAX_REMEMBERED` keys. */
const remember = (map, key, value) => {
map.delete(key);
map.set(key, value);
while (map.size > MAX_REMEMBERED) map.delete(map.keys().next().value);
};
/**
* The read path every proxied GET takes.
*
* `fetch(origin, path, cookie)` → `{ status, contentType, body: Buffer }`, throwing on transport
* failure; `verify(origin, cookie)` → whether that node accepts the token as a super_user, throwing
* likewise. `table()` returns the Harper table, or null where there is none (tests, or a host without
* the schema). `detach(fn)` runs `fn` outside the request's transaction (util/detach.js) — every table
* access goes through it: a write inside the request would be reaped with it, and a read would hold a
* read transaction open across the whole upstream wait, which Harper force-commits past 30s and then
* refuses further reads on. `workerId` names this worker's generation row.
*
* Every answer is `{ status, contentType, body, payload, cached, ageMs }`: `payload` is a parse of
* `body` that belongs to this caller alone (a merger may modify it), `undefined` for a non-JSON body,
* and null for an unreadable one; `ageMs` is how long ago the console fetched it, null when it just did.
*/
export function createReadCache({ table, fetch, verify, detach, workerId, enabled = () => true, now = Date.now, log }) {
const confirmed = new Map();
const confirming = new Map();
// When each token was last refused or signed out. Evidence gathered BEFORE that moment — a request
// already in flight at logout that answers 200 after it — never confirms the token again.
const revoked = new Map();
const inflight = new Map();
// This worker's own last write: its generation row's value, and seen by its own reads before the
// row has committed.
let localGeneration = 0;
const generationKey = `${GENERATION_PREFIX}${workerId}`;
const idOf = (origin, cookie) => createHash('sha256').update(origin).update('\0').update(cookie).digest('base64url');
/** `askedAt`: when the request that proved the token was SENT — what a revocation is compared with. */
const confirm = (origin, cookie, askedAt = now()) => {
if (!cookie) return;
const id = idOf(origin, cookie);
const refusedAt = revoked.get(id);
if (refusedAt !== undefined && askedAt <= refusedAt) return;
remember(confirmed, id, now());
};
const forget = (origin, cookie) => {
if (!cookie) return;
const id = idOf(origin, cookie);
confirmed.delete(id);
remember(revoked, id, now());
};
/**
* Whether `origin` currently accepts this token — remembered, else asked (once, however many wait).
* A transport failure REJECTS: the node did not answer, and the caller must not then send it a second
* request to time out on.
*/
const vouched = (origin, cookie) => {
if (!cookie) return false;
const id = idOf(origin, cookie);
const at = confirmed.get(id);
if (at !== undefined && now() - at < VERIFY_MS && now() >= at) return true;
let pending = confirming.get(id);
if (!pending) {
const askedAt = now();
pending = Promise.resolve()
.then(() => verify(origin, cookie))
.then((ok) => {
if (ok) confirm(origin, cookie, askedAt);
else forget(origin, cookie);
return !!ok;
})
.finally(() => confirming.delete(id));
confirming.set(id, pending);
}
return pending;
};
const answerOf = (raw, { cached = false, ageMs = null } = {}) => ({
status: raw.status,
contentType: raw.contentType,
body: raw.body,
payload: isJson(raw.contentType) ? parse(raw.body.toString('utf8')) : undefined,
cached,
ageMs,
});
const lookup = async (tableRef, key) => {
try {
return (await detach(() => tableRef.get(key))) ?? null;
} catch (e) {
log?.warn?.(`[prerender-console] read cache lookup failed; reading the node: ${e?.message ?? String(e)}`);
return null;
}
};
/** The newest write any worker on this host has stamped. */
const generationOf = async (tableRef) => {
let newest = localGeneration;
try {
await detach(async () => {
const rows = tableRef.search({
conditions: [{ attribute: 'key', comparator: 'starts_with', value: GENERATION_PREFIX }],
select: ['key', 'startedAt'],
});
for await (const row of rows) if (Number.isFinite(row?.startedAt)) newest = Math.max(newest, row.startedAt);
});
} catch (e) {
log?.warn?.(`[prerender-console] write generation unreadable; reading the node: ${e?.message ?? String(e)}`);
// Unknown generation: nothing cached may be trusted, so every entry reads as stale.
return Infinity;
}
return newest;
};
const hitOf = (entry, at) => {
const payload = parse(entry.body);
if (payload === null) return null;
let text = entry.body;
// The plugin's own age, carried forward: the footer's "cached Ns ago" must count the console's
// time on top, or a minute-old window would read as the plugin's few seconds.
if (ageOf(payload) !== null) {
payload.cacheAgeMs = at - Math.min(entry.bornAt, entry.startedAt);
text = JSON.stringify(payload);
}
return {
status: 200,
contentType: entry.contentType,
body: Buffer.from(text, 'utf8'),
payload,
cached: true,
ageMs: at - entry.startedAt,
};
};
const store = (key, raw, payload, startedAt) => {
const tableRef = table();
if (!tableRef || payload === null || payload === undefined) return;
if (raw.body.length > MAX_CACHED_BODY) return;
const text = raw.body.toString('utf8');
const record = {
body: text,
contentType: raw.contentType,
startedAt,
// A node whose clock runs behind can report a negative age; data is never younger than its fetch.
bornAt: startedAt - Math.max(0, ageOf(payload) ?? 0),
};
Promise.resolve(detach(() => tableRef.put(key, record))).catch((e) =>
log?.warn?.(`[prerender-console] read cache store failed for ${key}: ${e?.message ?? String(e)}`)
);
};
/** The leader's fetch: the one request that goes to the node, which every concurrent reader rides. */
const lead = async (key, origin, path, cookie) => {
const startedAt = now();
const flight = { startedAt, promise: fetch(origin, path, cookie) };
inflight.set(key, flight);
let raw;
try {
raw = await flight.promise;
} finally {
if (inflight.get(key) === flight) inflight.delete(key);
}
const answer = answerOf(raw);
if (raw.status === 200) {
confirm(origin, cookie, startedAt);
store(key, raw, answer.payload, startedAt);
} else if (raw.status === 401 || raw.status === 403) {
forget(origin, cookie);
}
return answer;
};
return {
confirm,
forget,
/** One read of one node: from the table, from a fetch already in flight, or from the node. */
async read(origin, route, query, cookie) {
const path = route + (query ? `?${query}` : '');
const ttl = ttlFor(route);
const tableRef = enabled() ? table() : null;
if (!tableRef || ttl <= 0) return answerOf(await fetch(origin, path, cookie));
const key = cacheKeyFor(origin, path);
// Asked at most once per read: a node that does not answer the session check costs one
// timeout, not one per path below.
let vouching = null;
const isVouched = () => (vouching ??= Promise.resolve(vouched(origin, cookie)));
const [entry, generation] = await Promise.all([lookup(tableRef, key), generationOf(tableRef)]);
if (isFresh(entry, { ttl, now: now(), generation }) && (await isVouched())) {
// Fresh when looked up is not enough: the session check can take up to a request timeout, and
// the entry must still be fresh — and its age true — at the moment it is served.
const servedAt = now();
const hit = isFresh(entry, { ttl, now: servedAt, generation }) ? hitOf(entry, servedAt) : null;
if (hit) return hit;
}
const flight = inflight.get(key);
if (flight && flight.startedAt >= generation && (await isVouched())) {
// A transport failure is the node's state at this moment, for every rider alike — rethrown,
// not retried. A 401/403 is about the LEADER's token, so a rider asks with its own.
const raw = await flight.promise;
if (raw.status !== 401 && raw.status !== 403) return answerOf(raw);
}
return lead(key, origin, path, cookie);
},
/**
* A write went through this console: every entry fetched before now is stale. Awaited by the
* POST before it answers, so the reload that follows cannot read the old generation.
*/
async noteWrite() {
localGeneration = Math.max(localGeneration, now());
inflight.clear();
const tableRef = table();
if (!tableRef) return;
const at = localGeneration;
try {
await detach(() => tableRef.put(generationKey, { body: null, contentType: null, startedAt: at, bornAt: at }));
} catch (e) {
log?.warn?.(
`[prerender-console] write generation not stored; other workers may serve pre-write reads for up to their TTL: ${e?.message ?? String(e)}`
);
}
},
};
}