Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
20 commits
Select commit Hold shift + click to select a range
28c547c
fork 0.5.1-fork.1: root-cause fixes from dist-patch era + --global si…
Sep 22, 2026
7db7ae2
fork: global-instance pipe lock (fixed name, cwd-independent) - secon…
Sep 22, 2026
3da65c3
fork: sweep stops workers whose project config.json is gone - deleted…
Sep 22, 2026
56d54d0
fork: unify pipe-lock helpers (global+project share one code path, gl…
Sep 22, 2026
cd0f3ab
fork: decisions --limit takes the NEWEST N (was oldest-N shown as '最近…
Sep 22, 2026
2f893ef
fork: decision plausibility guard at extraction (reject table rows / …
Sep 22, 2026
9683952
fork: unresolved --json actually emits JSON now (registration dropped…
Sep 22, 2026
c63b5da
fork (audit round): 5 fixes from pattern-scan subagent
Sep 22, 2026
2ab7390
test: defuse time-bomb in forget A1 (fixture date was absolute, '天前' …
Sep 22, 2026
1377b12
fork (global mode): dual-write worker logs to each project's own watc…
Sep 22, 2026
3b7d74f
fork (hardening round): 3 root-cures from deep sweep
Sep 22, 2026
8e49450
fork CRITICAL fix: discover cache was keyed by source-db mtime only -…
Sep 22, 2026
3e7679e
fix(capture): incremental decision refresh + workflow-origin exclusio…
Sep 24, 2026
f17f52d
fix(capture): applyExtraction re-includes existing user_tags in meta_…
Sep 24, 2026
3b69246
fix(query): dedupe unresolved questions across sessions (same questio…
Sep 24, 2026
3e00e97
fix(extract): anchor only the noisy weak verbs (不用/直接用/优先用/统一用) to se…
Sep 24, 2026
5d76b3f
perf(judge): pending-transition refresh becomes incremental (seq wate…
Sep 24, 2026
80e3d5a
fix(extract): question extraction gets a fragment guard too (list/tab…
Sep 24, 2026
e4b0ba6
fix(watch): wire the anonymous stats counter into runJudge (daemon-si…
Sep 24, 2026
691084f
test: regression locks for fork additions - origin=workflow exclusion…
Sep 24, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 7 additions & 5 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion package.json
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
{
"name": "@ewanjasper/sessionrelay",
"version": "0.5.1",
"version": "0.5.1-fork.1",
"description": "会话接力 SessionRelay — 属于项目、不属于任何厂商的本地记忆层(跨 Agent 会话记忆 / 中文检索 / MCP / HOP 交接协议)",
"license": "MIT",
"type": "module",
Expand Down
26 changes: 19 additions & 7 deletions src/adapters/claude-code/watcher.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,13 +29,11 @@ export function watchDir(
let watcher: fs.FSWatcher | null = null;
let pollTimer: NodeJS.Timeout | null = null;

try {
watcher = fs.watch(dir, { recursive: true }, (_ev, filename) => {
if (!filename) return;
onEvent(path.join(dir, filename.toString()));
});
} catch {
// Linux 递归不支持 → mtime 轮询兜底(技术方案 §9.1)
// [fork 0922] 轮询兜底提为函数:原本只在 fs.watch 同步创建失败(Linux)时进入,
// 但运行中的 error 事件(目录被删/句柄溢出/杀软干扰)无人监听会让 EventEmitter
// 直接 throw——全局守护一死全部项目捕获停摆。error 时降级轮询自愈。
const startPolling = () => {
if (pollTimer) return;
const snap = new Map<string, number>();
const scan = () => {
try {
Expand All @@ -59,6 +57,20 @@ export function watchDir(
};
scan(); // 建立基线
pollTimer = setInterval(scan, 1000);
};

try {
watcher = fs.watch(dir, { recursive: true }, (_ev, filename) => {
if (!filename) return;
onEvent(path.join(dir, filename.toString()));
});
watcher.on('error', () => {
try { watcher?.close(); } catch { /* 已关 */ }
watcher = null;
startPolling(); // 事件流断了就退化为 1s 轮询,捕获不断线
});
} catch {
startPolling();
}

return {
Expand Down
11 changes: 11 additions & 0 deletions src/adapters/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,14 @@ export interface SessionSourceAdapter {
* 返回该会话的 compaction 信息;无压缩返回 null
*/
detectCompaction?(ds: DiscoveredSession, config: AdapterConfig): CompactionInfo | null;

/**
* [fork 0922] 内容变更探针:一次调用返回"会话ID → 内容签名"全表快照。
* 供 runSync 水位线跳过未变更会话——每周期总共一次聚合查询,而非每会话逐条探测。
* 签名必须由内容决定(如 MAX(rowid)),不得用 mtime:内容变更未必 bump mtime(上游测试契约)。
* 未实现时 runSync 退回 mtimeMs+sizeBytes 签名(对文件型源已足够)。
*/
changeProbe?(config: AdapterConfig): Map<string, string>;
}

/** compaction 检测结果 */
Expand All @@ -65,6 +73,9 @@ export interface DiscoveredSession {
updatedAt?: string;
sizeBytes: number;
mtimeMs: number;
/** [fork 0924] 会话来源提示:适配器识别出"非项目正主"的会话(如工作流子代理)时标注,
* 入库时写入 origin 列,决策检索默认排除——防施工碎片淹没项目决策 */
originHint?: 'workflow';
}

/** 读取结果 */
Expand Down
45 changes: 41 additions & 4 deletions src/adapters/zcode/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,14 +27,25 @@ export function resetConn(): void {
cachedPath = null;
}

// [fork 0922] discover 结果缓存:源库未被写入(文件 mtime 未变)时直接复用上一轮结果,
// 空闲周期 0 SQL。原实现每周期对 ZCode 库全表 lower(directory) 扫描(无索引),
// 全局守护 N 个 worker × 每个事件都各扫一遍。
// ⚠️ 键必须含项目根:多个项目共享同一源库,只按 mtime 键会让 B 项目拿到 A 项目的
// 会话列表并错误入库(0922 晚实测污染:music 42→559、novel 15→495)。
const discoverCache = new Map<string, { mtimeMs: number; rows: DiscoveredSession[] }>();

export function discover(projectRoot: string, dbPath: string): DiscoveredSession[] {
if (!fs.existsSync(dbPath)) return [];
const st = fs.statSync(dbPath);
const cacheKey = `${path.resolve(projectRoot).toLowerCase()}\u0000${path.resolve(dbPath).toLowerCase()}`;
const hit = discoverCache.get(cacheKey);
if (hit && hit.mtimeMs === st.mtimeMs) return hit.rows;
const z = getConn(dbPath);
try {
const rows = z
.prepare('SELECT id, title, time_created, time_updated FROM session WHERE lower(directory) = lower(?)')
.all(path.resolve(projectRoot)) as Array<{ id: string; title: string; time_created: number; time_updated: number }>;
return rows.map((r) => ({
const out = rows.map((r) => ({
source: SOURCE_ID,
sourceSessionId: r.id,
sourceFile: `zcode:${r.id}`,
Expand All @@ -43,7 +54,12 @@ export function discover(projectRoot: string, dbPath: string): DiscoveredSession
updatedAt: new Date(r.time_updated).toISOString(),
sizeBytes: 0,
mtimeMs: r.time_updated,
// [fork 0924] 工作流子代理会话标注 origin=workflow:它们是施工日志不是项目会话,
// 决策提取/检索默认排除(mc整合包实测:88% 决策碎片来自它们)
originHint: /^workflow subagent/i.test(r.title ?? '') ? ('workflow' as const) : undefined,
}));
discoverCache.set(cacheKey, { mtimeMs: st.mtimeMs, rows: out });
return out;
} finally {
// 使用缓存连接,不关闭(由 resetConn 或进程退出时关闭)
}
Expand Down Expand Up @@ -153,7 +169,9 @@ export function readNew(ds: DiscoveredSession, dbPath: string, cursor: unknown):
// ── 改动 3:compaction 检测 ──
export function detectCompaction(ds: DiscoveredSession, dbPath: string): CompactionInfo | null {
if (!fs.existsSync(dbPath)) return null;
const z = new Database(dbPath, { readonly: true });
// [fork] 复用 getConn 缓存连接——原实现每个会话每周期 new Database 却从不关闭
// (注释"缓存连接"系自 getConn 误复制),连接全靠 GC 兜底 → 句柄泄漏 + 内存锯齿 + CPU 空烧(上游 #2 根因)
const z = getConn(dbPath);
try {
z.pragma('busy_timeout = 3000');
const comp = z.prepare(`
Expand All @@ -174,7 +192,7 @@ export function detectCompaction(ds: DiscoveredSession, dbPath: string): Compact
summaryMessageId: data.summaryMessageId,
};
} finally {
// 缓存连接,不关闭
// 复用缓存连接,不关闭(与注释语义一致)
}
}

Expand All @@ -197,7 +215,7 @@ export const adapter: SessionSourceAdapter = {
if (!fs.existsSync(dbPath)) return `数据库不存在:${dbPath}(未安装 ZCode 可忽略)`;
try {
const z = new Database(dbPath, { readonly: true });
// 缓存连接,不关闭
z.close(); // [fork] 探测完即关:一次性连接不留着等 GC(原实现泄漏句柄,doctor 路径)
return null;
} catch (e) {
return `只读探测失败:${(e as Error).message}`;
Expand All @@ -206,4 +224,23 @@ export const adapter: SessionSourceAdapter = {
detectCompaction(ds, config) {
return detectCompaction(ds, config.dbPath as string);
},

// [fork 0922] 内容签名:两条聚合查询拿全表 MAX(rowid)(消息+part,part 覆盖 compaction 追加)。
// 内容变更必动 rowid;mtime 不可靠(测试夹具证明追加可以不 bump time_updated)。
changeProbe(config) {
const dbPath = config.dbPath as string;
const out = new Map<string, string>();
if (!fs.existsSync(dbPath)) return out;
const z = getConn(dbPath);
try {
z.pragma('busy_timeout = 3000');
const msg = z.prepare('SELECT session_id, MAX(rowid) AS m FROM message GROUP BY session_id').all() as Array<{ session_id: string; m: number }>;
const part = z.prepare('SELECT session_id, MAX(rowid) AS m FROM part GROUP BY session_id').all() as Array<{ session_id: string; m: number }>;
for (const r of msg) out.set(r.session_id, String(r.m));
for (const r of part) out.set(r.session_id, `${out.get(r.session_id) ?? '0'}/${r.m}`);
return out;
} catch {
return out; // 探针失败 → 空签名 → 全部不命中水位线 → 全量 ingest(安全侧)
}
},
};
5 changes: 4 additions & 1 deletion src/bin/srelay.ts
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ program
.command('watch')
.description('守护捕获(前台运行;--install-service 注册系统服务)')
.option('--foreground', '前台运行(服务调用路径)')
.option('--global', '全局守护:一个进程看管项目注册表里的全部项目(动态收编新项目)')
.option('--install-service', '注册守护服务(Windows 计划任务)')
.option('--uninstall', '卸载守护服务')
.option('--status', '查看守护与服务状态')
Expand Down Expand Up @@ -220,7 +221,9 @@ program
.option('--json')
.action(async (opts) => {
const { cmdUnresolved } = await import('../cli/meta.js');
await cmdUnresolved({ limit: opts.limit ? Number(opts.limit) : undefined });
// [fork 0922] 原实现丢弃 opts.json(选项存在却没传)——unresolved --json 永远输出人话,
// 依赖 JSON 的调用方(开场简报)被迫走文本兜底、丧失一切结构化过滤
await cmdUnresolved({ limit: opts.limit ? Number(opts.limit) : undefined, json: opts.json });
});

program
Expand Down
9 changes: 6 additions & 3 deletions src/capture/archive.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,10 @@ function isProtected(s: SessionRow): string | null {
export function runArchive(db: DB, opts: ArchiveOptions): ArchiveResult {
const result: ArchiveResult = { archived: 0, skipped: 0, bytesFreed: 0, details: [] };
const now = new Date();
const cutoff = opts.days
// [fork 0922] days/sizeMb 为 0 是合法值("归档 0 天前"="全归档"边界语义由调用方定),
// 原实现的 truthy 判断把 0 当"未提供"→ cutoff 退化成 9999-12-31、预算退化成无限——
// 用户要 0 条实际归档全部(配 --hard 即灾难)。改显式 != null 判定。
const cutoff = opts.days != null && Number.isFinite(opts.days)
? new Date(now.getTime() - opts.days * 86400_000).toISOString()
: opts.before ?? '9999-12-31';

Expand Down Expand Up @@ -83,8 +86,8 @@ export function runArchive(db: DB, opts: ArchiveOptions): ArchiveResult {
});
}

// 按体积限制(如果指定了 sizeMb)
let bytesBudget = opts.sizeMb ? opts.sizeMb * 1024 * 1024 : Infinity;
// 按体积限制(如果指定了 sizeMb;0 也按 0 尊重,同上)
let bytesBudget = opts.sizeMb != null && Number.isFinite(opts.sizeMb) ? opts.sizeMb * 1024 * 1024 : Infinity;

for (const s of candidates) {
// 保护规则
Expand Down
10 changes: 8 additions & 2 deletions src/capture/judge.ts
Original file line number Diff line number Diff line change
@@ -1,7 +1,7 @@
// 判定 tick(方针 §6.1 两阶段提交;Phase 2 将在 confirmed 副作用中加元数据提取与 summary_rule)
// now 可注入(技术方案 T5:冷却期逻辑必须可测,不能真等 6 小时)
import type { DB } from '../store/db.js';
import { dueIdle, dueConfirm, markPending, confirmSession } from '../store/db.js';
import { dueIdle, dueConfirm, markPending, confirmSession, refreshExtraction, refreshDecisionsIncremental } from '../store/db.js';
import type { StatsCounter } from '../core/stats/counter.js';

export interface JudgeOptions {
Expand All @@ -19,7 +19,13 @@ export function runJudge(db: DB, o: JudgeOptions): JudgeResult {
const cooldownCutoff = new Date(o.now.getTime() - o.cooldownH * 3_600_000).toISOString();

const idle = dueIdle(db, o.projectId, idleCutoff);
for (const id of idle) markPending(db, id, o.now.toISOString());
for (const id of idle) {
markPending(db, id, o.now.toISOString());
// [fork 0924] 1a:pending 转换即刷新决策——长活会话不再冻结到 6h 确认。
// [fork 0924b] 改增量:只提取上次水位之后的新消息(jieba 只扫增量),
// 全量 topics/summary 重建留给 confirm——pending 时不再有 O(全会话) 的 CPU 突发
refreshDecisionsIncremental(db, id);
}

// confirmed 统一走 confirmSession(提取元数据 + summary_rule + meta_text,方针 §6.6)
const toConfirm = dueConfirm(db, o.projectId, cooldownCutoff);
Expand Down
36 changes: 33 additions & 3 deletions src/capture/sync.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,17 @@ function titleFromMessages(msgs: Array<{ role: string; content: string }>): stri
return first.content.replace(/\s+/g, ' ').trim().slice(0, 60);
}

// [fork 0922] 会话水位线:内容签名未变 且 会话仍在库里 → 整跳过(discover 仍跑,ingestOne 清零)。
// 背景:N 个项目守护各自监听同一批源目录,源库任何写入都触发所有守护各跑一轮 cycle,
// 每个 cycle 又对每个已知会话跑 3-4 条查询——活跃时段单守护 7.5% 单核,守护越多放大越狠。
// 签名铁律:内容由 adapter.changeProbe(zcode=MAX(rowid) 聚合,每周期 2 条 SQL)决定;
// 文件型源退回 mtimeMs+sizeBytes。不得只用 mtime——内容变更未必 bump mtime(上游测试契约)。
// "仍在库里"缺一不可(上游测试抓出的真缺陷):rebuild/forget/手动删行后库与签名脱钩,
// 没有它重建后的库永远收不到重摄。已知会话集=每周期一条本地聚合查询,代价毫秒级。
// 内存态即够:守护重启后首轮全量(本就是补账需求);10 分钟一次全量清扫兜底一切"以为没变其实变了"。
const syncWatermark = new Map<string, string>();
let lastFullSweep = 0;

export async function runSync(opts: SyncOptions): Promise<SyncStats> {
const root = opts.projectRoot;
const cfg = opts.config;
Expand All @@ -43,14 +54,22 @@ export async function runSync(opts: SyncOptions): Promise<SyncStats> {

if (mode === 'off') return result;

// 改动 1:注册表初始化(含 custom adapter 加载)
ensureRegistered(root);
// [fork 0922] 10 分钟全量清扫:清水位线跑一轮完整 ingest,堵一切"以为没变其实变了"的边角
if (Date.now() - lastFullSweep > 600_000) syncWatermark.clear();
lastFullSweep = Date.now();

// 改动 1:注册表初始化(含 custom adapter 加载)——[fork] 原 here 连调两次,
// 第二次因 customLoaded 守卫恒返回空结果,custom adapter 的加载错误从未上报
const customResult = ensureRegistered(root);
for (const err of customResult.errors) result.warnings.push(`custom adapter 加载失败:${err}`);

const own = !opts.db;
const db = opts.db ?? createDb(dbFile(root));
const projectId = cfg.identity.project_id ?? projectIdOf(root);
// [fork 0922] 已知会话集:库里现存 (source, source_session_id)——水位线跳过的第二个必要条件
const knownSessions = new Set<string>(
(db.prepare("SELECT source || ':' || source_session_id AS k FROM sessions WHERE project_id = ?").all(projectId) as Array<{ k: string }>).map(r => r.k)
);
const ignoreRules = loadIgnoreRules(root);
// forget 防复活次级防线(设计 v4 §3.2):入口整表载入一次,ingest 内 Set 判定
const tombstones = loadTombstones(db);
Expand All @@ -68,6 +87,8 @@ export async function runSync(opts: SyncOptions): Promise<SyncStats> {
}
const aConfig = adapterConfig(cfg, source);
const discovered = adapter.discover(root, aConfig);
// [fork 0922] 内容签名探针:每源每周期一次聚合查询;无探针的源退回 mtime+size
const probe = adapter.changeProbe?.(aConfig);

for (const ds of discovered) {
result.discovered++;
Expand All @@ -86,7 +107,15 @@ export async function runSync(opts: SyncOptions): Promise<SyncStats> {
continue;
}

// [fork 0922] 水位线:内容签名未变 且 会话仍在库里 → 跳过(rebuild/forget/删行后自动失效)
const wmKey = `${ds.source}:${ds.sourceSessionId}`;
const sig = probe
? probe.get(ds.sourceSessionId)
: (ds.mtimeMs !== undefined ? `${ds.mtimeMs}:${ds.sizeBytes}` : undefined);
if (sig !== undefined && syncWatermark.get(wmKey) === sig && knownSessions.has(wmKey)) continue;

await ingestOne(db, ds, { mode, projectId, cfg, result, stats: opts.stats, ignoreRules, tombstones, source, aConfig });
if (sig !== undefined) syncWatermark.set(wmKey, sig);
}
}
} finally {
Expand Down Expand Up @@ -202,7 +231,8 @@ async function ingestOne(
createdAt: ds.createdAt ?? lastEventAt,
lastEventAt,
sourceFile: ds.sourceFile,
origin: ctx.origin,
// [fork 0924] 适配器的来源提示(如 zcode 工作流子代理→workflow)优先于默认 auto
origin: ds.originHint ?? ctx.origin,
});
if (up.isNew) ctx.result.newSessions++;

Expand Down
Loading