{"version":3,"sources":["../src/index.ts","../src/mongo-lock.ts"],"sourcesContent":["/**\n * awaitly-mongo\n *\n * MongoDB persistence adapter for awaitly workflows.\n * Provides ready-to-use SnapshotStore backed by MongoDB.\n * Supports both WorkflowSnapshot and ResumeState (serialized via serializeResumeState).\n */\n\nimport type { Db, MongoClientOptions } from \"mongodb\";\nimport { MongoClient as MongoClientImpl } from \"mongodb\";\nimport type { WorkflowSnapshot, SnapshotStore } from \"awaitly/durable\";\nimport type { WorkflowLock } from \"awaitly/durable\";\nimport {\n  type ResumeState,\n  type StoreSaveInput,\n  type StoreLoadResult,\n  isWorkflowSnapshot,\n  isResumeState,\n  isSerializedResumeState,\n  serializeResumeState,\n  deserializeResumeState,\n} from \"awaitly/durable\";\nimport { createMongoLock, type MongoLockOptions } from \"./mongo-lock\";\n\n/** Document shape for the snapshots collection (string _id). */\ninterface SnapshotDoc {\n  _id: string;\n  snapshot: WorkflowSnapshot | import(\"awaitly/durable\").SerializedResumeState;\n  updatedAt: Date;\n}\n\nexport type { SnapshotStore, WorkflowSnapshot } from \"awaitly/durable\";\nexport type { WorkflowLock } from \"awaitly/durable\";\nexport type { MongoLockOptions } from \"./mongo-lock\";\nexport type { StoreSaveInput, StoreLoadResult } from \"awaitly/durable\";\n\n// =============================================================================\n// MongoOptions\n// =============================================================================\n\n/**\n * Options for the mongo() shorthand function.\n */\nexport interface MongoOptions {\n  /** MongoDB connection URL. */\n  url: string;\n  /** Database name. @default 'awaitly' */\n  database?: string;\n  /** Collection name for snapshots. @default 'awaitly_snapshots' */\n  collection?: string;\n  /** Key prefix for IDs. @default '' */\n  prefix?: string;\n  /** Bring your own client. */\n  client?: MongoClientImpl;\n  /** MongoDB client options. */\n  clientOptions?: MongoClientOptions;\n  /** Cross-process lock options. When set, the store implements WorkflowLock. */\n  lock?: MongoLockOptions;\n}\n\n/** Mongo store with widened save/load for WorkflowSnapshot and ResumeState. Compatible with SnapshotStore for snapshot-only usage. */\nexport interface MongoStore extends Partial<WorkflowLock> {\n  save(id: string, state: StoreSaveInput): Promise<void>;\n  load(id: string): Promise<StoreLoadResult>;\n  loadResumeState(id: string): Promise<ResumeState | null>;\n  delete(id: string): Promise<void>;\n  list(options?: { prefix?: string; limit?: number }): Promise<Array<{ id: string; updatedAt: string }>>;\n  close(): Promise<void>;\n}\n\n// =============================================================================\n// mongo() - One-liner Snapshot Store Setup\n// =============================================================================\n\n/**\n * Create a snapshot store backed by MongoDB.\n * Save accepts WorkflowSnapshot or ResumeState; load returns whichever was stored.\n * Use loadResumeState(id) for type-safe restore, or toResumeState(await store.load(id)).\n *\n * @example\n * ```typescript\n * import { mongo } from 'awaitly-mongo';\n * import { createWorkflow } from 'awaitly';\n *\n * const store = mongo('mongodb://localhost:27017/mydb');\n * const workflow = createWorkflow(deps);\n *\n * // Run and persist resume state\n * const { result, resumeState } = await workflow.runWithState(fn);\n * await store.save('wf-123', resumeState);\n *\n * // Restore\n * const resumeState = await store.loadResumeState('wf-123');\n * if (resumeState) await workflow.run(fn, { resumeState });\n * ```\n *\n * @example\n * ```typescript\n * // With options including cross-process locking\n * const store = mongo({\n *   url: 'mongodb://localhost:27017',\n *   database: 'myapp',\n *   collection: 'my_workflow_snapshots',\n *   prefix: 'orders:',\n *   lock: { lockCollectionName: 'my_workflow_locks' },\n * });\n * ```\n */\nexport function mongo(urlOrOptions: string | MongoOptions): MongoStore {\n  const opts = typeof urlOrOptions === \"string\" ? { url: urlOrOptions } : urlOrOptions;\n  const prefix = opts.prefix ?? \"\";\n\n  // Parse database from URL if provided\n  let databaseName = opts.database;\n  const urlMatch = opts.url.match(/mongodb(?:\\+srv)?:\\/\\/[^/]+\\/([^?]+)/);\n  if (!databaseName && urlMatch && urlMatch[1]) {\n    databaseName = urlMatch[1];\n  }\n  databaseName = databaseName ?? \"awaitly\";\n\n  const collectionName = opts.collection ?? \"awaitly_snapshots\";\n\n  const ownClient = !opts.client;\n  let client: MongoClientImpl | undefined = opts.client;\n  let db: Db | undefined;\n  let connected = false;\n  let lock: { tryAcquire: WorkflowLock[\"tryAcquire\"]; release: WorkflowLock[\"release\"]; renew: NonNullable<WorkflowLock[\"renew\"]> } | null = null;\n\n  const ensureConnected = async (): Promise<Db> => {\n    if (db && connected) return db;\n\n    if (!client) {\n      client = new MongoClientImpl(opts.url, {\n        directConnection: !opts.url.includes(\"mongodb+srv://\"),\n        ...opts.clientOptions,\n      });\n    }\n\n    await client.connect();\n    connected = true;\n    db = client.db(databaseName);\n\n    // Create index on updatedAt for list queries\n    const collection = db.collection<SnapshotDoc>(collectionName);\n    await collection.createIndex({ updatedAt: -1 }, { background: true }).catch(() => {\n      // Index may already exist, ignore error\n    });\n\n    if (opts.lock && !lock) {\n      lock = createMongoLock(db, opts.lock);\n    }\n\n    return db;\n  };\n\n  const store: MongoStore = {\n    async save(id: string, state: StoreSaveInput): Promise<void> {\n      const db = await ensureConnected();\n      const collection = db.collection<SnapshotDoc>(collectionName);\n      const fullId = prefix + id;\n      const toStore = isResumeState(state) ? serializeResumeState(state) : state;\n      await collection.updateOne(\n        { _id: fullId },\n        {\n          $set: {\n            snapshot: toStore,\n            updatedAt: new Date(),\n          },\n        },\n        { upsert: true }\n      );\n    },\n\n    async load(id: string): Promise<StoreLoadResult> {\n      const db = await ensureConnected();\n      const collection = db.collection<SnapshotDoc>(collectionName);\n      const fullId = prefix + id;\n      const doc = await collection.findOne({ _id: fullId });\n      if (!doc) return null;\n      const raw = doc.snapshot;\n      if (isSerializedResumeState(raw)) return deserializeResumeState(raw);\n      if (isWorkflowSnapshot(raw)) return raw;\n      return raw as WorkflowSnapshot;\n    },\n\n    async loadResumeState(id: string): Promise<ResumeState | null> {\n      const loaded = await store.load(id);\n      if (loaded === null) return null;\n      if (isResumeState(loaded)) return loaded;\n      return null;\n    },\n\n    async delete(id: string): Promise<void> {\n      const db = await ensureConnected();\n      const collection = db.collection<SnapshotDoc>(collectionName);\n      const fullId = prefix + id;\n      await collection.deleteOne({ _id: fullId });\n    },\n\n    async list(options?: { prefix?: string; limit?: number }): Promise<Array<{ id: string; updatedAt: string }>> {\n      const db = await ensureConnected();\n      const collection = db.collection<SnapshotDoc>(collectionName);\n      const filterPrefix = prefix + (options?.prefix ?? \"\");\n      const limit = options?.limit ?? 100;\n      const escaped = filterPrefix.replace(/[.*+?^${}()|[\\]\\\\]/g, \"\\\\$&\");\n\n      const cursor = collection\n        .find({ _id: { $regex: `^${escaped}` } })\n        .sort({ updatedAt: -1 })\n        .limit(limit);\n\n      const docs = await cursor.toArray();\n      return docs.map(doc => ({\n        id: String(doc._id).slice(prefix.length),\n        updatedAt: doc.updatedAt.toISOString(),\n      }));\n    },\n\n    async close(): Promise<void> {\n      if (ownClient && client) {\n        await client.close();\n        connected = false;\n        db = undefined;\n      }\n    },\n\n    // Lock methods are added dynamically below after first connection\n    async tryAcquire(id: string, options?: { ttlMs?: number }): Promise<{ ownerToken: string } | null> {\n      await ensureConnected();\n      if (!lock) return null;\n      return lock.tryAcquire(id, options);\n    },\n\n    async release(id: string, ownerToken: string): Promise<void> {\n      await ensureConnected();\n      if (!lock) return;\n      return lock.release(id, ownerToken);\n    },\n\n    async renew(id: string, ownerToken: string, opts?: { ttlMs?: number }): Promise<boolean> {\n      await ensureConnected();\n      if (!lock) return false;\n      return lock.renew(id, ownerToken, opts);\n    },\n  };\n\n  // Only include lock methods if lock is configured\n  if (!opts.lock) {\n    delete store.tryAcquire;\n    delete store.release;\n    delete store.renew;\n  }\n\n  return store;\n}\n","/**\n * MongoDB workflow lock (lease) for cross-process concurrency control.\n * Uses a lease (TTL) + owner token; release verifies the token.\n */\n\nimport type { Db, Collection } from \"mongodb\";\nimport { randomUUID } from \"node:crypto\";\n\nexport interface MongoLockOptions {\n  /**\n   * Collection name for workflow locks.\n   * @default 'workflow_lock'\n   */\n  lockCollectionName?: string;\n}\n\ninterface LockDocument {\n  _id: string;\n  ownerToken: string;\n  expiresAt: Date;\n}\n\n/**\n * Create tryAcquire and release functions that use a MongoDB lock collection.\n * Caller must pass the same Db used for state (so one connection).\n */\nexport function createMongoLock(\n  db: Db,\n  options: MongoLockOptions = {}\n): {\n  tryAcquire(\n    id: string,\n    opts?: { ttlMs?: number }\n  ): Promise<{ ownerToken: string } | null>;\n  release(id: string, ownerToken: string): Promise<void>;\n  renew(id: string, ownerToken: string, opts?: { ttlMs?: number }): Promise<boolean>;\n  ensureLockCollection(): Promise<void>;\n} {\n  const lockCollectionName = options.lockCollectionName ?? \"workflow_lock\";\n  const collection = db.collection<LockDocument>(lockCollectionName);\n\n  async function ensureLockCollection(): Promise<void> {\n    const collections = await db.listCollections({ name: lockCollectionName }).toArray();\n    if (collections.length === 0) {\n      await db.createCollection(lockCollectionName);\n    }\n  }\n\n  async function tryAcquire(\n    id: string,\n    opts?: { ttlMs?: number }\n  ): Promise<{ ownerToken: string } | null> {\n    const ttlMs = opts?.ttlMs ?? 60_000;\n    const ownerToken = randomUUID();\n    const expiresAt = new Date(Date.now() + ttlMs);\n\n    await ensureLockCollection();\n\n    // Atomic: insert or update only when no doc or doc is expired.\n    // When lock exists and is unexpired, filter won't match and upsert\n    // throws duplicate key error (E11000) - catch it and return null.\n    try {\n      const result = await collection.findOneAndUpdate(\n        {\n          _id: id,\n          $or: [\n            { expiresAt: { $lt: new Date() } },\n            { expiresAt: { $exists: false } },\n          ],\n        },\n        { $set: { ownerToken, expiresAt } },\n        { upsert: true, returnDocument: \"after\" }\n      );\n\n      if (result && result.ownerToken === ownerToken) {\n        return { ownerToken };\n      }\n      return null;\n    } catch (error: unknown) {\n      // Duplicate key error means lock exists and is not expired\n      if (\n        error &&\n        typeof error === \"object\" &&\n        \"code\" in error &&\n        error.code === 11000\n      ) {\n        return null;\n      }\n      throw error;\n    }\n  }\n\n  async function release(id: string, ownerToken: string): Promise<void> {\n    await collection.deleteOne({ _id: id, ownerToken });\n  }\n\n  async function renew(\n    id: string,\n    ownerToken: string,\n    opts?: { ttlMs?: number }\n  ): Promise<boolean> {\n    const ttlMs = opts?.ttlMs ?? 60_000;\n    const expiresAt = new Date(Date.now() + ttlMs);\n\n    const result = await collection.updateOne(\n      { _id: id, ownerToken },\n      { $set: { expiresAt } }\n    );\n\n    // matchedCount, not modifiedCount: two renewals in the same millisecond\n    // compute an identical expiresAt, and Mongo reports an unchanged document\n    // as not modified. Reading that as \"lease lost\" made the durable heartbeat\n    // abort a workflow that still owned its lease. Matching on the owner token\n    // already proves ownership, which is what renew answers.\n    return result.matchedCount > 0;\n  }\n\n  return { tryAcquire, release, renew, ensureLockCollection };\n}\n"],"mappings":"yaAAA,IAAAA,EAAA,GAAAC,EAAAD,EAAA,WAAAE,IAAA,eAAAC,EAAAH,GASA,IAAAI,EAA+C,mBAG/CC,EASO,2BCfP,IAAAC,EAA2B,kBAoBpB,SAASC,EACdC,EACAC,EAA4B,CAAC,EAS7B,CACA,IAAMC,EAAqBD,EAAQ,oBAAsB,gBACnDE,EAAaH,EAAG,WAAyBE,CAAkB,EAEjE,eAAeE,GAAsC,EAC/B,MAAMJ,EAAG,gBAAgB,CAAE,KAAME,CAAmB,CAAC,EAAE,QAAQ,GACnE,SAAW,GACzB,MAAMF,EAAG,iBAAiBE,CAAkB,CAEhD,CAEA,eAAeG,EACbC,EACAC,EACwC,CACxC,IAAMC,EAAQD,GAAM,OAAS,IACvBE,KAAa,cAAW,EACxBC,EAAY,IAAI,KAAK,KAAK,IAAI,EAAIF,CAAK,EAE7C,MAAMJ,EAAqB,EAK3B,GAAI,CACF,IAAMO,EAAS,MAAMR,EAAW,iBAC9B,CACE,IAAKG,EACL,IAAK,CACH,CAAE,UAAW,CAAE,IAAK,IAAI,IAAO,CAAE,EACjC,CAAE,UAAW,CAAE,QAAS,EAAM,CAAE,CAClC,CACF,EACA,CAAE,KAAM,CAAE,WAAAG,EAAY,UAAAC,CAAU,CAAE,EAClC,CAAE,OAAQ,GAAM,eAAgB,OAAQ,CAC1C,EAEA,OAAIC,GAAUA,EAAO,aAAeF,EAC3B,CAAE,WAAAA,CAAW,EAEf,IACT,OAASG,EAAgB,CAEvB,GACEA,GACA,OAAOA,GAAU,UACjB,SAAUA,GACVA,EAAM,OAAS,KAEf,OAAO,KAET,MAAMA,CACR,CACF,CAEA,eAAeC,EAAQP,EAAYG,EAAmC,CACpE,MAAMN,EAAW,UAAU,CAAE,IAAKG,EAAI,WAAAG,CAAW,CAAC,CACpD,CAEA,eAAeK,EACbR,EACAG,EACAF,EACkB,CAClB,IAAMC,EAAQD,GAAM,OAAS,IACvBG,EAAY,IAAI,KAAK,KAAK,IAAI,EAAIF,CAAK,EAY7C,OAVe,MAAML,EAAW,UAC9B,CAAE,IAAKG,EAAI,WAAAG,CAAW,EACtB,CAAE,KAAM,CAAE,UAAAC,CAAU,CAAE,CACxB,GAOc,aAAe,CAC/B,CAEA,MAAO,CAAE,WAAAL,EAAY,QAAAQ,EAAS,MAAAC,EAAO,qBAAAV,CAAqB,CAC5D,CDVO,SAASW,EAAMC,EAAiD,CACrE,IAAMC,EAAO,OAAOD,GAAiB,SAAW,CAAE,IAAKA,CAAa,EAAIA,EAClEE,EAASD,EAAK,QAAU,GAG1BE,EAAeF,EAAK,SAClBG,EAAWH,EAAK,IAAI,MAAM,sCAAsC,EAClE,CAACE,GAAgBC,GAAYA,EAAS,CAAC,IACzCD,EAAeC,EAAS,CAAC,GAE3BD,EAAeA,GAAgB,UAE/B,IAAME,EAAiBJ,EAAK,YAAc,oBAEpCK,EAAY,CAACL,EAAK,OACpBM,EAAsCN,EAAK,OAC3CO,EACAC,EAAY,GACZC,EAAuI,KAErIC,EAAkB,UAClBH,GAAMC,IAELF,IACHA,EAAS,IAAI,EAAAK,YAAgBX,EAAK,IAAK,CACrC,iBAAkB,CAACA,EAAK,IAAI,SAAS,gBAAgB,EACrD,GAAGA,EAAK,aACV,CAAC,GAGH,MAAMM,EAAO,QAAQ,EACrBE,EAAY,GACZD,EAAKD,EAAO,GAAGJ,CAAY,EAI3B,MADmBK,EAAG,WAAwBH,CAAc,EAC3C,YAAY,CAAE,UAAW,EAAG,EAAG,CAAE,WAAY,EAAK,CAAC,EAAE,MAAM,IAAM,CAElF,CAAC,EAEGJ,EAAK,MAAQ,CAACS,IAChBA,EAAOG,EAAgBL,EAAIP,EAAK,IAAI,IAG/BO,GAGHM,EAAoB,CACxB,MAAM,KAAKC,EAAYC,EAAsC,CAE3D,IAAMC,GADK,MAAMN,EAAgB,GACX,WAAwBN,CAAc,EACtDa,EAAShB,EAASa,EAClBI,KAAU,iBAAcH,CAAK,KAAI,wBAAqBA,CAAK,EAAIA,EACrE,MAAMC,EAAW,UACf,CAAE,IAAKC,CAAO,EACd,CACE,KAAM,CACJ,SAAUC,EACV,UAAW,IAAI,IACjB,CACF,EACA,CAAE,OAAQ,EAAK,CACjB,CACF,EAEA,MAAM,KAAKJ,EAAsC,CAE/C,IAAME,GADK,MAAMN,EAAgB,GACX,WAAwBN,CAAc,EACtDa,EAAShB,EAASa,EAClBK,EAAM,MAAMH,EAAW,QAAQ,CAAE,IAAKC,CAAO,CAAC,EACpD,GAAI,CAACE,EAAK,OAAO,KACjB,IAAMC,EAAMD,EAAI,SAChB,SAAI,2BAAwBC,CAAG,KAAU,0BAAuBA,CAAG,MAC/D,sBAAmBA,CAAG,EAAUA,EAEtC,EAEA,MAAM,gBAAgBN,EAAyC,CAC7D,IAAMO,EAAS,MAAMR,EAAM,KAAKC,CAAE,EAClC,OAAIO,IAAW,KAAa,QACxB,iBAAcA,CAAM,EAAUA,EAC3B,IACT,EAEA,MAAM,OAAOP,EAA2B,CAEtC,IAAME,GADK,MAAMN,EAAgB,GACX,WAAwBN,CAAc,EACtDa,EAAShB,EAASa,EACxB,MAAME,EAAW,UAAU,CAAE,IAAKC,CAAO,CAAC,CAC5C,EAEA,MAAM,KAAKK,EAAkG,CAE3G,IAAMN,GADK,MAAMN,EAAgB,GACX,WAAwBN,CAAc,EACtDmB,EAAetB,GAAUqB,GAAS,QAAU,IAC5CE,EAAQF,GAAS,OAAS,IAC1BG,EAAUF,EAAa,QAAQ,sBAAuB,MAAM,EAQlE,OADa,MALEP,EACZ,KAAK,CAAE,IAAK,CAAE,OAAQ,IAAIS,CAAO,EAAG,CAAE,CAAC,EACvC,KAAK,CAAE,UAAW,EAAG,CAAC,EACtB,MAAMD,CAAK,EAEY,QAAQ,GACtB,IAAIL,IAAQ,CACtB,GAAI,OAAOA,EAAI,GAAG,EAAE,MAAMlB,EAAO,MAAM,EACvC,UAAWkB,EAAI,UAAU,YAAY,CACvC,EAAE,CACJ,EAEA,MAAM,OAAuB,CACvBd,GAAaC,IACf,MAAMA,EAAO,MAAM,EACnBE,EAAY,GACZD,EAAK,OAET,EAGA,MAAM,WAAWO,EAAYQ,EAAsE,CAEjG,OADA,MAAMZ,EAAgB,EACjBD,EACEA,EAAK,WAAWK,EAAIQ,CAAO,EADhB,IAEpB,EAEA,MAAM,QAAQR,EAAYY,EAAmC,CAE3D,GADA,MAAMhB,EAAgB,EAClB,EAACD,EACL,OAAOA,EAAK,QAAQK,EAAIY,CAAU,CACpC,EAEA,MAAM,MAAMZ,EAAYY,EAAoB1B,EAA6C,CAEvF,OADA,MAAMU,EAAgB,EACjBD,EACEA,EAAK,MAAMK,EAAIY,EAAY1B,CAAI,EADpB,EAEpB,CACF,EAGA,OAAKA,EAAK,OACR,OAAOa,EAAM,WACb,OAAOA,EAAM,QACb,OAAOA,EAAM,OAGRA,CACT","names":["index_exports","__export","mongo","__toCommonJS","import_mongodb","import_durable","import_node_crypto","createMongoLock","db","options","lockCollectionName","collection","ensureLockCollection","tryAcquire","id","opts","ttlMs","ownerToken","expiresAt","result","error","release","renew","mongo","urlOrOptions","opts","prefix","databaseName","urlMatch","collectionName","ownClient","client","db","connected","lock","ensureConnected","MongoClientImpl","createMongoLock","store","id","state","collection","fullId","toStore","doc","raw","loaded","options","filterPrefix","limit","escaped","ownerToken"]}