ARTICLE DETAIL

资讯详情

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

Flue Sandbox Agent 框架的 MongoDB 持久化接入指南:@flue/mongodb 包深度解析

Flue Sandbox Agent 框架的 MongoDB 持久化接入指南:@flue/mongodb 包深度解析 Flue Sandbox Agent 框架的 MongoDB 持久化接入指南flue/mongodb 包深度解析【免费下载链接】flueThe sandbox agent framework.项目地址: https://gitcode.com/GitHub_Trending/flue1/flueflue/mongodb是 Flue 项目The sandbox agent framework中面向 Node 目标运行时Node-target projects的 MongoDB 持久化适配层。它不直接内置 MongoDB 驱动依赖而是通过一个由应用自持的MongoRunner完成全部存储读写本文将从安装、Runner 契约、事务重试语义、完整接入示例、Schema 迁移到存储模型系统讲解如何为 Flue Agent 接入 MongoDB 持久化并深入对应源码印证底层实现。适用前提为什么独立 MongoDB 实例会被直接拒绝flue/mongodb对部署形态有硬性要求必须是副本集replica set、Atlas 部署或支持事务的分片集群transaction-capable sharded cluster独立运行的 standalone MongoDB 会在写入 Schema 版本戳schema stamping之前就被拒绝。这一判断发生在迁移阶段。在 mongodb-adapter.ts 中migrate()会先调用 runner 的topology()探测拓扑类型const topology await runner.topology(); if (topology.kind standalone || !topology.transactions) throw new TypeError( flue/mongodb requires a replica set, Atlas, or a transaction-capable sharded cluster., );而topology()的实现README 示例代码中给出通过db.admin().command({ hello: 1 })读取服务器响应setName存在即为副本集msg isdbgrid即为分片集群否则为 standalone同时用logicalSessionTimeoutMinutes判断事务能力。原因在于 Flue 的持久化语义严重依赖 MongoDB 多文档事务跨集合原子写入、状态机转移standalone 无法满足这一约束。安装与依赖策略驱动由应用自持pnpm add flue/mongodb mongodb关键设计是该包没有生产环境驱动依赖no production driver dependency。查看 package.json其dependencies只有flue/runtime与ulidxmongodb驱动仅作为devDependencies用于开发测试。这意味着MongoDB 官方驱动的版本、连接池、TLS、认证、备份等配置完全由应用拥有的MongoClient决定flue/mongodb通过鸭子类型structural typing的MongoRunner接口消费驱动能力驱动升级不影响包本身包的导出面在 index.ts 中收敛为mongodb()工厂函数、runMongoTransactionWithRetry帮助函数以及MongoCollection、MongoOperations、MongoRunner等类型。MongoRunner 契约适配层与驱动的边界Runner 是flue/mongodb的核心抽象它必须满足以下运行约束README 原文要求使用snapshot 读关注snapshot read concern使用majority 写关注majority write concern每个回调操作callback operation绑定同一个ClientSession回调内操作严格串行sequential对TransientTransactionError做有界的整事务重试bounded whole-transaction retries对UnknownTransactionCommitResult做仅提交阶段重试commit-only retries。这一契约在源码中被完整类型化为 mongodb-runner.ts 的MongoRunner接口export interface MongoRunner extends MongoOperations { transactionT(fn: (tx: MongoOperations) PromiseT): PromiseT; topology(): PromiseMongoTopology; ensureCollection(spec: MongoCollectionSpec): Promisevoid; inspectCollection(name: string): Promise{ ... } | null; close(): void | Promisevoid; }同时暴露的还有MongoOperationscollection(name)工厂、MongoCollectionfindOne/find/insertOne/insertMany/updateOne/updateMany/findOneAndUpdate/deleteOne/deleteMany八类操作、MongoIndexSpec、MongoCollectionSpec含validator、validationLevel: strict、validationAction: error等类型。其中MongoTopology.kind的取值范围为replica_set | sharded | standalone | unknown。事务重试语义的官方参考实现README 要求“整事务重试仅针对TransientTransactionError提交重试仅针对UnknownTransactionCommitResult且两个循环都有上界”。包内提供了可直接复用的官方实现 runMongoTransactionWithRetry整事务最多尝试maxTransactionAttempts默认5次提交最多尝试maxCommitAttempts默认10次提交失败仅当错误带UnknownTransactionCommitResult标签时才重试且到达上限即抛出事务体失败仅当错误带TransientTransactionError标签时才重新开事务否则抛出每次失败都会先abort()再end()会话abort失败被静默吞掉.catch(() undefined)。这与 README 示例中手写 runner 的循环逻辑5 次事务尝试、10 次提交尝试完全一致两种方式都符合契约。完整接入示例应用侧 Runner 实现以下是 README 提供的完整接入代码是接入flue/mongodb的标准骨架import { mongodb, type MongoCollection, type MongoOperations, type MongoRunner, } from flue/mongodb; import { MongoClient } from mongodb; const client new MongoClient(process.env.MONGODB_URL!); await client.connect(); const db client.db(process.env.MONGODB_DATABASE); const operations (session?: import(mongodb).ClientSession): MongoOperations { let pending Promise.resolve(); const queue T(operation: () PromiseT): PromiseT { const next pending.then(operation, operation); pending next.then( () undefined, () undefined, ); return next; }; return { collection(name): MongoCollection { const collection db.collection(name); const options session ? { session } : {}; return { findOne: (filter, opts) queue(() collection.findOne(filter, { ...opts, ...options })), find: (filter {}, opts {}) queue(() collection.find(filter, { ...opts, ...options }).toArray()), insertOne: (document) queue(() collection.insertOne(document, options)), insertMany: (documents) queue(() collection.insertMany(documents, options)), updateOne: (filter, update, opts) queue(() collection.updateOne(filter, update, { ...opts, ...options })), updateMany: (filter, update) queue(() collection.updateMany(filter, update, options)), findOneAndUpdate: (filter, update, opts) queue(() collection.findOneAndUpdate(filter, update, { ...opts, ...options })), deleteOne: (filter) queue(() collection.deleteOne(filter, options)), deleteMany: (filter) queue(() collection.deleteMany(filter, options)), } as MongoCollection; }, }; }; const runner: MongoRunner { ...operations(), async transaction(fn) { for (let attempt 0; attempt 5; attempt) { const session client.startSession(); try { session.startTransaction({ readConcern: { level: snapshot }, writeConcern: { w: majority }, }); const result await fn(operations(session)); for (let commitAttempt 0; ; commitAttempt) { try { await session.commitTransaction(); break; } catch (error) { const retryable error instanceof Error hasErrorLabel in error typeof error.hasErrorLabel function error.hasErrorLabel(UnknownTransactionCommitResult); if (!retryable || commitAttempt 9) throw error; } } return result; } catch (error) { await session.abortTransaction().catch(() undefined); if ( !(error instanceof Error) || !(hasErrorLabel in error) || !(error as any).hasErrorLabel(TransientTransactionError) || attempt 4 ) throw error; } finally { await session.endSession(); } } throw new Error(unreachable); }, async topology() { const hello await db.admin().command({ hello: 1 }); const kind hello.setName ? replica_set : hello.msg isdbgrid ? sharded : standalone; return { kind, transactions: kind ! standalone hello.logicalSessionTimeoutMinutes ! null, }; }, async ensureCollection(spec) { const existing await db.listCollections({ name: spec.name }).hasNext(); if (!existing) { try { await db.createCollection(spec.name, { validator: spec.validator, validationLevel: spec.validationLevel, validationAction: spec.validationAction, }); } catch (error) { if ( !(error instanceof Error) || !(codeName in error) || error.codeName ! NamespaceExists ) throw error; } } await db.command({ collMod: spec.name, validator: spec.validator, validationLevel: spec.validationLevel, validationAction: spec.validationAction, }); for (const index of spec.indexes) await db.collection(spec.name).createIndex(index.key, index); }, async inspectCollection(name) { const info await db.listCollections({ name }).next(); if (!info) return null; const indexes (await db.collection(name).listIndexes().toArray()) .filter((index) index.name ! _id_) .map((index) ({ name: String(index.name), key: index.key as Recordstring, 1 | -1, ...(index.unique true ? { unique: true } : {}), ...(index.partialFilterExpression ? { partialFilterExpression: index.partialFilterExpression } : {}), ...(index.collation ? { collation: index.collation } : {}), })); return { validator: info.options.validator, validationLevel: info.options.validationLevel, validationAction: info.options.validationAction, indexes, }; }, close: () client.close(), }; export default mongodb(runner);该示例中值得注意的工程细节queue串行化所有集合操作通过一个pendingPromise 链排队执行确保同一事务回调内不会出现并发交叉这是 README 明确要求“callback operations 严格串行”的落地方式ensureCollection的幂等性先查listCollections不存在才createCollection对NamespaceExists错误并发创建静默容忍随后统一用collMod固化 validator/validationLevel/validationAction再逐个createIndex建索引inspectCollection用于校验返回当前 validator、校验级别与全部非_id_索引含 unique、partialFilterExpression、collation 元信息供迁移阶段做精确比对。迁移与连接先 migrate() 再 connect()接入的关键纪律是必须调用并await migrate()之后才能调用connect()。在 mongodb-adapter.ts 中migrate()承担了完整的初始化职责而connect()在migrated标志为 false 时直接抛出TypeError(flue/mongodb connect() requires a successful migrate() first.)见 mongodb-adapter.ts。migrate()的完整流程包括预检读取meta集合中的format_version存在旧式schema_version时检查是否可收养adoptable legacy version仅8可无迁移直接收养因为它与 format 1 的存储形态逐字节一致见 mongodb-adapter.ts既无版本戳又有未版本化数据时拒绝PersistedFormatVersionErrorstoredVersion 为unversioned拓扑检查如前述拒绝 standalone / 无事务拓扑迁移锁在meta集合上以migration_lock文档实现跨进程互斥lease 为 30 秒MIGRATION_LEASE_MS 30_000每MIGRATION_LEASE_MS / 3心跳续租一次锁丢失或租约过期则抛错finally中删除锁并等待心跳队列排空——保证多实例并发启动时只有一位执行迁移版本盖章新版本戳format_version先于旧schema_version删除写入中断也不会出现无戳状态Schema 校验调用ensureSchema逐集合核对 validator 与索引见下文收尾写入/复核format_version并执行一次collectGarbage()清理过期暂存值最后置migrated true。Schema 校验的“零宽容”策略ensureSchemaschema.ts对集合实施精确比对validator 采用规范化 JSON 序列化比对canonical()对键排序、剔除 undefinedvalidationLevel 必须为strict、validationAction 必须为error不一致即抛TypeError(... has incompatible validation options.)每个期望索引按名字在现有索引中查找用comparableIndex比对 key、unique、partialFilterExpression、collation。一个值得注意的细节是simple collationlocale: simple与无 collation 视为等价因为 MongoDB 对 simple 二进制排序索引持久化时会省略 collation 字段见 schema.ts 的注释。Schema v6 是重置边界README 明确指出Schema v6 是一个重置边界reset boundary。使用其他 schema 版本创建的存储stores在迁移前必须清空clear否则migrate()会以版本不兼容为由拒绝。这也解释了assertAdoptableLegacyVersion只放行schema_version: 8的原因——其余旧版本一概视为不可读存储。集合清单与索引设计schema(prefix)schema.ts定义了完整集合默认前缀flue_见mongodb-adapter.ts中options.collectionPrefix ?? flue_集合索引说明meta—格式版本戳、schema 版本戳、迁移锁counters—submission 序号计数器$inc自增value_generationsstate_created{ state: 1, createdAt: 1 }任意值分代登记表staged/published 状态机valuesgeneration_index{ generation: 1, index: 1 }unique分代值的分片数据体submissionssubmission_idunique, simple collation、status_sequence{ status: 1, sequence: 1 }、session_status_sequence{ sessionKey: 1, status: 1, sequence: 1 }simple collation、joined_into{ joinedInto: 1 }simple collationAgent 提交submission状态机与队列conversation_streams—会话流元数据identity、producer、offsetconversation_fold_checkpoints—折叠检查点缓存conversation_batchespath_offset{ path: 1, offset: 1 }unique, simple collation、producer_sequence{ path: 1, producerId: 1, producerEpoch: 1, producerSequence: 1 }unique, simple collation规范追加日志批次attachmentspath_attachment{ path: 1, attachmentId: 1 }unique, simple collation不可变外部附件存储模型Flue 在 MongoDB 上的数据语义README 对存储模型给出了精确定义结合源码可以逐条印证规范追加型会话流canonical append-only conversation streamsconversation_streamsconversation_batches构成唯一的会话记录sole transcript并且从头重放replayed from its beginningREADME 明确说明“replay acceleration 与 persisted-log compaction 被推迟deferred”即当前版本不做重放加速与日志压缩。会话在实例生命周期内持续追加没有按会话删除no per-session deletion。追加路径由 conversation-store.ts 的append()实现全部在事务内完成批次数据以JSON.stringify序列化后整段存入data字段受12MB 单批次上限约束MAX_BATCH_DATA_LENGTH 12 * 1024 * 1024见 conversation-store.ts。源码注释解释了该数值的依据2MB 的单条记录已接近 20 万 token超过 12MB 的批次意味着无界工具结果而非真实对话数据内置工具将结果截断在 50KB通过producerId/producerEpoch/incarnation/nextProducerSequence实现生产者围栏producer fencingacquireProducer以$inc: { producerEpoch: 1 }递增纪元并归还 incarnationappend校验四元组一致、sequence 为下一个期望值且最终条件更新modifiedCount ! 1即判并发篡夺同一 producer sequence 重放时校验数据与 submission 引用完全一致冲突则报错、一致则幂等返回{ offset, appended: false }读取支持游标offset支持-1从头、now游标到头部与数值偏移limit 被clampLimit限制在默认值与最大值之间。不可变外部附件immutable external attachmentsattachments集合存储附件字节与元数据mimeType、byteSize、digest、conversationId实现见 attachment-store.ts。put()在事务内先查重同path attachmentId已存在且引用与字节完全一致则幂等返回不一致则抛AttachmentConflictError并发下依赖唯一索引path_attachment捕获 11000 重复键错误后二次确认。字节以 BSON Binary 持久化读取时校验摘要verifyAttachmentBytes并拷贝返回。持久化提交durable submissions、认领与租约claims and leasessubmissions集合承载 Agent 提交的完整生命周期状态机queued → running → terminalizing → settled以及 turn-boundary 的joining/joined。从 submission-store.ts 可以看到认领claimSubmission在事务内先确保自己是同 sessionKey 队列中 sequence 最靠前的候选再用聚合管道$set原子完成queued → running转移、写入attemptId/ownerId/leaseExpiresAt、递增attemptCount并设置默认maxAttempts与timeoutAt租约续期与过期renewLeases批量延长运行中提交的租约listExpiredSubmissions找出leaseExpiresAt已过期的运行中提交以回收join 扇出settleJoinedSubmissions在宿主提交结算事务内把joined行以宿主结果结算把joining遗留行回退为queued避免交付凭空消失坏行隔离parseOperationalRows遇到无法解析的行时将该行强制转为settled并附带错误信息保证单行损坏不会卡死整个队列。任意值分代暂存value stagingFlue 需要持久化任意结构的值如 submission 的 payload 与附件块。value-store.ts 实现了一个两阶段方案暂存stage把JSON.stringify([value])切分为每段最大4MBPART_BYTES 4 * 1024 * 1024按字节边界回退切割避免截断多字节字符的不可变分片写入values集合并在value_generations登记state: staged与createdAt发布publish在一个短事务short transaction内把登记表的状态改为published并同时发布其 manifest——这正是 README 所说“任意值以有界不可变分片暂存持久化代际状态就绪后由短事务发布清单”的机制垃圾回收collectGarbage清理超过1 小时STAGED_MAX_AGE_MS 60 * 60 * 1000仍未发布的暂存代每批 100 条GC_BATCH由state_created索引支撑且在每次migrate()收尾时也会执行一次。读取时严格校验分片数量与 index 连续性不完整则抛TypeError(Persisted MongoDB value generation is incomplete.)。关于“删除”的定位README 特别提醒整实例的流删除与附件删除方法whole-instance stream and attachment deletion methods是底层原语low-level primitives不是公开的业务编排接口not public orchestration。这意味着应用不应把它们当作日常清理手段而应视作底层维护工具。多租户与运维建议README 最后给出两条运维建议尽可能使用独立数据库dedicated database否则必须通过options.collectionPrefix设置唯一命名空间。工厂函数签名mongodb(runner, options)中collectionPrefix默认flue_见 mongodb-adapter.ts所有集合名通过collectionName(prefix, name)拼接schema.ts因此共享库内多租户互不干扰凭据、TLS、连接池、备份与客户端生命周期credentials, TLS, pooling, backups, client lifecycle全部在应用拥有的驱动application-owned driver中配置——flue/mongodb只通过MongoRunner.close()委托关闭连接不接管任何连接管理。小结flue/mongodb以“无驱动依赖 应用自持 Runner”的轻量设计为 Flue Node 目标项目提供了符合强事务语义的 MongoDB 持久化严格的拓扑门槛、精确的 Schema 校验、有界的双标签事务重试、追加型会话流与分代暂存机制共同保证了 Agent 会话、提交与附件的持久可靠性。接入时只需遵循两条主线——实现一个满足契约的MongoRunner并始终先await migrate()再connect()运维上优先独立数据库共享库则务必设置唯一的collectionPrefix。如需进一步深入可直接阅读 mongodb-runner.tsRunner 契约与重试实现、schema.ts集合与索引定义、mongodb-adapter.ts迁移与连接流程以及 conversation-store.ts、submission-store.ts、attachment-store.ts、value-store.ts 四个存储实现。【免费下载链接】flueThe sandbox agent framework.项目地址: https://gitcode.com/GitHub_Trending/flue1/flue创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表