Monitor websites uptime using Cloudflare Workers
Something went wrong. Try again.
14 kB · 402 lines
TypeScript
at optimization
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403import { probeBatchRequestSchema, probeRequestSchema, type ObservationErrorCode, type ProbeBatchRequest, type ProbeRequest, type ProbeResponse,} from '@uptime/contracts';import { isRegionId, type RegionId } from '@uptime/regions';import { isForbiddenIpv6Literal, isIntInForbiddenIpv4Cidrs, ipv4StringToInt,} from '@uptime/contracts';
import { collectEndpointEvidence } from './endpoint-evidence.js';import { collectDnsCandidates } from './dns-candidates.js';
interface Env { PROBE_REGION: RegionId; PROBE_SIGNING_SECRET: string; PROBE_REQUEST_MAX_SKEW_SECONDS: string; PROBE_MAX_REQUEST_BYTES: string; PROBE_VERSION: string;}
const maxRedirects = 5;const maxBodyBytes = 65_536;const signatureVersion = 'v1';const batchProbeConcurrency = 2;
function json(body: unknown, status = 200): Response { return new Response(JSON.stringify(body), { status, headers: { 'content-type': 'application/json; charset=utf-8', 'cache-control': 'no-store' }, });}
function errorResponse(code: string, status: number): Response { return json({ error: { code } }, status);}
function isForbiddenIpv4(host: string): boolean { const parts = host.split('.'); if (parts.length !== 4 || parts.some((part) => !/^\d+$/.test(part))) return false; const parsed = ipv4StringToInt(host); if (parsed === null) return true; return isIntInForbiddenIpv4Cidrs(parsed);}
function isForbiddenIpv6(host: string): boolean { const normalized = host.toLowerCase().replace(/^\[|\]$/g, ''); if (!normalized.includes(':')) return false; return isForbiddenIpv6Literal(normalized);}
function validateTarget(raw: string): URL { let target: URL; try { target = new URL(raw); } catch { throw new Error('invalid_url'); } if ( (target.protocol !== 'http:' && target.protocol !== 'https:') || target.username || target.password || (target.port && !/^\d+$/.test(target.port)) ) { throw new Error('invalid_url'); } if (isForbiddenIpv4(target.hostname) || isForbiddenIpv6(target.hostname)) throw new Error('blocked_address'); return target;}
async function expectedSignature( secret: string, issuedAt: string, requestId: string, body: string,): Promise<string> { const encoder = new TextEncoder(); const key = await crypto.subtle.importKey( 'raw', encoder.encode(secret), { name: 'HMAC', hash: 'SHA-256' }, false, ['sign'], ); const data = encoder.encode(`${signatureVersion}\n${issuedAt}\n${requestId}\n${body}`); const signature = await crypto.subtle.sign('HMAC', key, data); const bytes = new Uint8Array(signature); return btoa(String.fromCharCode(...bytes)) .replace(/\+/g, '-') .replace(/\//g, '_') .replace(/=+$/g, '');}
function constantTimeEqual(left: string, right: string): boolean { if (left.length !== right.length) return false; let diff = 0; for (let i = 0; i < left.length; i += 1) diff |= left.charCodeAt(i) ^ right.charCodeAt(i); return diff === 0;}
async function consumeBounded( response: Response, limit: number,): Promise<{ bodyBytes: number; exceeded: boolean }> { const declaredLength = Number(response.headers.get('content-length')); if (Number.isFinite(declaredLength) && declaredLength > limit) { await response.body?.cancel().catch(() => undefined); return { bodyBytes: 0, exceeded: true }; } if (!response.body) return { bodyBytes: 0, exceeded: false }; const reader = response.body.getReader(); let total = 0; let exceeded = false; try { while (total <= limit) { const { done, value } = await reader.read(); if (done) break; total += value.byteLength; if (total > limit) { exceeded = true; break; } } } finally { await reader.cancel().catch(() => undefined); } return { bodyBytes: Math.min(total, limit), exceeded };}
function result(request: ProbeRequest, fields: Omit<ProbeResponse, 'regionId'>): ProbeResponse { return { regionId: request.regionId, ...fields };}
function failure( request: ProbeRequest, errorCode: ObservationErrorCode, detail: string, startedAt: string, env: Env, metadata: Partial< Pick< ProbeResponse, | 'responseMs' | 'totalMs' | 'placement' | 'colo' | 'finalUrl' | 'endpointEvidence' | 'dnsDiagnostic' | 'redirectCount' | 'bodyBytes' > > = {},): ProbeResponse { return result(request, { status: 'network_failure', success: false, httpStatus: null, responseMs: metadata.responseMs ?? null, totalMs: metadata.totalMs ?? null, errorCode, errorDetail: detail.slice(0, 500), placement: metadata.placement ?? null, colo: metadata.colo ?? null, finalUrl: metadata.finalUrl ?? null, endpointEvidence: metadata.endpointEvidence ?? null, dnsDiagnostic: metadata.dnsDiagnostic ?? null, redirectCount: metadata.redirectCount ?? null, bodyBytes: metadata.bodyBytes ?? null, probeVersion: env.PROBE_VERSION, startedAt, completedAt: new Date().toISOString(), });}
async function diagnosticFor(request: ProbeRequest, target: URL) { if (!request.dnsDiagnostic) return null; return collectDnsCandidates(target.hostname, request.dnsDiagnostic);}
async function mapWithConcurrency<T, Result>( values: readonly T[], concurrency: number, operation: (value: T) => Promise<Result>,): Promise<Result[]> { const results = new Array<Result>(values.length); let nextIndex = 0; async function worker(): Promise<void> { while (nextIndex < values.length) { const index = nextIndex; nextIndex += 1; results[index] = await operation(values[index]!); } } await Promise.all(Array.from({ length: Math.min(concurrency, values.length) }, () => worker())); return results;}
async function probeBatch( request: ProbeBatchRequest, env: Env, runtime: { readonly colo?: string | undefined; readonly placement?: string | undefined },) { const results = await mapWithConcurrency(request.items, batchProbeConcurrency, async (item) => { const response = await probe( { ...item, requestId: request.requestId, issuedAt: request.issuedAt, regionId: request.regionId, }, env, runtime, ); return { checkRunId: item.checkRunId, monitorId: item.monitorId, response }; }); return { requestId: request.requestId, regionId: request.regionId, results };}
async function probe( request: ProbeRequest, env: Env, runtime: { readonly colo?: string | undefined; readonly placement?: string | undefined },): Promise<ProbeResponse> { const started = Date.now(); const startedAt = new Date().toISOString(); let target: URL; try { target = validateTarget(request.url); } catch (error) { return failure(request, 'invalid_response', String(error), startedAt, env); } const controller = new AbortController(); const timeout = setTimeout(() => controller.abort(), request.timeoutMs); let redirects = 0; let responseMs: number | null = null; try { while (true) { const response = await fetch(target.toString(), { method: 'GET', redirect: 'manual', signal: controller.signal, headers: { 'cache-control': 'no-cache', pragma: 'no-cache' }, cf: { cacheTtl: 0 }, }); responseMs = Date.now() - started; if (response.status >= 300 && response.status < 400 && response.headers.get('location')) { if (redirects >= maxRedirects) return failure(request, 'redirect_limit', 'Maximum redirects exceeded', startedAt, env, { responseMs, totalMs: Date.now() - started, placement: runtime.placement ?? null, colo: runtime.colo ?? null, finalUrl: target.toString(), endpointEvidence: collectEndpointEvidence(target.toString(), new Headers()), dnsDiagnostic: await diagnosticFor(request, target), redirectCount: redirects, }); redirects += 1; target = validateTarget(new URL(response.headers.get('location')!, target).toString()); continue; } const evidenceStarted = Date.now(); const endpointEvidence = collectEndpointEvidence(target.toString(), response.headers); const evidenceElapsedMs = Date.now() - evidenceStarted; const consumed = await consumeBounded(response, maxBodyBytes); if (consumed.exceeded) { const totalMs = Math.max(0, Date.now() - started - evidenceElapsedMs); const dnsDiagnostic = await diagnosticFor(request, target); return failure( request, 'response_too_large', 'Response body exceeded 65536 bytes', startedAt, env, { responseMs, totalMs, placement: runtime.placement ?? null, colo: runtime.colo ?? null, finalUrl: target.toString(), endpointEvidence, dnsDiagnostic, redirectCount: redirects, bodyBytes: consumed.bodyBytes, }, ); } const success = response.status >= 200 && response.status <= 399; const totalMs = Math.max(0, Date.now() - started - evidenceElapsedMs); const dnsDiagnostic = await diagnosticFor(request, target); return result(request, { status: success ? 'success' : 'http_failure', success, httpStatus: response.status, responseMs, totalMs, errorCode: null, errorDetail: success ? null : `HTTP ${response.status}`, placement: runtime.placement ?? null, colo: runtime.colo ?? null, finalUrl: target.toString(), endpointEvidence, dnsDiagnostic, redirectCount: redirects, bodyBytes: consumed.bodyBytes, probeVersion: env.PROBE_VERSION, startedAt, completedAt: new Date().toISOString(), }); } } catch (error) { const message = error instanceof Error ? error.message : String(error); const code: ObservationErrorCode = controller.signal.aborted ? 'timeout' : /certificate|tls/i.test(message) ? 'tls' : 'connection'; const totalMs = Date.now() - started; const dnsDiagnostic = await diagnosticFor(request, target); return failure(request, code, message, startedAt, env, { responseMs, totalMs, placement: runtime.placement ?? null, colo: runtime.colo ?? null, finalUrl: target.toString(), endpointEvidence: collectEndpointEvidence(target.toString(), new Headers()), dnsDiagnostic, redirectCount: redirects, }); } finally { clearTimeout(timeout); }}
export default { async fetch(request: Request, env: Env, _context: ExecutionContext): Promise<Response> { if (request.method !== 'POST') return errorResponse('method_not_allowed', 405); const length = Number(request.headers.get('content-length') ?? 0); const limit = Number(env.PROBE_MAX_REQUEST_BYTES || maxBodyBytes); if (!Number.isFinite(length) || length > limit) return errorResponse('payload_too_large', 413); const body = await request.text(); if (new TextEncoder().encode(body).byteLength > limit) return errorResponse('payload_too_large', 413); const issuedAt = request.headers.get('x-uptime-issued-at'); const requestId = request.headers.get('x-uptime-request-id'); const signature = request.headers.get('x-uptime-signature'); if ( request.headers.get('x-uptime-signature-version') !== signatureVersion || !issuedAt || !requestId || !signature ) return errorResponse('invalid_signature', 401); const issuedAtMs = Date.parse(issuedAt); const skewMs = Number(env.PROBE_REQUEST_MAX_SKEW_SECONDS) * 1_000; if (!Number.isFinite(issuedAtMs) || Math.abs(Date.now() - issuedAtMs) > skewMs) return errorResponse('stale_request', 401); const expected = await expectedSignature(env.PROBE_SIGNING_SECRET, issuedAt, requestId, body); if (!constantTimeEqual(expected, signature)) return errorResponse('invalid_signature', 401); let decoded: unknown; try { decoded = JSON.parse(body); } catch { return errorResponse('invalid_payload', 400); } const runtime = { colo: (request.cf as unknown as { readonly colo?: string } | undefined)?.colo, placement: request.headers.get('cf-placement') ?? undefined, }; const batch = probeBatchRequestSchema.safeParse(decoded); if (batch.success) { if ( !isRegionId(env.PROBE_REGION) || batch.data.requestId !== requestId || batch.data.issuedAt !== issuedAt || batch.data.regionId !== env.PROBE_REGION ) return errorResponse('region_or_identity_mismatch', 403); return json(await probeBatch(batch.data, env, runtime)); } const single = probeRequestSchema.safeParse(decoded); if (!single.success) return errorResponse('invalid_payload', 400); const payload = single.data; if ( !isRegionId(env.PROBE_REGION) || payload.requestId !== requestId || payload.issuedAt !== issuedAt || payload.regionId !== env.PROBE_REGION ) return errorResponse('region_or_identity_mismatch', 403); return json(await probe(payload, env, runtime)); },} satisfies ExportedHandler<Env>;
// Workers cannot reveal the final socket address, so trusted public targets remain required.