diff --git a/src/slack/bot.ts b/src/slack/bot.ts index 7a20e635..6cd46948 100644 --- a/src/slack/bot.ts +++ b/src/slack/bot.ts @@ -22,7 +22,8 @@ import { SlackSocketClient, type SlackEnvelope } from './socket.js'; import { resolveEventText, shouldAttachSlack, shouldProcessSlackEvent, type SlackMessageEvent } from './events.js'; import { isThreadParticipated, markThreadParticipated, - claimThreadPrefetch, releaseThreadPrefetch, resetThreadPrefetchClaims, + claimThreadPrefetch, commitThreadPrefetch, + releaseThreadPrefetch, resetThreadPrefetchClaims, } from './thread-tracker.js'; import { sendSlackText, getSlackSendClient } from './send-only-client.js'; import { startSlackProgress, statusFromToolEvent } from './progress.js'; @@ -237,7 +238,11 @@ export async function processSlackMessageEvent( // rather than a return-by-return audit. let prefetchCommitted = false; try { - await runSlackMessageEvent(event, target, text, signal, opts, () => { prefetchCommitted = true; }); + await runSlackMessageEvent(event, target, text, signal, opts, () => { + prefetchCommitted = Boolean(opts.prefetchToken) && commitThreadPrefetch( + event.channel || '', event.thread_ts || '', opts.prefetchToken || 0, + ); + }); } finally { if (opts.prefetchToken && !prefetchCommitted) { releaseThreadPrefetch(event.channel || '', event.thread_ts || '', opts.prefetchToken); @@ -454,7 +459,8 @@ export async function handleSlackEnvelope(envelope: SlackEnvelope): Promise - processSlackMessageEvent(event, target, text, signal, { - prefetchToken, - ...(reservedEventKey ? { eventKey: reservedEventKey } : {}), - ...(reservationGeneration !== undefined ? { reservationGeneration } : {}), - })); + prefetchHandedOff = enqueueSlackIngress(slackIngressLaneKey(target), signal => + processSlackMessageEvent(event, target, text, signal, { + prefetchToken, + ...(reservedEventKey ? { eventKey: reservedEventKey } : {}), + ...(reservationGeneration !== undefined ? { reservationGeneration } : {}), + })); + } finally { + if (prefetchToken && !prefetchHandedOff) { + releaseThreadPrefetch(event.channel || '', event.thread_ts || '', prefetchToken); + } + } } // ─── Init / Shutdown ──────────────────────────────── diff --git a/src/slack/ingress.ts b/src/slack/ingress.ts index e437aebb..4f67b23e 100644 --- a/src/slack/ingress.ts +++ b/src/slack/ingress.ts @@ -180,8 +180,8 @@ export function slackIngressLaneKey(target: RemoteTarget): string { export function enqueueSlackIngress( laneKey: string, task: (signal: AbortSignal) => Promise, -): void { - if (resetting) return; +): boolean { + if (resetting) return false; const taskGeneration = generation; const controller = new AbortController(); controllers.add(controller); @@ -201,6 +201,7 @@ export function enqueueSlackIngress( void tail.then(() => { if (ingressTails.get(laneKey) === tail) ingressTails.delete(laneKey); }); + return true; } export type SlackRunContext = { diff --git a/src/slack/thread-tracker.ts b/src/slack/thread-tracker.ts index 9dafbe95..4c521662 100644 --- a/src/slack/thread-tracker.ts +++ b/src/slack/thread-tracker.ts @@ -92,8 +92,9 @@ export function isThreadParticipated(channel: string, threadTs: string): boolean // too, so re-injecting the thread's history is the RIGHT behavior. Persisting // the claim would leave a context-less session permanently without context. -/** key -> the token of the claim currently holding it. */ -const prefetchClaimed = new Map(); +type PrefetchClaim = { token: number; committed: boolean }; +/** key -> current owner and whether history was actually injected. */ +const prefetchClaimed = new Map(); const PREFETCH_CLAIM_CAP = 500; let prefetchToken = 0; @@ -115,17 +116,36 @@ export function claimThreadPrefetch(channel: string, threadTs: string): number { const key = threadKey(channel, threadTs); if (prefetchClaimed.has(key)) return 0; if (prefetchClaimed.size >= PREFETCH_CLAIM_CAP) { - // Oldest half by insertion order — a claim is never refreshed, so - // insertion order IS recency here. - for (const [stale] of [...prefetchClaimed].slice(0, Math.floor(PREFETCH_CLAIM_CAP / 2))) { + // Active owners are singleflight locks, not cache entries. Evicting one + // lets another envelope claim the same live thread and inject history + // twice. Only completed claims may give ground under pressure. + let removed = 0; + const target = Math.floor(PREFETCH_CLAIM_CAP / 2); + for (const [stale, claim] of prefetchClaimed) { + if (!claim.committed) continue; prefetchClaimed.delete(stale); + removed += 1; + if (removed >= target) break; } + // All bounded slots can legitimately be in flight. Decline rather than + // queue or violate singleflight; a later message can retry after one + // owner commits or releases. + if (prefetchClaimed.size >= PREFETCH_CLAIM_CAP) return 0; } const token = ++prefetchToken; - prefetchClaimed.set(key, token); + prefetchClaimed.set(key, { token, committed: false }); return token; } +/** Mark that this owner actually injected history; completed claims are evictable. */ +export function commitThreadPrefetch(channel: string, threadTs: string, token: number): boolean { + if (!channel || !threadTs || !token) return false; + const claim = prefetchClaimed.get(threadKey(channel, threadTs)); + if (!claim || claim.token !== token) return false; + claim.committed = true; + return true; +} + /** * Give a claim back when no history was actually injected. * @@ -138,7 +158,7 @@ export function claimThreadPrefetch(channel: string, threadTs: string): number { export function releaseThreadPrefetch(channel: string, threadTs: string, token: number): void { if (!channel || !threadTs || !token) return; const key = threadKey(channel, threadTs); - if (prefetchClaimed.get(key) !== token) return; + if (prefetchClaimed.get(key)?.token !== token) return; prefetchClaimed.delete(key); } diff --git a/structure/str_func.md b/structure/str_func.md index 0ce2c0d4..0b8ebddf 100644 --- a/structure/str_func.md +++ b/structure/str_func.md @@ -239,11 +239,11 @@ cli-jaw/ │ │ └── discord-file.ts ← Discord 파일 전송 (67L) │ ├── slack/ ← Slack 인터페이스 (20 files, Socket Mode + Web API, SDK 없음) │ │ ├── socket.ts ← Socket Mode client (apps.connections.open → wss, ack-before-work, envelope dedupe TTL, hello deadline, backoff 재연결) (372L) -│ │ ├── bot.ts ← Slack 봇 lifecycle + envelope routing + orchestrate 경로 + queued-result waiter (627L) +│ │ ├── bot.ts ← Slack 봇 lifecycle + envelope routing + orchestrate 경로 + queued-result waiter (639L) │ │ ├── api.ts ← Slack Web API fetch wrapper (HTTP 200 + ok:false를 실패로 처리, credential/URL redaction) (168L) │ │ ├── format.ts ← CommonMark → mrkdwn 변환 + code-fence 보존 chunking (62L) │ │ ├── events.ts ← inbound gating (self-echo/bot/subtype/allowlist/mention) + Block Kit 텍스트 추출 (216L) -│ │ ├── thread-tracker.ts ← 참여 스레드 영속 추적 (mention/봇응답 마킹, 캡드 셋, 무멘션 스레드 연속 대화 게이트 지원) (154L) +│ │ ├── thread-tracker.ts ← 참여 스레드 영속 추적 (mention/봇응답 마킹, 캡드 셋, 무멘션 스레드 연속 대화 게이트 지원) (174L) │ │ ├── enrichment-cache.ts ← 공용 동시성 프리미티브 (TTL/cap 캐시, 원인별 억제, 능력 잠금 단일 재탐침, in-flight 합류, 집계 취소, 세대 무효화) (425L) │ │ ├── conversation.ts ← 대화/스레드 컨텍스트 (conversations.info + replies, 참여자는 author 유도, method별 억제·시작률) (347L) │ │ ├── context.ts ← 프롬프트 컨텍스트 블록 조립 (채널 id·thread_ts 무절단, 섹션별 코드포인트 예산, 신뢰 경계 문구 보존) (239L) @@ -251,7 +251,7 @@ cli-jaw/ │ │ ├── attachment-recovery.ts ← app_mention 봉투에 없는 첨부를 channel+ts 재조회로 복구 (oldest+inclusive+limit=1) (53L) │ │ ├── commands.ts ← slash command → 공유 parseCommand/executeCommand 파이프라인 (148L) │ │ ├── slack-file.ts ← files.getUploadURLExternal → upload → completeUploadExternal 3단계 업로드 (97L) -│ │ ├── ingress.ts ← 세션별 ingress lane + admitSlackRun 동기 실행 예약(sessionLanes) + 전역 다운로드 세마포어 + shutdown abort/drain (275L) ✨ +│ │ ├── ingress.ts ← 세션별 ingress lane + admitSlackRun 동기 실행 예약(sessionLanes) + 전역 다운로드 세마포어 + shutdown abort/drain (276L) ✨ │ │ ├── inbound-file.ts ← 인바운드 첨부 단일 IO owner (files.info → 인증 스트리밍 다운로드 → saveUpload, 파일/메시지 바이트 예산, 고정 error code) (280L) ✨ │ │ ├── inbound-url.ts ← 인바운드 다운로드 URL 검증 (Slack host allowlist + https-only hop + 사설망 거부) (44L) ✨ │ │ ├── send-only-client.ts ← bot-token 전용 outbound + conversations.open DM 해석 (69L) diff --git a/tests/unit/safe-install.test.ts b/tests/unit/safe-install.test.ts index bcd53b71..7100dca4 100644 --- a/tests/unit/safe-install.test.ts +++ b/tests/unit/safe-install.test.ts @@ -459,10 +459,15 @@ test('SAF-004g: postinstall child processes use service-safe PATH consistently', }); test('SAF-004h: README scopes Windows installation support to WSL', () => { + const nativeWindowsBetaStart = readmeSrc.indexOf('Native Windows (PowerShell beta)'); + const nativeWindowsBetaEnd = readmeSrc.indexOf('', nativeWindowsBetaStart); + const readmeOutsideNativeWindowsBeta = nativeWindowsBetaStart >= 0 && nativeWindowsBetaEnd >= 0 + ? readmeSrc.slice(0, nativeWindowsBetaStart) + readmeSrc.slice(nativeWindowsBetaEnd + ''.length) + : readmeSrc; assert.ok(readmeSrc.includes('wsl --install'), 'README should document Windows setup through WSL'); assert.ok(readmeSrc.includes('wsl.exe -d Ubuntu -- bash -lc "jaw dashboard"'), 'README should document PowerShell-to-WSL login-shell invocation'); assert.ok(readmeSrc.includes('macOS / Linux / WSL with Node.js 22+ already installed'), 'README default npm install block should be OS-scoped'); - assert.equal(readmeSrc.includes('Get-Command jaw'), false, 'README must not troubleshoot native PowerShell jaw resolution as a supported path'); + assert.equal(readmeOutsideNativeWindowsBeta.includes('Get-Command jaw'), false, 'README must keep native PowerShell jaw resolution inside the explicitly scoped beta section'); assert.equal(localizedReadmeSrc.includes('$env:JAW_SAFE="1"; npm install -g cli-jaw'), false, 'localized READMEs must not advertise native PowerShell safe install'); assert.equal(localizedReadmeSrc.includes('# Windows PowerShell'), false, 'localized READMEs must not present native PowerShell install snippets'); assert.equal(localizedReadmeSrc.includes('npm bin -g'), false, 'localized README troubleshooting should use npm prefix -g, not removed npm bin -g'); diff --git a/tests/unit/slack-thread-prefetch.test.ts b/tests/unit/slack-thread-prefetch.test.ts index e7835526..6cb2c1b6 100644 --- a/tests/unit/slack-thread-prefetch.test.ts +++ b/tests/unit/slack-thread-prefetch.test.ts @@ -9,17 +9,99 @@ // thread BEFORE the ingress task runs, so a participation check inside that task // is a dead branch (#316). -import test from 'node:test'; +import test, { mock } from 'node:test'; import assert from 'node:assert/strict'; +import { rmSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { settings } from '../../src/core/config.ts'; import { claimThreadPrefetch, + commitThreadPrefetch, releaseThreadPrefetch, resetThreadPrefetchClaims, + resetThreadTrackerForTest, } from '../../src/slack/thread-tracker.ts'; import { buildThreadPreamble, PREAMBLE_TOTAL_CAP } from '../../src/slack/context.ts'; -test.beforeEach(() => resetThreadPrefetchClaims()); +let recoverAttachments: () => Promise = async () => []; + +mock.module('../../src/slack/attachment-recovery.ts', { + namedExports: { + recoverSlackAttachments: async () => recoverAttachments(), + }, +}); + +mock.module('../../src/orchestrator/gateway.ts', { + namedExports: { + submitMessage: () => ({ action: 'started', requestId: 'R-prefetch' }), + }, +}); + +mock.module('../../src/orchestrator/collect.ts', { + namedExports: { orchestrateAndCollect: async () => 'reply' }, +}); + +mock.module('../../src/slack/send-only-client.ts', { + namedExports: { + getSlackSendClient: () => ({ token: 'xoxb-test' }), + sendSlackText: async () => ({ ok: true }), + }, +}); + +mock.module('../../src/slack/forwarder.ts', { + namedExports: { + createSlackForwarder: () => () => { }, + relaySlackImages: async () => { }, + }, +}); + +const { handleSlackEnvelope } = await import('../../src/slack/bot.ts'); +const { enqueueSlackIngress, resetSlackIngress } = await import('../../src/slack/ingress.ts'); + +const trackerPath = join(tmpdir(), `cli-jaw-prefetch-${process.pid}.json`); + +test.beforeEach(async () => { + await resetSlackIngress(); + resetThreadPrefetchClaims(); + resetThreadTrackerForTest(trackerPath); + recoverAttachments = async () => []; + settings.slack.channelIds = []; + settings.slack.mentionOnly = true; + settings.slack.threadRequireMention = false; +}); + +test.after(() => { + resetThreadTrackerForTest(); + rmSync(trackerPath, { force: true }); + rmSync(`${trackerPath}.tmp`, { force: true }); +}); + +function deferred(): { promise: Promise; resolve: (value: T) => void } { + let resolve!: (value: T) => void; + const promise = new Promise(done => { resolve = done; }); + return { promise, resolve }; +} + +function threadedEnvelope(text: string, suffix: string) { + return { + envelope_id: `E-${suffix}`, + type: 'events_api', + payload: { + event: { + type: 'app_mention', channel: `C-${suffix}`, user: 'U1', text, + ts: `${suffix}.2`, thread_ts: `${suffix}.1`, + }, + }, + } as const; +} + +function assertClaimReleased(suffix: string): void { + const token = claimThreadPrefetch(`C-${suffix}`, `${suffix}.1`); + assert.ok(token > 0, `thread ${suffix} remained claimed`); + releaseThreadPrefetch(`C-${suffix}`, `${suffix}.1`, token); +} test('a thread is claimable exactly once', () => { const first = claimThreadPrefetch('C1', '100.1'); @@ -74,6 +156,88 @@ test('reset clears every claim', () => { assert.ok(claimThreadPrefetch('C1', '100.1') > 0, 'a new runtime re-injects history'); }); +test('capacity pressure never evicts an active claim', () => { + const tokens: number[] = []; + for (let i = 0; i < 500; i += 1) { + tokens.push(claimThreadPrefetch('C1', `${i}.1`)); + } + assert.ok(tokens.every(token => token > 0)); + assert.equal( + claimThreadPrefetch('C1', 'overflow.1'), 0, + 'a new prefetch must degrade while every bounded slot is active', + ); + assert.equal( + claimThreadPrefetch('C1', '0.1'), 0, + 'the oldest live owner must remain claimed under pressure', + ); +}); + +test('capacity pressure may evict completed claims but preserves active ones', () => { + const active = claimThreadPrefetch('C1', 'active.1'); + for (let i = 0; i < 499; i += 1) { + const ts = `done-${i}.1`; + const token = claimThreadPrefetch('C1', ts); + assert.ok(commitThreadPrefetch('C1', ts, token)); + } + assert.ok(claimThreadPrefetch('C1', 'new.1') > 0, 'completed entries make bounded room'); + assert.equal(claimThreadPrefetch('C1', 'active.1'), 0, 'the live owner is never evicted'); + releaseThreadPrefetch('C1', 'active.1', active); +}); + +test('an accepted envelope that becomes empty releases its prefetch claim', async () => { + // The gate accepts whitespace as a present text field, but normalization + // below the claim turns it into an empty prompt and returns before enqueue. + await handleSlackEnvelope(threadedEnvelope(' ', 'empty')); + assertClaimReleased('empty'); +}); + +test('a reset handled before enqueue releases its prefetch claim', async () => { + await handleSlackEnvelope(threadedEnvelope('reset', 'reset')); + assertClaimReleased('reset'); +}); + +test('an attachment-recovery exception releases its prefetch claim', async () => { + recoverAttachments = async () => { throw new Error('recovery failed'); }; + await assert.rejects( + handleSlackEnvelope(threadedEnvelope('inspect attachment', 'recover-error')), + /recovery failed/, + ); + assertClaimReleased('recover-error'); +}); + +test('an ingress reset that refuses handoff releases the caller-owned claim', async () => { + const blocker = deferred(); + assert.equal( + enqueueSlackIngress('prefetch-reset-blocker', async () => blocker.promise), true, + 'the blocker must be accepted before reset starts', + ); + await Promise.resolve(); + + const recoveryEntered = deferred(); + const recoveryResult = deferred(); + recoverAttachments = async () => { + recoveryEntered.resolve(); + return recoveryResult.promise; + }; + + const handling = handleSlackEnvelope(threadedEnvelope('continue', 'reset-race')); + await recoveryEntered.promise; + const resetting = resetSlackIngress(); + try { + assert.equal( + enqueueSlackIngress('prefetch-reset-probe', async () => { }), false, + 'ingress must report that it refused ownership during reset', + ); + recoveryResult.resolve([]); + await handling; + assertClaimReleased('reset-race'); + } finally { + recoveryResult.resolve([]); + blocker.resolve(); + await resetting; + } +}); + // ─── preamble rendering ───────────────────────────── test('the preamble is delimited and labelled with the reply count', () => {