首页
/ Pi Session Search 深度解析:跨后端会话搜索的统一查询 API、扫描实现与 SQLite FTS 后端

Pi Session Search 深度解析:跨后端会话搜索的统一查询 API、扫描实现与 SQLite FTS 后端

2026-09-06 12:04:51作者:蔡怀权

Pi 的 Session Search 是 @earendil-works/pi-agent-core 提供的一个轻量级查询接口,作用域限定在**已提交的会话条目(committed session entries)**之上:共享契约只返回稳定的命中标识 (sessionId, entryId),而片段、时间戳、分数、排名等展示数据由各具体后端自行扩展。读完本篇,你将掌握该搜索契约的完整 API 形状、扫描式(Scanning)默认实现的参数细节、SQLite FTS5 后端的懒加载机制与触发器同步原理,以及如何为 JSONL 会话自行搭建 Elasticsearch 索引层。

核心 API:最小化的共享契约

搜索契约定义在 search/index.ts,并通过 包入口export * from "./search/index.ts")作为 @earendil-works/pi-agent-core 的公共 API 导出。核心类型共三个:

export interface SessionSearchHit {
  /** Logical identifier of the session that owns the entry. */
  readonly sessionId: string;

  /** Logical identifier of the entry within that session. */
  readonly entryId: string;
}

export interface SessionSearchOptions {
  /** Restrict results to specific canonical entry types. */
  readonly entryTypes?: readonly Entry["type"][];

  /** Maximum number of hits to return. Backends may return fewer, not more. */
  readonly limit?: number;

  /** Abort signal for cancellation, e.g. search-as-you-type. */
  readonly signal?: AbortSignal;
}

export interface SessionSearch<T extends SessionSearchHit = SessionSearchHit> {
  search(text: string, options?: SessionSearchOptions): AsyncIterable<T>;
}

基础命中(base hit)被刻意设计得极简:(sessionId, entryId) 是能在 JSONL、内存、SQLite FTS 和远程索引之间移植的逻辑身份。片段(snippet)、时间戳、相关度分数、元数据、偏移量与排序语义,全部归具体实现所有。这种"契约只管身份、后端自扩展"的拆分,使得上层 UI 代码可以只依赖 SessionSearch<T> 接口编程,而换后端时只需替换泛型参数 T

三个选项参数的语义边界在源码中可以直接验证:

  • entryTypes:按规范条目类型过滤。源码中 entryTypes 为空数组(length === 0)时直接返回空结果,见 scanning.ts 的 search 实现SQLite 后端的 search 方法
  • limit:命中上限,"只可能更少、不可能更多"(Backends may return fewer, not more)。扫描实现达到上限后立即 return 结束迭代(scanning.ts#L171);SQLite 后端则透传给 SQL 的 LIMIT,缺省值为 -1(不限)(search-backend.ts#L173);
  • signal:取消原语,两个内置后端都会在对每个 readable/每一行结果前检查 signal.aborted 并抛出 AbortError(扫描端见 throwIfAborted)。

为什么选择 AsyncIterable

AsyncIterable 让消费方可以做三件事:尽早渲染早期结果、提前停止迭代(拿够就停)、以及用 AbortSignal 取消在途工作。需要强调的边界是:防抖(debounce)始终是 UI/调用方的职责,API 只提供取消原语本身。官方推荐的"边输入边搜(search-as-you-type)"调用范式:

let currentAbortController: AbortController | undefined;

async function updateResults(query: string) {
  currentAbortController?.abort();
  const controller = new AbortController();
  currentAbortController = controller;

  try {
    for await (const hit of search.search(query, { limit: 10, signal: controller.signal })) {
      render(hit);
    }
  } catch (error) {
    if (!(error instanceof Error) || error.name !== "AbortError") throw error;
  }
}

在实现侧,取消的落点很具体:扫描后端在切换每个会话源之前和产出每个候选之前各检查一次信号;SQLite 后端则在打开数据库后、以及结果游标产出的每一行之前检查。取消时若 signal.reason 本身是 Error 则直接抛出它,否则构造一个 name === "AbortError" 的标准错误。测试侧也有对应断言:预 abort 的信号会让迭代以 AbortError 拒绝,见 search.test.ts 中 "honors entry type filters and abort signals in scanning search"

默认实现一:扫描式搜索(Scanning Search)

扫描器(reusable scanner)把"类会话的可读源"(session-like readables)——即具备 getMetadatafindEntriesgetLabel 三个能力的对象——适配为投影条目。相关类型定义在 scanning.ts

export interface SessionSearchCandidate {
  readonly entryId: string;
  readonly seq: number;
  readonly type: Entry["type"];
  readonly timestamp: number;
  readonly text: string;
  readonly fields?: Record<string, unknown>;
}

export interface ScanningSessionSearchHit extends SessionSearchHit {
  readonly timestamp: number;
  readonly snippet: string;
}

SessionSearchCandidate匹配之前的扫描器输入:包含可搜索文本、类型、序号(seq)和可选的投影字段;扫描器把匹配成功的候选转换为对外发布的 hit(扫描后默认 hit 追加 timestampsnippet,见 createDefaultScanningHit)。

直接扫描已打开的会话

const search = createScanningSessionSearch(sessions);

for await (const hit of search.search("authentication", { limit: 10 })) {
  const session = sessionsById.get(hit.sessionId)!;
  const entry = await session.getEntry(hit.entryId);
  console.log(entry);
}

createScanningSessionSearch 的第一个参数既可以是已打开会话的只读数组,也可以是一个 ScanningReadableSource 函数——每次搜索时按需返回 AsyncIterable<ScanningReadable>。这个函数形态正是为"发现/加载"逻辑解耦服务的。

JSONL 无需单独的公共搜索适配器

JSONL 支撑的代码可以保留本地的发现/加载逻辑,再把加载出的 storage 传给同一个扫描器:

async function* jsonlReadables(jsonl: JsonlSessionRepoOptions, query: JsonlSessionListOptions = {}) {
  for (const metadata of await listJsonlSessionMetadata(jsonl, query)) {
    yield loadJsonlSessionStorage(jsonl, metadata);
  }
}

const search = createScanningSessionSearch((query) => jsonlReadables(jsonl, query));

仓库测试中正是这样做的:jsonlReadables 包装 listJsonlSessionMetadata + loadJsonlSessionStorage,在真实磁盘目录上创建 JSONL 会话后搜索 "auth",验证两个会话各命中一条,且 label 文本 "disk label" 也可被搜到,见 search.test.ts 中 "scans JSONL sessions from disk"

写者租约注意事项:扫描源不得对 harness 拥有的会话调用 SessionRepo.open(),如果该操作可能会占用写者租约(writer lease)。JSONL 应使用只读的加载辅助函数;已经打开的会话/storage 可以直接被扫描。

扫描器的源码级行为细节

从源码结构看,扫描器还有几个文档未展开、但直接影响行为的选择:

  • 默认文本投影defaultSearchText 把条目整体 JSON.stringify,若存在 label 再追加在尾部(scanning.ts#L48-L54)。也就是说 label 参与搜索——测试 "includes labels in memory scanning projections" 验证了 setLabel 后的 "important label" 能被 "important" 命中;
  • 默认匹配:大小写不敏感的子串包含(toLowerCase().includesscanning.ts#L111-L113),查询文本进入匹配前会先 trim().toLowerCase()
  • 分页游标scanReadableEntriesfindEntries({ order: "oldestFirst", limit: pageSize, cursor: { afterSeq } }) 循环翻页,pageSize 缺省 100,短页(length < pageSize)即终止(scanning.ts#L56-L89)。注意 pageSize 取的是 query.limit ?? options.pageSize ?? 100——limit 在此影响的是读页大小而非命中截断,命中截断由 hitCount >= limit 单独完成;
  • 可插拔点ScanningSessionSearchOptions 提供 projectText(自定义文本投影)、match(自定义匹配谓词)、createHit(自定义命中构造)、sourceOptions(把归一化后的查询文本与搜索选项透传给源函数,用于服务端下推过滤),见 scanning.ts#L38-L46

默认实现二:SQLite FTS 后端

SQLite 后端位于独立的 @earendil-works/pi-session-backend-sqlite-node 包,实现见 search-backend.ts。它暴露一个扩展 hit:

export interface SqliteSessionSearchHit extends SessionSearchHit {
  readonly metadata: SqliteSessionMetadata;
  readonly timestamp: number;
  readonly score: number;
}
const search = createSqliteSessionSearch({ env, sqlite, databasePath });

for await (const hit of search.search("auth", {
  entryTypes: ["message", "compaction"],
  limit: 20,
})) {
  console.log(hit.sessionId, hit.entryId, hit.score);
}

参数上,SqliteSessionSearchOptions 需要 env(取 absolutePathcreateDir 两个能力,用于解析数据库路径并确保父目录存在)、sqlite(数据库工厂)与 databasePathsearch-backend.ts#L49-L53)。

FTS 表与触发器的懒加载

FTS 表及其触发器在第一次非空白搜索时才创建。首次创建时,SQLite 会执行一次从规范 entries 表到 FTS 索引的全量重建(INSERT INTO session_search_fts(session_search_fts) VALUES('rebuild'));此后,触发器负责让 FTS 与规范条目的插入、删除和 payload 更新保持同步。这使得 SQLite 搜索在提交后即是新鲜的,但代价是:一旦该数据库启用了搜索,FTS 触发器失败会连带回滚规范的 SQLite 写入

从源码看,实际的 DDL 是一个 FTS5 虚拟表加三个触发器(ensureSearchSchema):

CREATE VIRTUAL TABLE IF NOT EXISTS session_search_fts USING fts5(
  payload,
  content = 'entries',
  content_rowid = 'rowid',
  tokenize = 'trigram remove_diacritics 1'
);
CREATE TRIGGER ... AFTER INSERT ON entries ...;        -- 新增同步
CREATE TRIGGER ... AFTER DELETE ON entries ...;        -- 删除同步
CREATE TRIGGER ... AFTER UPDATE OF payload ON entries ...; -- payload 更新 = 先删旧再插新

几个值得注意的实现事实:

  • 分词器是 trigram(三元组),因此支持子串匹配而非整词匹配,remove_diacritics 1 忽略变音符号。对应测试断言了 trigram 匹配行为,见 sqlite-node 的 search.test.ts"matches trigrams");
  • 查询注入防护:用户查询文本被整体双引号包裹且内部引号翻倍("${queryText.replaceAll('"', '""')}"),保证用户输入不会暴露 FTS 查询语法,有专门的回归测试 "handles quoted search text without exposing FTS syntax"
  • 排序ORDER BY bm25(session_search_fts),即 hit 中的 score 来自 SQLite 的 BM25 相关度函数;
  • 连接配置:每次打开搜索数据库都会执行 PRAGMA journal_mode=WALsynchronous=FULLbusy_timeout=5000,随后应用常规迁移(openDatabase),初始化失败时确保 db.close()
  • 边界行为:空白查询、limit <= 0、空 entryTypes 均直接返回空迭代,不触碰数据库——因此"未搜索过的数据库不会被初始化出 FTS",测试 "does not initialize FTS for canonical writes or blank searches" 对此有断言。

SQLite 后端 README 还明确了一个架构约定:仓库(repository)只惰性拥有一个共享数据库连接,而搜索是同一个规范数据库之上的独立服务,repository 本身不暴露 search() 方法。

回滚边界是真实存在的:测试套件中专门验证了 "co-located FTS trigger writes fail 时,规范 append 被回滚" 以及 "FTS cleanup 失败时,规范删除被回滚" 两条用例("rolls back canonical appends when co-located FTS trigger writes fail" 等),这印证了文档中"FTS 触发器失败可回滚规范写入"的警告。

自建索引后端:以 JSONL + Elasticsearch 为例

搜索索引是后端自有的派生状态(backend-owned derived state)。共享包只导出查询 API;应用或后端包在需要显式索引维护时,可以自行定义 writer/feed 契约。

以 JSONL 会话接入 Elasticsearch 为例——这是应用拥有的胶水代码:core 提供查询契约与 JSONL 会话发现,Elastic 的 writer 契约本地化在这个适配器里。完整参考实现如下:

import { Client } from "@elastic/elasticsearch";
import {
  scanningEntries,
  type JsonlSessionMetadata,
  type JsonlSessionRepoOptions,
  type SessionSearch,
  type SessionSearchHit,
  type SessionSearchOptions,
} from "@earendil-works/pi-agent-core";

// JSONL-backed code can provide this locally from existing JSONL list/load helpers.
async function* jsonlReadables(jsonl: JsonlSessionRepoOptions, options: { cwd?: string } = {}) {
  for (const metadata of await listJsonlSessionMetadata(jsonl, options)) {
    yield loadJsonlSessionStorage(jsonl, metadata);
  }
}

interface SearchIndexWriter<TItem> {
  apply(items: TItem[]): Promise<void>;
  flush?(): Promise<void>;
}

interface IndexedSessionSearch<T extends SessionSearchHit, TItem>
  extends SessionSearch<T>, SearchIndexWriter<TItem> {}

type ElasticSessionFeedItem =
  | { type: "upsert"; id: string; body: ElasticSessionDoc }
  | { type: "delete"; id: string };

interface ElasticSessionDoc {
  sessionId: string;
  entryId: string;
  seq: number;
  timestamp: number;
  cwd: string;
  text: string;
  metadata: JsonlSessionMetadata;
  fields?: Record<string, unknown>;
}

interface ElasticSessionSearchHit extends SessionSearchHit {
  readonly timestamp: number;
  readonly snippet: string;
  readonly score?: number;
}

class ElasticSessionSearch
  implements IndexedSessionSearch<ElasticSessionSearchHit, ElasticSessionFeedItem>
{
  constructor(
    private readonly client: Client,
    private readonly index: string,
  ) {}

  async apply(items: ElasticSessionFeedItem[]): Promise<void> {
    const operations = items.flatMap((item) => {
      if (item.type === "delete") {
        return [{ delete: { _index: this.index, _id: item.id } }];
      }
      return [{ index: { _index: this.index, _id: item.id } }, item.body];
    });

    if (operations.length > 0) await this.client.bulk({ operations });
  }

  async flush(): Promise<void> {
    await this.client.indices.refresh({ index: this.index });
  }

  async *search(
    text: string,
    options: SessionSearchOptions = {},
  ): AsyncIterable<ElasticSessionSearchHit> {
    const result = await this.client.search<ElasticSessionDoc>({
      index: this.index,
      size: options.limit ?? 20,
      query: {
        bool: {
          must: [{ match: { text } }],
        },
      },
    });

    for (const hit of result.hits.hits) {
      if (!hit._source) continue;
      if (options.signal?.aborted) throw options.signal.reason;
      yield {
        sessionId: hit._source.sessionId,
        entryId: hit._source.entryId,
        timestamp: hit._source.timestamp,
        snippet: hit._source.text,
        score: hit._score ?? undefined,
      };
    }
  }
}

注意两个设计点:search() 虽然由一次远程调用返回结果,但仍以 AsyncIterable 逐步 yield,并在每个 hit 前检查 options.signal?.aborted,与共享契约的取消语义保持一致;apply/flush 组合成了 SearchIndexWriter 的批量写入契约,flush 通过 indices.refresh 使索引立即可见。

补齐/重建任务(catch-up/rebuild job)可以不获取写者租约,直接把 JSONL 投影喂给 Elasticsearch——这里复用扫描器导出的 scanningEntries 逐候选投影:

async function indexJsonlSessionsIntoElastic(
  jsonl: JsonlSessionRepoOptions,
  elastic: ElasticSessionSearch,
  options: { cwd?: string } = {},
): Promise<void> {
  for await (const session of jsonlReadables(jsonl, { cwd: options.cwd })) {
    const metadata = await session.getMetadata();
    for await (const candidate of scanningEntries(session)) {
      await elastic.apply([{
        type: "upsert",
        id: `${metadata.id}:${candidate.entryId}`,
        body: {
          sessionId: metadata.id,
          entryId: candidate.entryId,
          seq: candidate.seq,
          timestamp: candidate.timestamp,
          cwd: metadata.cwd,
          text: candidate.text,
          metadata,
          fields: candidate.fields,
        },
      }]);
    }
  }

  await elastic.flush();
}

文档 ID 采用 ${sessionId}:${entryId} 拼接,与共享契约的命中身份一一对应,天然幂等——重复执行 rebuild 只会覆盖同一批文档。

正确性与失败边界

最后梳理搜索系统的失败边界,这些约定决定了你的应用该如何容错:

  1. 索引是派生状态,不是权威数据。共享 API 层面,应用可以重试、重建,或将搜索标记为过期(stale)。各后端可以做出不同取舍——SQLite FTS 使用同库触发器,因此 FTS 故障可能回滚规范 SQLite 写入(前提是搜索已初始化了触发器)。
  2. 扫描源必须对重复 sessionId 快速失败。因为基础命中身份是 (sessionId, entryId),同一 sessionId 出现两次会使身份失去唯一性。扫描器确实实现了这一点:每个会话源读取 metadata 后先查 seenSessionIds,重复即抛出 Error("Duplicate sessionId: ...")scanning.ts#L159)。索引型后端则通常在存储/索引层自行强制唯一性(如 Elasticsearch 文档 _id)。
  3. 搜索的按需启用仍需要一层同步/索引设施。文档指出了一个后续改进方向:增加一个默认 no-op 的搜索索引 sink(例如 NOOP_SEARCH_INDEX_SINK),让规范写入点可以无条件地发出索引事件——这与遥测模块在遥测关闭时使用 no-op 实现的做法同构。在这一点落地之前,需要增量索引的应用(如 Elasticsearch 方案)要自行维护 feed 链路。

延伸阅读

登录后查看全文
热门项目推荐
相关项目推荐