mirror of
https://github.com/mksglu/context-mode.git
synced 2026-10-03 04:38:25 +08:00
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:
co-authored by
Claude Fable 5
parent
589d8214d5
commit
c94e8fc440
@@ -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`);
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user