ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flue 官方 Postgres 持久化适配器(@flue/postgres)完整实战指南

Flue 官方 Postgres 持久化适配器(@flue/postgres)完整实战指南 Flue 官方 Postgres 持久化适配器flue/postgres完整实战指南【免费下载链接】flueThe sandbox agent framework.项目地址: https://gitcode.com/GitHub_Trending/flue1/flueFlue 的flue/postgres是面向Node.js 目标的第一方 Postgres 持久化适配器它以自带驱动Bring Your Own Driver的方式为 Flue 应用提供基于 PostgreSQL 的持久化能力会话的 canonical 对话流、不可变附件、已接受的 prompt 与dispatch(...)提交记录都可以在进程重启与多副本之间可靠存活。读完本文你将掌握如何在src/db.ts中封装自己的 Postgres 驱动node-postgres 或 porsagerpostgres、理解适配器在启动时自动建表与格式版本校验的机制并能在多副本场景下正确使用这套持久化方案同时明确它不该承担的业务数据边界。一、适配器定位Flue 运行时状态的持久化层flue/postgres在 packages/postgres/README.md 中被明确定义为 Postgres-backed durable persistence for Flue applications on the Node.js target。它不是一个通用的 ORM而是把 Flue 运行时的内部状态落到 Postgres 中具体持久化四类数据每个 agent 实例的canonical 追加式对话流append-only conversation stream这是唯一的会话转录来源并从流的最开始处整体重放对话记录引用的不可变外部附件已接受的直接 prompt与dispatch(...)提交submission带有持久化的 claims认领与 leases租约用于故障恢复与防重复执行workflow-run 记录、事件流与运行索引。从 postgres-adapter.ts 的postgres()工厂函数可以看到一个适配器实例通过connect()一次性返回三件套存储connect() { return { submissionStore: new PgSubmissionStore(runner), conversationStreamStore: createPgConversationStreamStore(runner), attachmentStore: new PgAttachmentStore(runner), }; }分别对应提交存储、对话流存储与附件存储。设计上的取舍与边界README 明确给出几个值得注意的语义约束canonical 流是唯一转录且总是从头重放重放加速replay acceleration与持久化日志压缩persisted-log compaction目前是延后的Session 在其实例生命周期内持续追加没有按 session 删除的接口整实例的流与附件删除方法是底层原语不是面向用户的公开编排能力它不存储你的业务数据。客户记录、工单、支付流水等必须放在你自己的表里。二、自带驱动Bring Your Own Driver设计flue/postgres不挑选、不捆绑任何数据库驱动。它运行在一个你围绕自己配置好的驱动封装的runner之上因此驱动选择、连接池、TLS 以及所有连接选项完全由你掌控。一个 runner 就是三个函数定义见 postgres-adapter.ts 中的PostgresRunner接口函数职责query(text, params)执行带编号$N占位符的 SQL 字符串与位置参数把结果解析为普通对象数组Recordstring, unknown[]transaction(fn)在单条连接上的单个事务内运行fnresolve 时提交、throw 时回滚传给fn的tx只需要queryclose()关闭底层驱动参数类型PostgresParameter限定为string \| number \| boolean \| Uint8Array \| null查询函数签名是PostgresQuery (text: string, params?: PostgresParameter[]) PromiseSqlRow[]。以 porsagerpostgres为例README 给出的第一个例子封装的是postgresporsager驱动db.unsafe(text, params)直接对应 runner 的querydb.begin(...)天然提供单连接事务// src/db.ts import { postgres, type PostgresQuery } from flue/postgres; import sql from postgres; const db sql(process.env.DATABASE_URL!); export default postgres({ query: (text, params) db.unsafe(text, params), transaction: T(fn: (tx: { query: PostgresQuery }) PromiseT) db.begin((tx) fn({ query: (text, params) tx.unsafe(text, params) })) as PromiseT, close: () db.end(), });以 node-postgrespg为例事务必须单连接README 特别强调了一个容易踩的坑如果用pg的Pool连接池无法在任意连接之间跨连接运行事务所以transaction必须自行检出单个 client并亲自发出BEGIN/COMMIT/ROLLBACKimport { postgres } from flue/postgres; import { Pool } from pg; const pool new Pool({ connectionString: process.env.DATABASE_URL }); export default postgres({ query: async (text, params) (await pool.query(text, params)).rows, transaction: async (fn) { const client await pool.connect(); try { await client.query(BEGIN); const result await fn({ query: async (t, p) (await client.query(t, p)).rows }); await client.query(COMMIT); return result; } catch (error) { await client.query(ROLLBACK); throw error; } finally { client.release(); } }, close: () pool.end(), });这段代码同时也是官方 blueprint blueprints/database--postgres.md 中推荐的主选驱动写法pg^8.21.0并配types/pg^8.20.0开发依赖。注意query返回的是.rows数组这正是 runner 期望的形状。三、安装与接入flue add database postgres安装命令是flue add database postgresflue add database postgres会安装flue/postgres包、帮助你挑选驱动并写入db.ts。从 CLI 实现看flue add属于 blueprint 命令体系见 blueprints.ts它会从 blueprint 注册表拉取database/postgres的实现指南来指导接入。接入的核心约定来自 README 与 blueprint文件位置在源码根目录root/.flue/、root/src/或root/中先存在者创建src/db.ts并default-export适配器构建期发现Flue 在构建时发现db.ts并把默认导出接入生成的 Node server自动迁移适配器的migrate()钩子在启动时执行一次幂等地创建所需表不需要单独的迁移步骤也不要为了注册数据库而额外添加app.ts凭据驱动在运行时读取DATABASE_URL不要硬编码连接串也不要发明一个——连接串由环境提供。flue run默认加载项目的.env--env file可指定备选.env格式文件而vite dev与构建后的 server 读取的是 shell 环境变量。包信息速览从 packages/postgres/package.json 可以看到包名为flue/postgresEULA 为 Apache-2.0engines要求Node.js 22.19.0唯一运行时依赖是flue/runtimeworkspace 引用。开发依赖中的electric-sql/pglite说明测试可以在不启动真实 Postgres 服务的情况下用 PGlite 替换 runner 完成契约测试。四、何时使用单机 vs 多副本README 的决策建议非常明确当状态必须在宿主机被替换后存活或需要在多个应用副本之间共享时——例如另一个 Node 进程必须在宿主机故障后恢复已接受的工作或者多个副本需要看到同一份 workflow-run 历史——就使用flue/postgres如果只是单台宿主机内置的、文件后端的sqlite()适配器来自flue/runtime/node就足够了。换句话说flue/postgres的价值在于跨进程、跨宿主的共享与恢复而不是单机下的简单持久化。五、目标支持Node.js 专属flue/postgres只面向 Node.js 目标。Cloudflare 目标会自动使用 Durable Object SQLite并且在构建期拒绝db.ts文件因此数据库适配器在 Cloudflare 目标上不适用。这一点在 blueprint blueprints/database--postgres.md 中被列为先检查目标的第一步如果项目目标是 Cloudflare应直接停下并告知用户没有可添加的内容。六、源码级深度建表、格式版本与存储实现6.1 启动即迁移幂等建表与格式版本门禁postgres()工厂的migrate()调用ensureTablespostgres-adapter.ts。它把全部 schema 搭建放在单个事务里Postgres 的 DDL 是事务性的因此部分失败不会留下半迁移的数据库。整个过程分三层版本打点创建flue_metakey/value表写入format_version。打开数据库时如果发现未知或更新的版本会抛出PersistedFormatVersionError同时兼容 nightly 时代以schema_version8打点的存储存储形状与 format 1 逐字节一致直接原地改标即可任何其他旧值都会拒绝读写。建表幂等创建下述flue_*表。建索引为提交表创建三个关键索引。这与PersistenceAdapter接口的生命周期约定一致——见 agent-execution-store.ts 的接口文档框架在启动时调用一次migrate()如存在把存储带到当前格式版本然后 await 一次connect()获取全部存储数据库不可达或配置错误会在启动时失败而不是在第一个请求里关机时调用close()释放资源。6.2 七张表的结构由ensureTables的 DDL 可以得到完整的表结构flue_metakey TEXT PRIMARY KEY、value TEXT NOT NULL存format_version打点。flue_submission_chunks提交的大对象分块主键(submission_id, item_id, chunk_index)chunk_count记录总分块数data TEXT存块内容。flue_agent_submissions提交主表主键sequence BIGINT GENERATED ALWAYS AS IDENTITYsubmission_id唯一。核心列包括session_key、kinddispatch/direct、payload、status、accepted_at、canonical_ready_at、attempt_id、input_applied_at、abort_requested_at、started_at、joined_into、settled_at、error、attempt_count、max_attempts、timeout_at、owner_id、lease_expires_at、settlement_record_id、settlement_record。三个索引分别覆盖(status, sequence)、(session_key, status, sequence)以及joined_into的部分索引WHERE joined_into IS NOT NULL。flue_conversation_streams每个对话流的元信息主键path含identity_json、next_offset、producer_id、producer_epoch、next_producer_sequence、incarnation。flue_conversation_stream_batches流批次数据主键(path, seq)并有UNIQUE (path, producer_id, producer_epoch, producer_sequence)保证生产者幂等data TEXT存批次内容。flue_conversation_fold_checkpoints流折叠检查点主键path含head_offset、incarnation、format_version、data。flue_attachments附件二进制主键(stream_path, attachment_id)bytes BYTEA存二进制byte_size BIGINT NOT NULL CHECK (byte_size 0)并记录mime_type、digest、conversation_id、created_at。6.3 对话流存储复用 SQL 通用实现postgres-conversation-store.ts 只做了十几行配置调用defineSqlConversationStreamStore来自flue/runtime/adapter注入 Postgres 方言的细节——占位符用$Nplaceholder: (index) \$${index}、行锁用FOR UPDATE、插入忽略用ON CONFLICT (path) DO NOTHING、支持RETURNING。也就是说对话流存储的核心逻辑追加、重放、并发控制由 runtime 的 SQL 通用层实现Postgres 适配器只负责方言适配。6.4 附件存储校验 冲突检测postgres-attachment-store.ts 的put在插入前先verifyAttachmentBytes校验字节与摘要插入用ON CONFLICT ... DO NOTHING随后回读并对比matchesInput不一致则抛AttachmentConflictError保证同一(stream_path, attachment_id)下的内容不可变且幂等。get会按conversationId过滤、再次校验字节后才返回。6.5 提交存储CAS、租约与 join 机制PgSubmissionStorepostgres-adapter.ts是这套适配器最重的部分几个值得展开的实现事实认领的并发安全claimSubmission用 CTE 先锁定候选行再按主键更新并在外层WHERE再次校验status queued。源码注释解释了原因Postgres 不支持 SQLite 那种带自引用NOT EXISTS子查询的UPDATE ... AS alias写法且 READ COMMITTED 隔离级别下必须在 CTE 快照与 UPDATE 之间防止并发认领重复领取同一行。会话内顺序保证listRunnableSubmissions只返回每个session_key下最靠前的可运行提交用NOT EXISTS排除前面仍有queued/running/terminalizing/joining/joined的更早行从而保证同一会话的提交按接收顺序执行。租约与超时renewLeases按owner_id批量续约lease_expires_atlistExpiredSubmissions找出租约过期仍为running的提交交给上层重新认领这是宿主机故障后由另一进程恢复已接受工作的关键支撑。join 交付扇出claimJoinableSubmissions用FOR UPDATE锁住宿主行与排队行按 contiguous prefix 规则把同一会话、同一 agent 的排队提交以joining状态并入宿主 attempt宿主 settle 时settleJoinedSubmissions在同一个事务里把joined行以宿主结局结清、把joining残留行abort 或崩溃窗口回退为queued以免任务凭空消失。格式健壮性解析行时会做严格校验如 Postgres 把BIGINT作为字符串返回需Number()强转queued行不得带运行期字段、running/terminalizing行必须有attemptId与startedAt等发现畸形行会将其以错误结算而不是阻塞整个队列parseOperationalRows/failSubmissionSequence。七、验证与运行blueprint blueprints/database--postgres.md 给出了完整的验收步骤可以直接照做类型检查npx tsc --noEmit安全。构建为 Node 目标执行vite build确认适配器被发现并接入生成的 server。启动验证把DATABASE_URL指向可达的 Postgres本地容器即可启动 server 确认能正常 boot——首次运行migrate()会创建全部flue_*表重启后确认既有状态是被重新加载而非重建。红线不要用生产数据库来做测试。总结flue/postgres以极简的 runner 契约把 Postgres 接入 Flue 的持久化体系三函数 runner 让你完全掌控驱动与连接migrate()启动即建表、无需额外迁移步骤flue_meta格式版本门禁与幂等建表保证了存储的前向安全。它持久化的是 Flue 运行时状态——canonical 对话流、不可变附件、带 claims/leases 的提交与 workflow 记录而业务数据请留在你自己的表中。当你的 Flue 应用需要在多副本之间共享状态、或必须从宿主机故障中恢复未完成工作时它就是 Node 目标下的首选接入方式单机场景则优先考虑flue/runtime/node内置的文件后端sqlite()适配器。【免费下载链接】flueThe sandbox agent framework.项目地址: https://gitcode.com/GitHub_Trending/flue1/flue创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表