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
11 changes: 9 additions & 2 deletions AGENT-INSTALL.md
Original file line number Diff line number Diff line change
Expand Up @@ -574,8 +574,15 @@ carries no detections and would otherwise read as one. Those counts stay in your
sent to report them.

`protection.stop()` stops everything the guard has running in the background — the rule-refresh loop, the
block-log reporter, the detection reporter — and flushes what is buffered. `protection.stopRefresh()` is
the same method under its older name. Call it on shutdown; it is safe to call twice.
block-log reporter, the detection reporter — and flushes what is buffered. With `egress: true` it also
removes this guard's outbound-request screening; once no guard in the process is screening, `fetch` and
`node:http`/`node:https` are restored. Call it on shutdown; it is safe to call twice.
`protection.stopRefresh()` stops only the rule refresh: the reporters and outbound-request screening keep
running.

When more than one guard in a process has `egress: true`, an outbound call is checked by each of them and
refused if any one refuses it. A host listed in one guard's `allowHosts` is still refused when another guard
refuses it.

Two more endpoints the package can call, for completeness:

Expand Down
138 changes: 108 additions & 30 deletions src/protect/egress.js
Original file line number Diff line number Diff line change
Expand Up @@ -10,12 +10,50 @@
* onBlock?: (info:{url:string,host:string|null,method:string})=>void,
* dnsScreen?: boolean,
* lookup?: Function }} opts
* @returns {Promise<() => void>} uninstall (restores every patched surface)
* @returns {Promise<() => void>} uninstall (removes this screen; the last one out restores the patched surfaces)
*/
import { notify } from './notify.js';

// fetch and node:http(s) are process-wide, so the guard on them is too: one wrapper per surface, shared
// by every protection in the process (including another copy of this package), each registering its own
// screen. A call is refused when any registered screen refuses it. Keyed on a global symbol so that two
// copies of this module share one registry rather than each deciding the other's wrapper is enough.
const REGISTRY = Symbol.for('patchstack.connect.egress-guard');

function egressRegistry() {
const existing = globalThis[REGISTRY];
if (existing && existing.screens instanceof Set && existing.surfaces instanceof Map) return existing;
const created = { screens: new Set(), surfaces: new Map() };
Object.defineProperty(globalThis, REGISTRY, { value: created, configurable: true, writable: true });
return created;
}

const refusal = (host) => new Error(`Patchstack blocked an outbound request to a disallowed address: ${host}`);

// The destination of a fetch call whose arguments `Request` would not accept, when one can be read: a
// URL string, a URL object, or an object carrying `url`/`href`. Null when there is no parseable URL.
function readableDestination(input, init) {
try {
const raw = typeof input === 'string' || input instanceof URL ? String(input) : input?.url ?? input?.href;
if (typeof raw !== 'string' && !(raw instanceof URL)) return null;
const url = new URL(String(raw)).href;
const method = String(init?.method ?? input?.method ?? 'GET').toUpperCase();
return { url, method };
} catch {
return null;
}
}

// Asks every screen, so each one reports its own refusal, then answers whether any refused.
function anyRefuses(screens, url, host, method) {
let refused = false;
for (const screen of screens) {
if (screen.block(url, host, method)) refused = true;
}
return refused;
}

export async function installEgressGuard({ shouldBlock, onBlock, onSkip, dnsScreen = true, lookup, allowHosts } = {}) {
const restores = [];
if (typeof shouldBlock !== 'function') return () => {};
const exempt = new Set((allowHosts ?? []).map((h) => String(h).toLowerCase()));

Expand Down Expand Up @@ -89,6 +127,10 @@ export async function installEgressGuard({ shouldBlock, onBlock, onSkip, dnsScre
}
});

const registry = egressRegistry();
const own = { block, prescreen: resolvesToDisallowed, dns: screen, skip };
registry.screens.add(own);

// 1. global fetch — synchronous install, so it's active the instant this returns (no startup race).
const originalFetch = globalThis.fetch;
if (typeof originalFetch === 'function' && !originalFetch.__patchstackGuarded) {
Expand Down Expand Up @@ -158,26 +200,33 @@ export async function installEgressGuard({ shouldBlock, onBlock, onSkip, dnsScre
}
};

// Screen one outbound URL: hostname/allowlist/literal-IP check, then a DNS-resolution check for
// real hostnames. Throws if the destination is disallowed.
// Screen one outbound URL against every registered screen: hostname/allowlist/literal-IP check,
// then a DNS-resolution check for real hostnames. Throws if any screen disallows the destination.
const screenUrl = async (u, method) => {
let host = null;
try {
host = new URL(u).hostname;
} catch {
host = null;
}
if (block(u, host, method) || (await resolvesToDisallowed(u, host, method))) {
throw new Error(`Patchstack blocked an outbound request to a disallowed address: ${host ?? u}`);
}
const screens = [...registry.screens];
if (anyRefuses(screens, u, host, method)) throw refusal(host ?? u);
const resolved = await Promise.all(screens.map((each) => each.prescreen(u, host, method)));
if (resolved.includes(true)) throw refusal(host ?? u);
};

const guarded = async (input, init) => {
let cur;
try {
cur = new Request(input, { ...(init || {}), redirect: 'manual' });
} catch {
return originalFetch(input, init); // odd input we can't normalize — fail open, don't break the caller
// An input this runtime's Request refuses is handed to the underlying fetch as it came, which
// decides whether it is a request at all. Its destination is still screened when it can be read;
// when it cannot, the call goes out unscreened and is counted as such.
const destination = readableDestination(input, init);
if (destination) await screenUrl(destination.url, destination.method);
else for (const each of registry.screens) each.skip('unrecognised-request', {});
return originalFetch(input, init);
}
const callerRedirect = (init && init.redirect) || (input && input.redirect) || 'follow';

Expand Down Expand Up @@ -247,8 +296,15 @@ export async function installEgressGuard({ shouldBlock, onBlock, onSkip, dnsScre
};
guarded.__patchstackGuarded = true;
globalThis.fetch = guarded;
restores.push(() => {
if (globalThis.fetch === guarded) globalThis.fetch = originalFetch;
// Released only while it is still the global: a wrapper layered on top later (an APM agent, …)
// keeps calling this one, which then stays registered and screens with whichever screens are
// registered at the time — none, until a protection registers again.
registry.surfaces.set(guarded, {
release() {
if (globalThis.fetch !== guarded) return false;
globalThis.fetch = originalFetch;
return true;
},
});
}

Expand All @@ -270,12 +326,16 @@ export async function installEgressGuard({ shouldBlock, onBlock, onSkip, dnsScre
for (const moduleName of ['node:http', 'node:https']) {
try {
const mod = await import(moduleName);
const restore = patchHttpModule(mod.default ?? mod, block, screen, skip);
if (restore) {
if (registry.surfaces.has(moduleName)) continue;
const release = patchHttpModule(mod.default ?? mod, registry);
if (release) {
patchedAny = true;
restores.push(() => {
restore();
syncBuiltins();
registry.surfaces.set(moduleName, {
release() {
if (!release()) return false;
syncBuiltins();
return true;
},
});
}
} catch {
Expand All @@ -292,10 +352,14 @@ export async function installEgressGuard({ shouldBlock, onBlock, onSkip, dnsScre
// and a hostname-only check would over-promise the control. Outbound screening covers fetch and
// node:http/https.

// Removes this screen only. The last screen to leave releases the surfaces it can; one that is no
// longer the outermost wrapper stays registered, screening nothing until a screen registers again.
return () => {
for (const restore of restores) {
registry.screens.delete(own);
if (registry.screens.size > 0) return;
for (const [name, surface] of registry.surfaces) {
try {
restore();
if (surface.release()) registry.surfaces.delete(name);
} catch {
/* ignore */
}
Expand All @@ -305,7 +369,7 @@ export async function installEgressGuard({ shouldBlock, onBlock, onSkip, dnsScre

// Wrap http(s).request/get — and, on node:http, the ClientRequest constructor they build — so a
// blocked destination throws before the socket opens.
function patchHttpModule(http, block, screen, skip) {
function patchHttpModule(http, registry) {
if (!http || typeof http.request !== 'function' || http.__patchstackGuarded) return null;
const originalRequest = http.request;
const originalGet = http.get;
Expand All @@ -314,14 +378,23 @@ function patchHttpModule(http, block, screen, skip) {
// The arguments to hand on, after screening them. Throws when the destination is refused.
const guardArgs = (args) => {
const target = extractHttpTarget(args);
if (target && block(target.url, target.host, target.method)) {
throw new Error(`Patchstack blocked an outbound request to a disallowed address: ${target.host ?? target.url}`);
}
// DNS screen: only for real hostnames (a literal IP was already covered by the check above),
// and skip an explicitly allowlisted host (the operator trusts it — don't second-guess its DNS).
if (target && screen && target.host && screen.isIP(target.host) === 0 && !screen.isExempt(target.host)) {
if (!target) return args;
const screens = [...registry.screens];
if (anyRefuses(screens, target.url, target.host, target.method)) throw refusal(target.host ?? target.url);
// DNS screen: only for real hostnames (a literal IP was already covered by the check above), and
// not for a screen that allowlists this host (the operator trusts it — don't second-guess its DNS).
// One resolution serves every screen that wants one, so the connection is pinned to addresses that
// all of them checked; it goes through the resolver of the earliest of those screens.
const resolving = target.host
? screens.filter((each) => each.dns && each.dns.isIP(target.host) === 0 && !each.dns.isExempt(target.host))
: [];
if (resolving.length > 0) {
const block = (url, host, method) => anyRefuses(resolving, url, host, method);
const skip = (reason, detail) => {
for (const each of resolving) each.skip(reason, detail);
};
try {
return withScreeningLookup(args, target, block, screen.lookup, skip);
return withScreeningLookup(args, target, block, resolving[0].dns.lookup, skip);
} catch {
// The call goes on with the arguments it came with. Nothing is counted as a fail-open bypass:
// the only thing here that can throw is reading the caller's options, and Node copies that
Expand Down Expand Up @@ -365,13 +438,18 @@ function patchHttpModule(http, block, screen, skip) {
}
http.__patchstackGuarded = true;

// Released only when every wrapper is still the module's own export — don't clobber a wrapper another
// library (an APM agent, etc.) layered on top of us after install. Otherwise nothing is restored and
// the module keeps calling through ours, so it is never left half-guarded.
return () => {
// Only restore if our wrapper is still installed — don't clobber a wrapper another library
// (an APM agent, etc.) layered on top of us after install.
if (http.request === guardedRequest) http.request = originalRequest;
if (guardedGet && http.get === guardedGet) http.get = originalGet;
if (GuardedClientRequest && http.ClientRequest === GuardedClientRequest) http.ClientRequest = OriginalClientRequest;
if (http.request !== guardedRequest) return false;
if (guardedGet && http.get !== guardedGet) return false;
if (GuardedClientRequest && http.ClientRequest !== GuardedClientRequest) return false;
http.request = originalRequest;
if (guardedGet) http.get = originalGet;
if (GuardedClientRequest) http.ClientRequest = OriginalClientRequest;
delete http.__patchstackGuarded;
return true;
};
}

Expand Down
12 changes: 9 additions & 3 deletions src/protect/protect.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,11 @@ export interface Protection {
screenResponse(response: Response, request?: Request): Promise<Response>;
express(options?: { screenResponses?: boolean }): (req: unknown, res: unknown, next: () => void) => void;
node(options?: { maxBodyBytes?: number; screenResponses?: boolean }): (req: unknown, res: unknown, next: () => void) => void;
/** Present when `egress: true` — restores the original global fetch. */
/** Present when `egress: true` — removes this protection's outbound screen. Outbound calls are
* screened by every protection that has one registered, and any one of them can refuse a call:
* a host in one protection's `allowHosts` is still refused when another protection refuses it.
* When the last one leaves, the original `fetch` and `node:http`/`node:https` functions are
* restored. `stop()` calls this too; `stopRefresh()` does not. */
uninstallEgress?: () => void;
/** Present with a live source — re-fetch + hot-swap the rules once (used by the loop + push).
* Resolves with the outcome of the attempt: `ok: false` means the rules in force came from the
Expand All @@ -46,7 +50,8 @@ export interface Protection {
* the configured refresh secret (a push/zero-day trigger). No secret set → the handler 404s. */
refreshHandler?: () => (request: Request) => Promise<Response>;
/** Stops everything with a timer or a buffer behind it: the refresh loop, the block log, the
* detection reporter (flushing what it holds). Always present, and safe to call twice. */
* detection reporter (flushing what it holds), and this protection's outbound screen. Always
* present, and safe to call twice. */
/**
* Stop everything holding a timer or a buffer.
*
Expand All @@ -67,7 +72,8 @@ export interface Protection {
* to a failed token exchange, a failed post, or a shutdown that ran out of time is reported nowhere.
*/
stop: () => Promise<void>;
/** Alias of `stop`, under the name callers already have. */
/** Stops the rule refresh only — the poll loop and its recovery retries. The reporters and this
* protection's outbound screening keep running; use `stop()` to end those too. */
stopRefresh: () => Promise<void>;
/**
* Where the rules in force came from, and whether the most recent resolution was clean — the same
Expand Down
21 changes: 15 additions & 6 deletions src/protect/runtime.js
Original file line number Diff line number Diff line change
Expand Up @@ -1685,17 +1685,20 @@ export async function createProtection(options = {}) {
const canAsk = Boolean(options.token || pulseAuth);
recovery = live && canAsk && !loop && ruleSource.ok === false ? startRecovery(refreshTick, { onError }) : null;

// One method, always present, that reaches everything holding a timer or a buffer: the refresh loop,
// the block log, the detection reporter. Always present because a lifecycle method that exists only
// for some configurations is one a caller cannot rely on — and each of these components can be the
// only one installed, so any of them can be the one left running.
// One method, always present, that reaches everything holding a timer, a buffer or a process-wide
// hook: the refresh loop, the block log, the detection reporter, the outbound screen. Always present
// because a lifecycle method that exists only for some configurations is one a caller cannot rely on
// — and each of these components can be the only one installed, so any of them can be the one left
// running.
//
// Returns a promise that settles when the reporter has finished draining, so a host shutting down can
// await it rather than racing the last batch against process exit. Bounded and best-effort — a runtime
// that terminates regardless still wins — and ignoring the return behaves exactly as before.
protection.stop = () => {
loop?.stop();
recovery?.stop();
// This protection's outbound screen leaves the shared guard; other protections keep theirs.
protection.uninstallEgress?.();
// Both reporters, because the promise says every buffer this reaches is finished with. Waiting only
// for one would resolve while the other still had records outstanding — and resolve immediately in a
// configuration where the one being waited for was never built.
Expand All @@ -1705,8 +1708,14 @@ export async function createProtection(options = {}) {

return Promise.all(outstanding).then(() => undefined);
};
// The name callers already have, kept as an alias for it.
protection.stopRefresh = protection.stop;
// The rule refresh only: the poll loop and the recovery retries. The reporters and this protection's
// outbound screen keep running; `stop()` ends those as well.
protection.stopRefresh = () => {
loop?.stop();
recovery?.stop();

return Promise.resolve();
};
// Which of the three states reporting is in: requested and running, requested but undeliverable, or
// not requested. A boolean would collapse the middle one into "off", which is the reassuring reading.
// A getter, because the state follows refreshes: a property assigned once would report the boot value
Expand Down
2 changes: 1 addition & 1 deletion tests/protect/callback-containment.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -114,7 +114,7 @@ describe('a throwing host callback cannot break the guard', () => {
// poll loop is how a long-lived process dies hours after the mistake was made.
await expect(p.refresh()).resolves.not.toThrow();

p.stopRefresh?.();
p.stop?.();
});

it('keeps serving when onSkip throws', async () => {
Expand Down
10 changes: 5 additions & 5 deletions tests/protect/detections.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -322,7 +322,7 @@ describe('wiring', () => {
const posted = fetchMock.mock.calls.filter(([url]) => String(url).includes('/detections/'));
expect(posted.length, 'no detection report without reportDetections: true').toBe(0);

p.stopRefresh?.();
p.stop?.();
});
});

Expand Down Expand Up @@ -369,7 +369,7 @@ describe('declaring the capability', () => {
// The first fetch of a site with no cached bundle honestly reports that it holds no managed rules
// yet; the state that follows the resolution is asserted separately below.
expect(p.detectionReporting).toBe('on');
p.stopRefresh?.();
p.stop?.();
});

it('says nothing when reporting is off', async () => {
Expand All @@ -388,7 +388,7 @@ describe('declaring the capability', () => {
const p: any = await createProtection({ siteUuid: 'site-1', pulseRulesUrl: 'https://x.test/monitor/pulse' });

expect(seen.every((h) => h['X-Patchstack-Detections'] === undefined)).toBe(true);
p.stopRefresh?.();
p.stop?.();
});
});

Expand Down Expand Up @@ -427,7 +427,7 @@ describe('the wiring actually runs', () => {
});

await p.fetchGuard()(new Request('https://app.test/api/x?q=boom'));
p.stopRefresh?.();
p.stop?.();
await new Promise((resolve) => setTimeout(resolve, 5));

expect(posted.some((url) => url.includes('/detections/site-1'))).toBe(true);
Expand Down Expand Up @@ -464,7 +464,7 @@ describe('the capability claim is only made when it carries weight', () => {
.toBeUndefined();
}

p.stopRefresh?.();
p.stop?.();
});
});

Expand Down
Loading
Loading