mirror of
https://github.com/mksglu/context-mode.git
synced 2026-10-02 04:14:38 +08:00
This commit is contained in:
@@ -0,0 +1,111 @@
|
||||
/**
|
||||
* ctx_search flood-guard — per-agent-context progressive throttle.
|
||||
*
|
||||
* Background (#79 / #155 / #697): ctx_search carries a progressive throttle
|
||||
* so a single actor cannot spam dozens of individual searches and flood the
|
||||
* context window instead of batching via ctx_batch_execute. The original
|
||||
* implementation kept ONE module-global counter on the MCP server process.
|
||||
*
|
||||
* Issue #769: a parallel multi-agent fan-out (Claude Code Task/Workflow)
|
||||
* runs N subagents concurrently against the SAME per-session MCP server
|
||||
* process. With a single global counter their independent calls are summed
|
||||
* into one budget, so legitimate fan-out ("10 agents x 2 calls") trips the
|
||||
* guard that was only ever meant to catch ONE actor spamming. The budget is
|
||||
* tool-availability state that is logically per-agent-context, so the counter
|
||||
* must be keyed per agent-context — NOT removed. Single-actor flood
|
||||
* protection is preserved exactly; only the bucketing changes.
|
||||
*
|
||||
* This module is pure and transport-free so the policy is unit-testable
|
||||
* without spinning up the MCP server. `src/server.ts` owns the singleton and
|
||||
* supplies the per-call agent key (the session/agent id from
|
||||
* currentAttribution()).
|
||||
*/
|
||||
|
||||
export interface FloodGuardConfig {
|
||||
/** Rolling window length in ms. After this elapses a key's counter resets. */
|
||||
windowMs: number;
|
||||
/** After this many calls in the window, results taper to 1 per query. */
|
||||
softCapAfter: number;
|
||||
/** After this many calls in the window, the call is hard-blocked. */
|
||||
blockAfter: number;
|
||||
}
|
||||
|
||||
export interface FloodDecision {
|
||||
/** This key's call count within the current rolling window (1-based). */
|
||||
count: number;
|
||||
/** Window start timestamp (ms) for this key — used for the "in Ns" message. */
|
||||
windowStart: number;
|
||||
/** True once count exceeds blockAfter — caller must refuse the search. */
|
||||
blocked: boolean;
|
||||
/** True once count exceeds softCapAfter — caller trims to 1 result/query. */
|
||||
softCapped: boolean;
|
||||
}
|
||||
|
||||
interface Bucket {
|
||||
count: number;
|
||||
windowStart: number;
|
||||
}
|
||||
|
||||
/**
|
||||
* A rolling-window call counter bucketed per agent-context key. Each key gets
|
||||
* an independent window + counter, so concurrent subagents do not consume one
|
||||
* another's budget while a single greedy actor is still throttled and blocked
|
||||
* exactly as before.
|
||||
*/
|
||||
export class FloodGuard {
|
||||
readonly #cfg: FloodGuardConfig;
|
||||
readonly #buckets = new Map<string, Bucket>();
|
||||
/**
|
||||
* Hard ceiling on tracked keys — a defensive bound so a pathological host
|
||||
* that mints unbounded distinct agent ids cannot grow the map without limit.
|
||||
* When exceeded, the oldest-window bucket is evicted (its actor simply gets
|
||||
* a fresh window on its next call — fail-open, never a false block).
|
||||
*/
|
||||
readonly #maxKeys: number;
|
||||
|
||||
constructor(cfg: FloodGuardConfig, maxKeys = 4096) {
|
||||
this.#cfg = cfg;
|
||||
this.#maxKeys = Math.max(1, maxKeys);
|
||||
}
|
||||
|
||||
/**
|
||||
* Record one ctx_search call for `key` at time `now` (ms) and return the
|
||||
* throttle decision. Pure aside from the internal per-key counter state.
|
||||
*/
|
||||
record(key: string, now: number = Date.now()): FloodDecision {
|
||||
let bucket = this.#buckets.get(key);
|
||||
|
||||
if (!bucket || now - bucket.windowStart > this.#cfg.windowMs) {
|
||||
bucket = { count: 0, windowStart: now };
|
||||
this.#buckets.set(key, bucket);
|
||||
this.#evictIfNeeded();
|
||||
}
|
||||
|
||||
bucket.count++;
|
||||
|
||||
return {
|
||||
count: bucket.count,
|
||||
windowStart: bucket.windowStart,
|
||||
blocked: bucket.count > this.#cfg.blockAfter,
|
||||
softCapped: bucket.count > this.#cfg.softCapAfter,
|
||||
};
|
||||
}
|
||||
|
||||
/** Test/diagnostics helper — number of distinct keys currently tracked. */
|
||||
size(): number {
|
||||
return this.#buckets.size;
|
||||
}
|
||||
|
||||
#evictIfNeeded(): void {
|
||||
if (this.#buckets.size <= this.#maxKeys) return;
|
||||
let oldestKey: string | undefined;
|
||||
let oldestStart = Infinity;
|
||||
for (const [k, b] of this.#buckets) {
|
||||
if (b.windowStart < oldestStart) {
|
||||
oldestStart = b.windowStart;
|
||||
oldestKey = k;
|
||||
}
|
||||
}
|
||||
if (oldestKey !== undefined) this.#buckets.delete(oldestKey);
|
||||
}
|
||||
}
|
||||
+35
-13
@@ -58,6 +58,7 @@ import {
|
||||
CTX_SEARCH_SHARED_MODE,
|
||||
resolveProjectScope,
|
||||
} from "./search/ctx-search-schema.js";
|
||||
import { FloodGuard } from "./search/flood-guard.js";
|
||||
import { buildNodeCommand, type HookAdapter, type PlatformId, isInProcessPluginPlatform } from "./adapters/types.js";
|
||||
import { detectPlatform, getSessionDirSegments } from "./adapters/detect.js";
|
||||
import { parseCodexContextModePluginRoot } from "./adapters/codex/index.js";
|
||||
@@ -2350,12 +2351,36 @@ function readPositiveEnv(name: string, defaultValue: number): number {
|
||||
return Number.isFinite(parsed) && parsed > 0 ? parsed : defaultValue;
|
||||
}
|
||||
|
||||
let searchCallCount = 0;
|
||||
let searchWindowStart = Date.now();
|
||||
const SEARCH_WINDOW_MS = readPositiveEnv("CONTEXT_MODE_SEARCH_WINDOW_MS", 60_000);
|
||||
const SEARCH_MAX_RESULTS_AFTER = readPositiveEnv("CONTEXT_MODE_SEARCH_MAX_RESULTS_AFTER", 3); // after N calls: 1 result per query
|
||||
const SEARCH_BLOCK_AFTER = readPositiveEnv("CONTEXT_MODE_SEARCH_BLOCK_AFTER", 8); // after N calls: refuse, demand batching
|
||||
|
||||
// #769: progressive throttle bucketed PER agent-context, not machine-global.
|
||||
// Concurrent subagents share ONE MCP server process; a single global counter
|
||||
// summed their independent searches into one budget and hard-blocked
|
||||
// legitimate parallel fan-out. The guard keys each actor's window separately
|
||||
// so single-actor flood protection is preserved while fan-out is not starved.
|
||||
const searchFloodGuard = new FloodGuard({
|
||||
windowMs: SEARCH_WINDOW_MS,
|
||||
softCapAfter: SEARCH_MAX_RESULTS_AFTER,
|
||||
blockAfter: SEARCH_BLOCK_AFTER,
|
||||
});
|
||||
|
||||
/**
|
||||
* Per-agent flood-guard key. Each concurrent subagent in a Claude Code
|
||||
* Task/Workflow fan-out runs under its own session id (written to SessionDB
|
||||
* via hooks), so currentAttribution().sessionId is the per-agent discriminator
|
||||
* already available MCP-side. Falls back to a single shared bucket when no
|
||||
* identity is resolvable (preserves today's single-threaded behaviour).
|
||||
*/
|
||||
function searchFloodGuardKey(): string {
|
||||
try {
|
||||
return currentAttribution()?.sessionId ?? "__default__";
|
||||
} catch {
|
||||
return "__default__";
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Defensive coercion: parse stringified JSON arrays, AND lift a bare
|
||||
* non-empty string into a single-element array.
|
||||
@@ -2511,20 +2536,17 @@ EXAMPLE: ctx_search(queries: ["last user prompt", "active skills", "open blocker
|
||||
() => getProjectDir(),
|
||||
);
|
||||
|
||||
// Progressive throttling: track calls in time window
|
||||
// Progressive throttling: track calls per agent-context window (#769).
|
||||
const now = Date.now();
|
||||
if (now - searchWindowStart > SEARCH_WINDOW_MS) {
|
||||
searchCallCount = 0;
|
||||
searchWindowStart = now;
|
||||
}
|
||||
searchCallCount++;
|
||||
const flood = searchFloodGuard.record(searchFloodGuardKey(), now);
|
||||
const searchCallCount = flood.count;
|
||||
|
||||
// After SEARCH_BLOCK_AFTER calls: refuse
|
||||
if (searchCallCount > SEARCH_BLOCK_AFTER) {
|
||||
// After SEARCH_BLOCK_AFTER calls (for THIS agent): refuse
|
||||
if (flood.blocked) {
|
||||
return trackResponse("ctx_search", {
|
||||
content: [{
|
||||
type: "text" as const,
|
||||
text: `BLOCKED: ${searchCallCount} search calls in ${Math.round((now - searchWindowStart) / 1000)}s. ` +
|
||||
text: `BLOCKED: ${searchCallCount} search calls in ${Math.round((now - flood.windowStart) / 1000)}s. ` +
|
||||
"You're flooding context. STOP making individual search calls. " +
|
||||
"Use ctx_batch_execute(commands, queries) for your next research step.",
|
||||
}],
|
||||
@@ -2533,8 +2555,8 @@ EXAMPLE: ctx_search(queries: ["last user prompt", "active skills", "open blocker
|
||||
}
|
||||
|
||||
// Determine per-query result limit based on throttle level
|
||||
const effectiveLimit = searchCallCount > SEARCH_MAX_RESULTS_AFTER
|
||||
? 1 // after 3 calls: only 1 result per query
|
||||
const effectiveLimit = flood.softCapped
|
||||
? 1 // after soft cap: only 1 result per query
|
||||
: Math.min(limit, 2); // normal: max 2
|
||||
|
||||
const MAX_TOTAL = 40 * 1024; // 40KB total cap
|
||||
|
||||
+23
-3
@@ -1217,9 +1217,14 @@ export class SessionDB extends SQLiteBase {
|
||||
.slice(0, 16)
|
||||
.toUpperCase();
|
||||
const attribution = attributions?.[i];
|
||||
const projectDir = String(
|
||||
// #827: store project_dir in canonical path shape so the search-time
|
||||
// allow-set lookup (getSessionIdsForProject) matches regardless of the
|
||||
// separator / trailing-slash form the host adapter happened to emit.
|
||||
// normalizeWorktreePath is the same rule used for project-hash stability.
|
||||
const rawProjectDir = String(
|
||||
attribution?.projectDir ?? event.project_dir ?? this._getSessionProjectDir(sessionId) ?? "",
|
||||
).trim();
|
||||
const projectDir = rawProjectDir === "" ? "" : normalizeWorktreePath(rawProjectDir);
|
||||
const attributionSource = String(
|
||||
attribution?.source ?? event.attribution_source ?? "unknown",
|
||||
);
|
||||
@@ -1407,13 +1412,28 @@ export class SessionDB extends SQLiteBase {
|
||||
*/
|
||||
getSessionIdsForProject(projectDir: string): string[] {
|
||||
try {
|
||||
// #827: match by canonical path shape, not raw bytes. The host adapter
|
||||
// may store `project_dir` in a different separator / trailing-slash
|
||||
// shape than the search path resolves the scope in — most visibly on
|
||||
// Windows, where attribution often carries `C:\Users\me\proj` while the
|
||||
// server resolves `C:/Users/me/proj`. An exact `project_dir = ?` match
|
||||
// then returned an EMPTY allow-set and ctx_search reported "No results
|
||||
// found" even though the content was present. We fold BOTH sides through
|
||||
// the same canonical rule used for project-hash stability
|
||||
// (normalizeWorktreePath): backslash → forward slash, then strip the
|
||||
// trailing slash. Normalizing in SQL (RTRIM(REPLACE(...))) covers rows
|
||||
// already written un-normalized without a migration, while the JS-side
|
||||
// normalize keeps the bound parameter in the identical shape. This
|
||||
// preserves the #737 project scope — distinct directories still differ
|
||||
// after normalization, so cross-project isolation is intact.
|
||||
const normalized = normalizeWorktreePath(projectDir);
|
||||
const rows = this.db
|
||||
.prepare(
|
||||
`SELECT DISTINCT session_id
|
||||
FROM session_events
|
||||
WHERE project_dir = ?`,
|
||||
WHERE RTRIM(REPLACE(project_dir, '\\', '/'), '/') = ?`,
|
||||
)
|
||||
.all(projectDir) as Array<{ session_id: string }>;
|
||||
.all(normalized) as Array<{ session_id: string }>;
|
||||
return rows.map((r) => r.session_id);
|
||||
} catch {
|
||||
return [];
|
||||
|
||||
@@ -0,0 +1,106 @@
|
||||
/**
|
||||
* Issue #769: ctx_search flood-guard counter is shared across concurrent
|
||||
* subagents.
|
||||
*
|
||||
* The progressive throttle (introduced in 103b41dd, hardened in #697/#698)
|
||||
* keeps a single module-global counter on the per-session MCP server
|
||||
* process. In a parallel multi-agent fan-out (Claude Code Task/Workflow),
|
||||
* N subagents share that ONE process, so their independent search calls are
|
||||
* summed into a single budget. The guard — designed to stop ONE actor
|
||||
* spamming individual searches (#79/#155) — then misclassifies legitimate
|
||||
* parallel fan-out as flooding and hard-blocks subagents collectively.
|
||||
*
|
||||
* Fix: bucket the counter per agent-context key (the per-call session/agent
|
||||
* id) so each actor gets its own rolling window, WITHOUT removing the
|
||||
* single-actor flood protection.
|
||||
*
|
||||
* These tests pin the extracted pure flood-guard (src/search/flood-guard.ts)
|
||||
* so the policy is testable without spinning up the MCP transport.
|
||||
*/
|
||||
|
||||
import { describe, test, expect } from "vitest";
|
||||
import { FloodGuard } from "../../src/search/flood-guard.js";
|
||||
|
||||
const CFG = { windowMs: 60_000, softCapAfter: 3, blockAfter: 8 };
|
||||
|
||||
describe("Issue #769: flood-guard is per-agent-context, not machine-global", () => {
|
||||
test("single actor: blocks after blockAfter calls in the window (protection preserved)", () => {
|
||||
const guard = new FloodGuard(CFG);
|
||||
const now = 1_000;
|
||||
const agent = "agent-solo";
|
||||
|
||||
// 8 calls allowed (count 1..8 <= blockAfter), 9th blocked.
|
||||
for (let i = 1; i <= CFG.blockAfter; i++) {
|
||||
const d = guard.record(agent, now + i);
|
||||
expect(d.blocked).toBe(false);
|
||||
expect(d.count).toBe(i);
|
||||
}
|
||||
const ninth = guard.record(agent, now + 9);
|
||||
expect(ninth.blocked).toBe(true);
|
||||
expect(ninth.count).toBe(CFG.blockAfter + 1);
|
||||
});
|
||||
|
||||
test("RED→GREEN: concurrent subagents each get their OWN budget — fan-out is not collectively starved", () => {
|
||||
const guard = new FloodGuard(CFG);
|
||||
const now = 1_000;
|
||||
|
||||
// 10 distinct subagents, each makes 2 search calls inside the SAME 16s
|
||||
// window (the issue's exact "10 agents x 2 calls" scenario). Aggregate
|
||||
// is 20 calls — well over the global budget of 8 — but NO agent should
|
||||
// be blocked, because each has its own rolling counter.
|
||||
let anyBlocked = false;
|
||||
for (let call = 1; call <= 2; call++) {
|
||||
for (let a = 0; a < 10; a++) {
|
||||
const d = guard.record(`subagent-${a}`, now + call * 1000);
|
||||
if (d.blocked) anyBlocked = true;
|
||||
// Per-agent count must reflect only that agent's own calls.
|
||||
expect(d.count).toBe(call);
|
||||
}
|
||||
}
|
||||
expect(anyBlocked).toBe(false);
|
||||
});
|
||||
|
||||
test("one greedy actor does not consume another actor's budget", () => {
|
||||
const guard = new FloodGuard(CFG);
|
||||
const now = 1_000;
|
||||
|
||||
// Greedy agent floods to its hard block.
|
||||
for (let i = 1; i <= CFG.blockAfter + 1; i++) {
|
||||
guard.record("greedy", now + i);
|
||||
}
|
||||
const greedy = guard.record("greedy", now + CFG.blockAfter + 2);
|
||||
expect(greedy.blocked).toBe(true);
|
||||
|
||||
// A different agent, first call, must be untouched and unblocked.
|
||||
const fresh = guard.record("innocent", now + CFG.blockAfter + 3);
|
||||
expect(fresh.blocked).toBe(false);
|
||||
expect(fresh.count).toBe(1);
|
||||
});
|
||||
|
||||
test("rolling window resets per agent after windowMs elapses", () => {
|
||||
const guard = new FloodGuard(CFG);
|
||||
const agent = "agent-x";
|
||||
|
||||
guard.record(agent, 1_000);
|
||||
guard.record(agent, 2_000);
|
||||
// Jump beyond the window — counter resets for this agent.
|
||||
const afterReset = guard.record(agent, 1_000 + CFG.windowMs + 1);
|
||||
expect(afterReset.count).toBe(1);
|
||||
expect(afterReset.blocked).toBe(false);
|
||||
});
|
||||
|
||||
test("soft cap: effective per-query limit tapers after softCapAfter calls (1 actor)", () => {
|
||||
const guard = new FloodGuard(CFG);
|
||||
const agent = "agent-y";
|
||||
const now = 1_000;
|
||||
|
||||
// calls 1..3 are at/under soft cap → not soft-capped
|
||||
for (let i = 1; i <= CFG.softCapAfter; i++) {
|
||||
const d = guard.record(agent, now + i);
|
||||
expect(d.softCapped).toBe(false);
|
||||
}
|
||||
// call 4 exceeds soft cap → soft-capped (server trims to 1 result/query)
|
||||
const fourth = guard.record(agent, now + CFG.softCapAfter + 1);
|
||||
expect(fourth.softCapped).toBe(true);
|
||||
});
|
||||
});
|
||||
@@ -361,3 +361,75 @@ describe("Slice 5: resolveProjectScope", () => {
|
||||
expect(scope2).toBeUndefined();
|
||||
});
|
||||
});
|
||||
|
||||
// ═══════════════════════════════════════════════════════════
|
||||
// Slice 6: Issue #827 — project_dir matching is path-shape stable
|
||||
//
|
||||
// On Windows the host adapter writes session_events.project_dir in the
|
||||
// shape it observed (often backslash form, e.g. C:\Users\me\proj), while
|
||||
// the search path resolves the scope from the MCP server's getProjectDir()
|
||||
// which can differ by separator (C:/Users/me/proj) or trailing slash. The
|
||||
// original getSessionIdsForProject did an EXACT `project_dir = ?` match, so
|
||||
// the allow-set came back EMPTY and every chunk failed the allow-set test —
|
||||
// ctx_search returned "No results found" even though the content was there
|
||||
// (reproduced by Adriftnote/aaddrr on Windows 11, v1.0.162).
|
||||
//
|
||||
// The fix must NOT disable the #737 project scope (that would re-break the
|
||||
// shared-DB isolation the filter exists for). Instead, BOTH the stored
|
||||
// project_dir and the query project_dir are normalized canonically (the
|
||||
// same normalizeWorktreePath rule already used for project-hash stability,
|
||||
// src/session/db.ts:310) so equal directories compare equal regardless of
|
||||
// separator / trailing-slash shape.
|
||||
// ═══════════════════════════════════════════════════════════
|
||||
|
||||
describe("Slice 6: getSessionIdsForProject path-shape normalization (#827)", () => {
|
||||
test("matches when stored project_dir uses backslashes but query uses forward slashes (Windows shape)", () => {
|
||||
const db = createSessionDB();
|
||||
const sid = `s-win-${randomUUID()}`;
|
||||
|
||||
// Host adapter on Windows stored the backslash form...
|
||||
db.ensureSession(sid, "C:\\Users\\me\\proj");
|
||||
db.insertEvent(
|
||||
sid,
|
||||
{ type: "x", category: "x", data: "win-evt", priority: 2 },
|
||||
"PostToolUse",
|
||||
{ projectDir: "C:\\Users\\me\\proj", source: "env", confidence: 1 },
|
||||
);
|
||||
|
||||
// ...but the search path resolved the scope with forward slashes.
|
||||
const ids = db.getSessionIdsForProject("C:/Users/me/proj");
|
||||
expect(ids).toEqual([sid]);
|
||||
});
|
||||
|
||||
test("matches across trailing-slash difference", () => {
|
||||
const db = createSessionDB();
|
||||
const sid = `s-trail-${randomUUID()}`;
|
||||
|
||||
db.ensureSession(sid, "/home/me/proj");
|
||||
db.insertEvent(
|
||||
sid,
|
||||
{ type: "x", category: "x", data: "trail-evt", priority: 2 },
|
||||
"PostToolUse",
|
||||
{ projectDir: "/home/me/proj/", source: "env", confidence: 1 },
|
||||
);
|
||||
|
||||
const ids = db.getSessionIdsForProject("/home/me/proj");
|
||||
expect(ids).toEqual([sid]);
|
||||
});
|
||||
|
||||
test("still isolates distinct projects (does not over-match after normalization)", () => {
|
||||
const db = createSessionDB();
|
||||
const sA = `s-a-${randomUUID()}`;
|
||||
const sB = `s-b-${randomUUID()}`;
|
||||
|
||||
db.ensureSession(sA, "C:\\work\\alpha");
|
||||
db.ensureSession(sB, "C:\\work\\beta");
|
||||
db.insertEvent(sA, { type: "x", category: "x", data: "_", priority: 2 },
|
||||
"PostToolUse", { projectDir: "C:\\work\\alpha", source: "env", confidence: 1 });
|
||||
db.insertEvent(sB, { type: "x", category: "x", data: "_", priority: 2 },
|
||||
"PostToolUse", { projectDir: "C:\\work\\beta", source: "env", confidence: 1 });
|
||||
|
||||
expect(db.getSessionIdsForProject("C:/work/alpha")).toEqual([sA]);
|
||||
expect(db.getSessionIdsForProject("C:/work/beta")).toEqual([sB]);
|
||||
});
|
||||
});
|
||||
|
||||
Reference in New Issue
Block a user