Skip to content

Commit 51896d7

Browse files
committed
fix(clv2): serialize observer signal-counter to stop dropped increments
observe.sh bumps the SIGUSR1 throttle counter in ${PROJECT_DIR}/.observer-signal-counter with an unlocked read-modify-write. The hook runs on every tool call, so concurrent invocations read the same value, both increment, and lose a write, signaling the observer at unpredictable intervals and defeating the #521 throttle. Serialize the read-modify-write under a lock, and only ever bump the counter while that lock is held: - Prefer flock (a kernel advisory lock the OS auto-releases when the fd closes or the process dies, so there is no stale lock and no lost increment). - Fall back to an atomic mkdir lock on platforms without flock. It releases via a trap on EXIT/INT/TERM (so an async-timeout SIGTERM cannot strand the lock) and uses a bounded spin so the hook never blocks. If the lock cannot be acquired in the budget it skips the tick rather than racing on an unlocked counter; no hand-rolled PID stale-reclaim (which is racy and can delete a live re-acquirer's lock). - Guard the counter read against a corrupt (non-integer) file that would abort the hook under set -e. Add tests/hooks/observe-signal-counter-race.test.js: 20 concurrent observe.sh invocations must not lose increments (exact under flock; at most one dropped on the mkdir fallback), the runner rejects on any hook execution failure or hang, plus content guards for the lock and the corrupt-counter handling. Fixes #2296
1 parent 2bc924f commit 51896d7

2 files changed

Lines changed: 326 additions & 8 deletions

File tree

‎skills/continuous-learning-v2/hooks/observe.sh‎

Lines changed: 63 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -477,21 +477,76 @@ fi
477477
# which caused runaway parallel Claude analysis processes.
478478
SIGNAL_EVERY_N="${ECC_OBSERVER_SIGNAL_EVERY_N:-20}"
479479
SIGNAL_COUNTER_FILE="${PROJECT_DIR}/.observer-signal-counter"
480+
SIGNAL_COUNTER_LOCK="${SIGNAL_COUNTER_FILE}.lock"
480481
ACTIVITY_FILE="${PROJECT_DIR}/.observer-last-activity"
481482

482483
touch "$ACTIVITY_FILE" 2>/dev/null || true
483484

485+
# Serialize the throttle-counter read-modify-write. observe.sh runs on every
486+
# tool call (which can fire every second), so concurrent invocations previously
487+
# raced on this counter: both read the same value, both incremented, and one
488+
# write was lost, signaling the observer at unpredictable intervals (#2296).
489+
# Prefer flock (a kernel advisory lock the OS releases automatically if the hook
490+
# is killed); fall back to the atomic mkdir lock this script already uses for
491+
# the lazy-start path above. Both wrap the same read-modify-write below.
484492
should_signal=0
485-
if [ -f "$SIGNAL_COUNTER_FILE" ]; then
486-
counter=$(cat "$SIGNAL_COUNTER_FILE" 2>/dev/null || echo 0)
487-
counter=$((counter + 1))
488-
if [ "$counter" -ge "$SIGNAL_EVERY_N" ]; then
489-
should_signal=1
490-
counter=0
493+
494+
_ecc_bump_signal_counter() {
495+
if [ -f "$SIGNAL_COUNTER_FILE" ]; then
496+
counter=$(cat "$SIGNAL_COUNTER_FILE" 2>/dev/null || echo 0)
497+
# Guard against a corrupt counter file: a non-integer value would abort the
498+
# hook under `set -e` at the arithmetic below.
499+
case "$counter" in
500+
''|*[!0-9]*) counter=0 ;;
501+
esac
502+
counter=$((counter + 1))
503+
if [ "$counter" -ge "$SIGNAL_EVERY_N" ]; then
504+
should_signal=1
505+
counter=0
506+
fi
507+
echo "$counter" > "$SIGNAL_COUNTER_FILE"
508+
else
509+
echo "1" > "$SIGNAL_COUNTER_FILE"
510+
fi
511+
}
512+
513+
if command -v flock >/dev/null 2>&1 && exec 8>"$SIGNAL_COUNTER_LOCK" 2>/dev/null; then
514+
# flock blocks until acquired and is auto-released when fd 8 closes or the
515+
# process dies, so there is no stale lock and no lost increment. Only bump
516+
# the counter while the lock is actually held.
517+
if flock 8 2>/dev/null; then
518+
_ecc_bump_signal_counter
519+
flock -u 8 2>/dev/null || true
491520
fi
492-
echo "$counter" > "$SIGNAL_COUNTER_FILE"
521+
exec 8>&- 2>/dev/null || true
493522
else
494-
echo "1" > "$SIGNAL_COUNTER_FILE"
523+
# No flock (e.g. macOS): atomic mkdir lock with a bounded spin so the hook
524+
# never blocks indefinitely. A trap releases the lock on every exit path --
525+
# including the async-timeout SIGTERM -- so a killed hook does not strand the
526+
# directory. We deliberately do NOT hand-roll PID-based stale reclaim:
527+
# re-verifying then removing another process's lock is racy and can delete a
528+
# live re-acquirer's directory, reintroducing the very race this fixes.
529+
_signal_lock_held=0
530+
_signal_lock_spins=0
531+
while [ "$_signal_lock_spins" -lt 100 ]; do
532+
if mkdir "$SIGNAL_COUNTER_LOCK" 2>/dev/null; then
533+
trap 'rmdir "$SIGNAL_COUNTER_LOCK" 2>/dev/null || true' EXIT INT TERM
534+
_signal_lock_held=1
535+
break
536+
fi
537+
_signal_lock_spins=$((_signal_lock_spins + 1))
538+
sleep 0.02
539+
done
540+
if [ "$_signal_lock_held" -eq 1 ]; then
541+
# Bump only under the held lock -- never an unlocked read-modify-write.
542+
_ecc_bump_signal_counter
543+
rmdir "$SIGNAL_COUNTER_LOCK" 2>/dev/null || true
544+
trap - EXIT INT TERM
545+
fi
546+
# If the lock could not be acquired within the spin budget we skip this tick
547+
# rather than racing on an unlocked counter. Dropping one throttle tick under
548+
# extreme contention only delays the next observer signal slightly; it never
549+
# corrupts the counter or signals spuriously.
495550
fi
496551

497552
# Signal observer if running and throttle allows (check both project-scoped and global observer, deduplicate)
Lines changed: 263 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,263 @@
1+
/**
2+
* Regression tests for the SIGUSR1 throttle-counter race in observe.sh (#2296)
3+
*
4+
* observe.sh runs on every tool call and bumps a throttle counter in
5+
* ${PROJECT_DIR}/.observer-signal-counter so the observer is signaled only
6+
* every N observations (#521). The bump used a plain read-modify-write with no
7+
* locking, so concurrent hook invocations could read the same value, both
8+
* increment, and lose a write — the observer then fired at unpredictable
9+
* intervals. The fix serializes the read-modify-write with an atomic mkdir
10+
* lock.
11+
*
12+
* These tests drive the real observe.sh (reusing the stub harness from
13+
* observer-memory.test.js) and assert the lock's invariant: with the reset
14+
* threshold set high enough that no reset fires, the final counter must equal
15+
* the number of invocations — i.e. no increment is ever lost, even under heavy
16+
* concurrency.
17+
*
18+
* Run with: node tests/hooks/observe-signal-counter-race.test.js
19+
*/
20+
21+
const assert = require('assert');
22+
const path = require('path');
23+
const fs = require('fs');
24+
const os = require('os');
25+
const { spawn, spawnSync } = require('child_process');
26+
27+
let passed = 0;
28+
let failed = 0;
29+
30+
function test(name, fn) {
31+
try {
32+
fn();
33+
console.log(` ✓ ${name}`);
34+
passed++;
35+
} catch (err) {
36+
console.log(` ✗ ${name}`);
37+
console.log(` Error: ${err.message}`);
38+
failed++;
39+
}
40+
}
41+
42+
async function asyncTest(name, fn) {
43+
try {
44+
await fn();
45+
console.log(` ✓ ${name}`);
46+
passed++;
47+
} catch (err) {
48+
console.log(` ✗ ${name}`);
49+
console.log(` Error: ${err.message}`);
50+
failed++;
51+
}
52+
}
53+
54+
function createTempDir() {
55+
return fs.mkdtempSync(path.join(os.tmpdir(), 'ecc-signal-race-'));
56+
}
57+
58+
function cleanupDir(dir) {
59+
try {
60+
fs.rmSync(dir, { recursive: true, force: true });
61+
} catch {
62+
// ignore cleanup errors
63+
}
64+
}
65+
66+
const repoRoot = path.resolve(__dirname, '..', '..');
67+
const observeShPath = path.join(repoRoot, 'skills', 'continuous-learning-v2', 'hooks', 'observe.sh');
68+
69+
const isWindows = process.platform === 'win32';
70+
const hasPython = !isWindows && spawnSync('python3', ['--version']).status === 0;
71+
// When the runner has flock the lock is exact (blocking, kernel auto-release);
72+
// without it observe.sh uses a best-effort mkdir spin that may drop at most one
73+
// increment under pathological contention.
74+
const hasFlock = !isWindows && spawnSync('bash', ['-c', 'command -v flock']).status === 0;
75+
76+
// Build a self-contained observe.sh sandbox (stub detect-project.sh +
77+
// homunculus-dir.sh, SKILL_ROOT patched to the sandbox) and return its paths.
78+
function buildSandbox() {
79+
const testDir = createTempDir();
80+
const projectDir = path.join(testDir, 'project');
81+
fs.mkdirSync(projectDir, { recursive: true });
82+
83+
const skillRoot = path.join(testDir, 'skill');
84+
const scriptsDir = path.join(skillRoot, 'scripts');
85+
const scriptsLibDir = path.join(scriptsDir, 'lib');
86+
const hooksDir = path.join(skillRoot, 'hooks');
87+
fs.mkdirSync(scriptsDir, { recursive: true });
88+
fs.mkdirSync(scriptsLibDir, { recursive: true });
89+
fs.mkdirSync(hooksDir, { recursive: true });
90+
91+
fs.writeFileSync(
92+
path.join(scriptsDir, 'detect-project.sh'),
93+
[
94+
'#!/bin/bash',
95+
'PROJECT_ID="test-project"',
96+
'PROJECT_NAME="test-project"',
97+
`PROJECT_ROOT="${projectDir}"`,
98+
`PROJECT_DIR="${projectDir}"`,
99+
'CLV2_PYTHON_CMD="python3"',
100+
''
101+
].join('\n')
102+
);
103+
fs.writeFileSync(
104+
path.join(scriptsLibDir, 'homunculus-dir.sh'),
105+
[
106+
'#!/bin/bash',
107+
'_ecc_resolve_homunculus_dir() { printf "%s\\n" "$HOME/.local/share/ecc-homunculus"; }',
108+
''
109+
].join('\n')
110+
);
111+
112+
let observeContent = fs.readFileSync(observeShPath, 'utf8');
113+
observeContent = observeContent.replace(
114+
'SKILL_ROOT="$(cd "$SCRIPT_DIR/.." && pwd)"',
115+
`SKILL_ROOT="${skillRoot}"`
116+
);
117+
const testObserve = path.join(hooksDir, 'observe.sh');
118+
fs.writeFileSync(testObserve, observeContent, { mode: 0o755 });
119+
120+
return { testDir, projectDir, testObserve };
121+
}
122+
123+
// Run observe.sh once against the sandbox. Resolves when the process exits.
124+
function runObserve(testObserve, projectDir) {
125+
const input = JSON.stringify({
126+
tool_name: 'Read',
127+
tool_input: { file_path: '/tmp/test.txt' },
128+
session_id: 'test-session',
129+
cwd: projectDir
130+
});
131+
return new Promise((resolve, reject) => {
132+
const child = spawn('bash', [testObserve, 'post'], {
133+
env: {
134+
...process.env,
135+
HOME: projectDir,
136+
CLAUDE_CODE_ENTRYPOINT: 'cli',
137+
ECC_HOOK_PROFILE: 'standard',
138+
ECC_SKIP_OBSERVE: '0',
139+
CLAUDE_PROJECT_DIR: projectDir,
140+
// Reset threshold far above the invocation count, so no reset fires and
141+
// the final counter equals the number of invocations.
142+
ECC_OBSERVER_SIGNAL_EVERY_N: '100000'
143+
},
144+
stdio: ['pipe', 'ignore', 'pipe']
145+
});
146+
let stderr = '';
147+
// Fail the test on a hung hook rather than waiting forever.
148+
const timer = setTimeout(() => {
149+
child.kill('SIGKILL');
150+
reject(new Error('observe.sh timed out'));
151+
}, 20000);
152+
child.stderr.on('data', (chunk) => { stderr += chunk; });
153+
// A broken observe.sh must fail the test, not be silently swallowed.
154+
child.on('close', (code, signal) => {
155+
clearTimeout(timer);
156+
if (code === 0 && signal === null) {
157+
resolve();
158+
} else {
159+
reject(new Error(`observe.sh failed code=${code} signal=${signal}: ${stderr.trim()}`));
160+
}
161+
});
162+
child.on('error', (err) => {
163+
clearTimeout(timer);
164+
reject(err);
165+
});
166+
child.stdin.end(input);
167+
});
168+
}
169+
170+
function readCounter(projectDir) {
171+
const counterFile = path.join(projectDir, '.observer-signal-counter');
172+
if (!fs.existsSync(counterFile)) {
173+
return null;
174+
}
175+
return parseInt(fs.readFileSync(counterFile, 'utf8').trim(), 10);
176+
}
177+
178+
console.log('\n=== observe.sh signal-counter race regression (#2296) ===\n');
179+
180+
test('observe.sh uses a lock around the throttle-counter update', () => {
181+
const content = fs.readFileSync(observeShPath, 'utf8');
182+
assert.ok(
183+
content.includes('SIGNAL_COUNTER_LOCK'),
184+
'observe.sh should define a lock for the signal counter'
185+
);
186+
assert.ok(
187+
/flock 8\b/.test(content) || /mkdir "\$SIGNAL_COUNTER_LOCK"/.test(content),
188+
'observe.sh should acquire the counter lock via flock or an atomic mkdir'
189+
);
190+
});
191+
192+
test('observe.sh guards against a corrupt (non-integer) counter file', () => {
193+
const content = fs.readFileSync(observeShPath, 'utf8');
194+
assert.ok(
195+
/''\|\*\[!0-9\]\*\) counter=0/.test(content),
196+
'observe.sh should reset a non-integer counter to 0 before incrementing'
197+
);
198+
});
199+
200+
async function runSequential() {
201+
const { testDir, projectDir, testObserve } = buildSandbox();
202+
try {
203+
const N = 5;
204+
for (let i = 0; i < N; i++) {
205+
await runObserve(testObserve, projectDir);
206+
}
207+
const counter = readCounter(projectDir);
208+
assert.notStrictEqual(counter, null, 'counter file should exist after runs');
209+
assert.strictEqual(counter, N, `sequential counter should be ${N}, got ${counter}`);
210+
} finally {
211+
cleanupDir(testDir);
212+
}
213+
}
214+
215+
async function runConcurrent() {
216+
const { testDir, projectDir, testObserve } = buildSandbox();
217+
try {
218+
const K = 20;
219+
// Spawn all K before awaiting any, so they genuinely contend on the counter.
220+
const runs = [];
221+
for (let i = 0; i < K; i++) {
222+
runs.push(runObserve(testObserve, projectDir));
223+
}
224+
await Promise.all(runs);
225+
const counter = readCounter(projectDir);
226+
assert.notStrictEqual(counter, null, 'counter file should exist after concurrent runs');
227+
if (hasFlock) {
228+
// flock serializes every invocation, so no increment is ever lost. The
229+
// pre-fix unlocked code drops increments under this same contention.
230+
assert.strictEqual(
231+
counter,
232+
K,
233+
`with flock the counter must be exactly ${K}, got ${counter}`
234+
);
235+
} else {
236+
// mkdir fallback is best-effort: it may drop at most one increment if its
237+
// bounded spin is exhausted, but never the multi-increment loss the
238+
// unlocked code exhibited.
239+
assert.ok(
240+
counter >= K - 1,
241+
`mkdir fallback should keep the counter >= ${K - 1}, got ${counter}`
242+
);
243+
}
244+
} finally {
245+
cleanupDir(testDir);
246+
}
247+
}
248+
249+
(async () => {
250+
if (!isWindows && hasPython) {
251+
await asyncTest('sequential invocations increment the counter exactly once each', runSequential);
252+
await asyncTest('concurrent invocations never lose a counter increment', runConcurrent);
253+
} else {
254+
console.log(' - skipping shell-execution tests (requires non-Windows + python3)');
255+
}
256+
257+
console.log('\n=== Test Results ===');
258+
console.log(`Passed: ${passed}`);
259+
console.log(`Failed: ${failed}`);
260+
console.log(`Total: ${passed + failed}`);
261+
262+
process.exit(failed > 0 ? 1 : 0);
263+
})();

0 commit comments

Comments
 (0)