fix(forward): suppress platform forwards for 24h after HTTP 402

A churned org's bridge kept firing a POST on every single event forever —
the 402 ("Subscription required") response fell through the status handling
(only 401 and 429 were handled) and the fire-and-forget caller never backed
off.

Now a 402 persists `suppressed_until` (epoch ms, now + 24h) into
platform.json — the same config file the bridge already owns — and prints
one concise stderr note. maybeForward gates on the marker BEFORE any network
call; an expired marker is cleared and forwarding resumes automatically, and
any 2xx response clears a lingering marker for immediate resume after
reactivation. Marker survives restarts and is shared across concurrent
sessions via the config file; in-flight 402 bursts dedupe to a single write
and note.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Mert Koseoglu
2026-07-19 17:36:55 +03:00
co-authored by Claude Fable 5
parent 589d8214d5
commit c94e8fc440
2 changed files with 192 additions and 3 deletions
+59 -2
View File
@@ -11,6 +11,7 @@ import { execSync } from "node:child_process";
const CACHE_TTL_MS = 60_000;
const FETCH_TIMEOUT_MS = 2_000;
const SUPPRESS_TTL_MS = 24 * 60 * 60 * 1000; // 402 back-off window
const MAX_FIELD_LEN = 200;
const MAX_DEPTH = 4;
@@ -54,7 +55,35 @@ function normalizeConfig(raw) {
if (!platform_url && raw.events_url) platform_url = String(raw.events_url).replace(/\/events$/, "");
if (typeof api_key !== "string" || !api_key.startsWith("ctxm_")) return null;
if (typeof platform_url !== "string" || !platform_url) return null;
return { api_key, platform_url: platform_url.replace(/\/$/, "") };
const cfg = { api_key, platform_url: platform_url.replace(/\/$/, "") };
// 402 suppress marker (churned-org back-off) rides the same file so it
// survives restarts and is shared across concurrent sessions.
if (typeof raw.suppressed_until === "number" && Number.isFinite(raw.suppressed_until)) {
cfg.suppressed_until = raw.suppressed_until;
}
return cfg;
}
// === 402 suppress marker (churned-org back-off) ===
// A 402 ("Subscription required") means the org churned — hammering the
// platform on every event forever is pure waste. Persist `suppressed_until`
// (epoch ms) INTO platform.json itself: same file, same read/write ownership,
// no sibling cache to invent. null → clear the marker.
function persistSuppressMarker(untilMs) {
// In-memory first — suppression must hold within this process even if the
// file write fails (read-only FS, concurrent uninstall).
if (_cache && _cache !== NO_CONFIG) {
if (untilMs != null) _cache.suppressed_until = untilMs;
else delete _cache.suppressed_until;
}
const cfgPath = configPath();
try {
const raw = JSON.parse(fs.readFileSync(cfgPath, "utf8"));
if (untilMs != null) raw.suppressed_until = untilMs;
else if (raw.suppressed_until === undefined) return; // nothing to clear — skip the write
else delete raw.suppressed_until;
fs.writeFileSync(cfgPath, JSON.stringify(raw, null, 2) + "\n");
} catch { /* unreadable/unwritable — in-memory state already updated */ }
}
function readConfig() {
@@ -269,6 +298,15 @@ export async function maybeForward(event, platform, opts = {}) {
const cfg = readConfig();
if (!cfg) return;
// 402 suppression gate — BEFORE any allocation or network call. Active
// marker → the org's subscription is inactive; stay silent for 24h.
// Expired marker → clear it and proceed, so reactivated orgs resume
// automatically within 24h even without a fresh login.
if (typeof cfg.suppressed_until === "number") {
if (cfg.suppressed_until > Date.now()) return;
persistSuppressMarker(null);
}
// Project identity must be resolved from the RAW projectDir — the resolver
// reads `git config` against the actual filesystem path. After sanitize,
// $HOME-normalization would break the lookup. We overlay the resolved id
@@ -306,7 +344,26 @@ export async function maybeForward(event, platform, opts = {}) {
}),
signal: ctrl.signal,
});
if (res.status === 401) { _cache = null; _cacheLoadedAt = 0; }
if (res.ok) {
// Reactivation: a successful forward clears any lingering suppress
// marker (e.g. written by a concurrent session) — immediate resume.
if (_cache !== null && _cache !== NO_CONFIG && _cache.suppressed_until !== undefined) {
persistSuppressMarker(null);
}
}
else if (res.status === 401) { _cache = null; _cacheLoadedAt = 0; }
else if (res.status === 402) {
// Subscription inactive (churned org) — pause ALL forwards for 24h.
// Dedupe: a burst of in-flight events all passed the gate before the
// first 402 landed; only the first response writes + notes.
const alreadySuppressed = _cache !== null && _cache !== NO_CONFIG
&& typeof _cache.suppressed_until === "number"
&& _cache.suppressed_until > Date.now();
if (!alreadySuppressed) {
persistSuppressMarker(Date.now() + SUPPRESS_TTL_MS);
process.stderr.write("context-mode: platform subscription inactive — forwards paused for 24h\n");
}
}
else if (res.status === 429) {
process.stderr.write(`[context-mode-platform] rate limited (retry after ${res.headers.get("Retry-After")}s)\n`);
}
+133 -1
View File
@@ -11,7 +11,7 @@
*/
import { describe, test, beforeEach, afterEach, expect, vi } from "vitest";
import { mkdtempSync, rmSync, writeFileSync, mkdirSync } from "node:fs";
import { mkdtempSync, rmSync, writeFileSync, mkdirSync, readFileSync } from "node:fs";
import { dirname, join } from "node:path";
import { tmpdir } from "node:os";
import { execSync } from "node:child_process";
@@ -613,3 +613,135 @@ describe("platform-bridge — project identity resolution", () => {
});
});
});
// ─────────────────────────────────────────────────────────
// 402 subscription suppression — churned orgs must not be
// hammered with a forward on every event forever. A 402
// persists `suppressed_until` (epoch ms) into platform.json
// and pauses ALL forwards for 24h; expiry or a 2xx resumes.
// ─────────────────────────────────────────────────────────
describe("platform-bridge — 402 subscription suppression", () => {
const DAY_MS = 24 * 60 * 60 * 1000;
const config = {
api_key: "ctxm_suppress_test",
platform_url: "https://example.test/api/v1",
};
const event = { type: "tool_use", category: "edit", data: "x" };
let fakeHome: string;
let origHome: string | undefined;
let origXdg: string | undefined;
let origAppData: string | undefined;
let fetchSpy: ReturnType<typeof vi.spyOn>;
beforeEach(() => {
fakeHome = mkdtempSync(join(tmpdir(), "ctx-bridge-402-"));
origHome = process.env.HOME;
origXdg = process.env.XDG_CONFIG_HOME;
origAppData = process.env.APPDATA;
process.env.HOME = fakeHome;
delete process.env.XDG_CONFIG_HOME;
process.env.APPDATA = join(fakeHome, "AppData", "Roaming");
fetchSpy = vi.spyOn(globalThis, "fetch").mockResolvedValue(
new Response(null, { status: 200 }),
);
});
afterEach(() => {
if (origHome !== undefined) process.env.HOME = origHome;
else delete process.env.HOME;
if (origXdg !== undefined) process.env.XDG_CONFIG_HOME = origXdg;
else delete process.env.XDG_CONFIG_HOME;
if (origAppData !== undefined) process.env.APPDATA = origAppData;
else delete process.env.APPDATA;
try { rmSync(fakeHome, { recursive: true, force: true }); } catch {}
vi.resetModules();
vi.restoreAllMocks();
});
test("402 writes suppressed_until (~now+24h) and pauses subsequent forwards — no fetch", async () => {
const cfgFile = writePlatformConfig(fakeHome, config);
fetchSpy.mockResolvedValue(new Response("Subscription required", { status: 402 }));
const stderrSpy = vi.spyOn(process.stderr, "write").mockImplementation(() => true);
const { bridge } = await importFresh();
bridge._internal.resetState();
const before = Date.now();
const res = await bridge.maybeForward(event, "claude-code");
expect(res).toEqual({ ok: false, status: 402 });
// Marker persisted in the SAME config file the bridge already owns.
const persisted = JSON.parse(readFileSync(cfgFile, "utf8"));
expect(typeof persisted.suppressed_until).toBe("number");
expect(persisted.suppressed_until).toBeGreaterThanOrEqual(before + DAY_MS);
expect(persisted.suppressed_until).toBeLessThanOrEqual(Date.now() + DAY_MS);
// One concise operator note.
const pauseNotes = stderrSpy.mock.calls.filter((c) =>
String(c[0]).includes("forwards paused for 24h"));
expect(pauseNotes.length).toBe(1);
// Subsequent forward: early return BEFORE any network call.
fetchSpy.mockClear();
const res2 = await bridge.maybeForward(event, "claude-code");
expect(res2).toBeUndefined();
expect(fetchSpy).not.toHaveBeenCalled();
});
test("active marker survives process restart (fresh import) — still no fetch", async () => {
writePlatformConfig(fakeHome, {
...config,
suppressed_until: Date.now() + 60 * 60 * 1000,
});
const { bridge } = await importFresh();
bridge._internal.resetState();
await bridge.maybeForward(event, "claude-code");
expect(fetchSpy).not.toHaveBeenCalled();
});
test("expired marker: forward proceeds, 2xx clears the marker, forwarding resumes", async () => {
const cfgFile = writePlatformConfig(fakeHome, {
...config,
suppressed_until: Date.now() - 1000, // reactivated org, stale pause
});
const { bridge } = await importFresh();
bridge._internal.resetState();
const res = await bridge.maybeForward(event, "claude-code");
expect(res).toEqual({ ok: true, status: 200 });
expect(fetchSpy).toHaveBeenCalledTimes(1);
// Marker cleared from the config file — immediate durable resume.
const persisted = JSON.parse(readFileSync(cfgFile, "utf8"));
expect(persisted).not.toHaveProperty("suppressed_until");
expect(persisted.api_key).toBe(config.api_key); // rest of config untouched
await bridge.maybeForward(event, "claude-code");
expect(fetchSpy).toHaveBeenCalledTimes(2);
});
test("burst of in-flight 402s: marker written once, ONE stderr note total", async () => {
writePlatformConfig(fakeHome, config);
fetchSpy.mockResolvedValue(new Response("Subscription required", { status: 402 }));
const stderrSpy = vi.spyOn(process.stderr, "write").mockImplementation(() => true);
const { bridge } = await importFresh();
bridge._internal.resetState();
// Fire-and-forget loop shape: all events pass the gate before the first
// 402 response lands. Dedupe must collapse the notes to exactly one.
await Promise.all([
bridge.maybeForward(event, "claude-code"),
bridge.maybeForward(event, "claude-code"),
bridge.maybeForward(event, "claude-code"),
]);
const pauseNotes = stderrSpy.mock.calls.filter((c) =>
String(c[0]).includes("forwards paused for 24h"));
expect(pauseNotes.length).toBe(1);
});
});