{"version":3,"sources":["../src/context.ts"],"names":["Middleware","Run","LazyPromise","tasks","handler","runContext","Symbol","toStringTag","observe","fn","push","emitter","context","value","Object","assign","middleware","fns","isFunction","bind","before","executeSequentially","splice","Infinity","asyncIterator","response","then","isAsyncIterable","Error","RunContext","Serializable","AsyncLocalStorage","controller","runId","groupId","parentId","runParams","createdAt","signal","abort","reason","instance","input","parent","Date","params","createRandomHash","omit","AbortController","registerSignals","child","trace","id","parentRunId","pipe","destroy","FrameworkError","enter","getStore","namespace","creator","internal","finishEvent","startEvent","output","undefined","emit","result","Promise","race","run","_","reject","addEventListener","setTimeout","_e","e","ensure","register","createSnapshot","shallowCopy","loadSnapshot","snapshot"],"mappings":";;;;;;;;;;;;;;AA4BO,MAAeA,UAAAA,CAAAA;EA5BtB;;;AA8BA;AAOO,MAAMC,YAA+CC,uBAAAA,CAAAA;EArC5D;;;;AAsCqBC,EAAAA,KAAAA,GAAiC,EAAA;AAEpD,EAAA,WAAA,CACEC,SACmBC,UAAAA,EACnB;AACA,IAAA,KAAA,CAAMD,OAAAA,CAAAA,EAAAA,IAAAA,CAFaC,UAAAA,GAAAA,UAAAA;AAGrB,EAAA;EAES,CAACC,MAAAA,CAAOC,WAAW,IAAI,SAAA;AAEhCC,EAAAA,OAAAA,CAAQC,EAAAA,EAAmD;AACzD,IAAA,IAAA,CAAKN,MAAMO,IAAAA,CAAK,YAAYD,GAAG,IAAA,CAAKJ,UAAAA,CAAWM,OAAO,CAAA,CAAA;AACtD,IAAA,OAAO,IAAA;AACT,EAAA;AAEAC,EAAAA,OAAAA,CAAQC,KAAAA,EAAe;AACrB,IAAA,IAAA,CAAKV,KAAAA,CAAMO,KAAK,YAAA;AACdI,MAAAA,MAAAA,CAAOC,MAAAA,CAAO,IAAA,CAAKV,UAAAA,CAAWO,OAAAA,EAASC,KAAAA,CAAAA;AACvCC,MAAAA,MAAAA,CAAOC,MAAAA,CAAO,IAAA,CAAKV,UAAAA,CAAWM,OAAAA,CAAQC,SAASC,KAAAA,CAAAA;IACjD,CAAA,CAAA;AACA,IAAA,OAAO,IAAA;AACT,EAAA;AAEAG,EAAAA,UAAAA,CAAAA,GAAcC,GAAAA,EAA0B;AACtC,IAAA,KAAA,MAAWR,MAAMQ,GAAAA,EAAK;AACpB,MAAA,IAAA,CAAKd,KAAAA,CAAMO,IAAAA,CAAK,YACdQ,iBAAAA,CAAWT,EAAAA,CAAAA,GAAMA,EAAAA,CAAG,IAAA,CAAKJ,UAAU,CAAA,GAAII,EAAAA,CAAGU,IAAAA,CAAK,IAAA,CAAKd,UAAU,CAAA,CAAA;AAElE,IAAA;AACA,IAAA,OAAO,IAAA;AACT,EAAA;AAEA,EAAA,MAAgBe,MAAAA,GAAwB;AACtC,IAAA,MAAM,MAAMA,MAAAA,EAAAA;AACZ,IAAA,MAAMC,gCAAoB,IAAA,CAAKlB,KAAAA,CAAMmB,MAAAA,CAAO,CAAA,EAAGC,QAAAA,CAAAA,CAAAA;AACjD,EAAA;;EAGA,QAAQjB,MAAAA,CAAOkB,aAAa,CAAA,GAAiE;AAC3F,IAAA,MAAMC,QAAAA,GAAW,MAAM,IAAA,CAAKC,IAAAA,EAAI;AAChC,IAAA,IAAIC,0BAAAA,CAAgBF,QAAAA,CAAAA,EAAW;AAC7B,MAAA,OAAOA,QAAAA;IACT,CAAA,MAAO;AACL,MAAA,MAAM,IAAIG,MAAM,yBAAA,CAAA;AAClB,IAAA;AACF,EAAA;AACF;AAOO,MAAMC,mBAAmDC,6BAAAA,CAAAA;EA5FhE;;;;;EA6FE,OAAO,QAAA,GAAW,IAAIC,kCAAAA,EAAAA;AAEHC,EAAAA,UAAAA;AACHC,EAAAA,KAAAA;AACAC,EAAAA,OAAAA;AACAC,EAAAA,QAAAA;AACAxB,EAAAA,OAAAA;AACAC,EAAAA,OAAAA;AACTwB,EAAAA,SAAAA;AACSC,EAAAA,SAAAA;AAEhB,EAAA,IAAIC,MAAAA,GAAS;AACX,IAAA,OAAO,KAAKN,UAAAA,CAAWM,MAAAA;AACzB,EAAA;AAEAC,EAAAA,KAAAA,CAAMC,MAAAA,EAAgB;AACpB,IAAA,IAAA,CAAKR,UAAAA,CAAWO,MAAMC,MAAAA,CAAAA;AACxB,EAAA;EAEA,WAAA,CACkBC,QAAAA,EACGC,OACnBC,MAAAA,EACA;AACA,IAAA,KAAA,EAAK,EAAA,IAAA,CAJWF,QAAAA,GAAAA,QAAAA,EAAAA,KACGC,KAAAA,GAAAA,KAAAA;AAInB,IAAA,IAAA,CAAKL,SAAAA,uBAAgBO,IAAAA,EAAAA;AACrB,IAAA,IAAA,CAAKR,YAAYM,KAAAA,CAAMG,MAAAA;AACvB,IAAA,IAAA,CAAKZ,KAAAA,GAAQa,0BAAiB,CAAA,CAAA;AAC9B,IAAA,IAAA,CAAKX,WAAWQ,MAAAA,EAAQV,KAAAA;AACxB,IAAA,IAAA,CAAKC,OAAAA,GAAUS,MAAAA,EAAQT,OAAAA,IAAWY,yBAAAA,EAAAA;AAClC,IAAA,IAAA,CAAKlC,OAAAA,GAAUmC,WAAAA,CAAMJ,MAAAA,EAAQ/B,OAAAA,IAAW,EAAC,EAAW;AAAC,MAAA,IAAA;AAAM,MAAA;AAAW,KAAA,CAAA;AAEtE,IAAA,IAAA,CAAKoB,UAAAA,GAAa,IAAIgB,eAAAA,EAAAA;AACtBC,IAAAA,gCAAAA,CAAgB,KAAKjB,UAAAA,EAAY;MAACU,KAAAA,CAAMJ,MAAAA;MAAQK,MAAAA,EAAQL;AAAO,KAAA,CAAA;AAE/D,IAAA,IAAA,CAAK3B,OAAAA,GAAU8B,QAAAA,CAAS9B,OAAAA,CAAQuC,KAAAA,CAAyB;AACvDtC,MAAAA,OAAAA,EAAS,IAAA,CAAKA,OAAAA;MACduC,KAAAA,EAAO;AACLC,QAAAA,EAAAA,EAAI,IAAA,CAAKlB,OAAAA;AACTD,QAAAA,KAAAA,EAAO,IAAA,CAAKA,KAAAA;AACZoB,QAAAA,WAAAA,EAAaV,MAAAA,EAAQV;AACvB;KACF,CAAA;AACA,IAAA,IAAIU,MAAAA,EAAQ;AACV,MAAA,IAAA,CAAKhC,OAAAA,CAAQ2C,IAAAA,CAAKX,MAAAA,CAAOhC,OAAO,CAAA;AAClC,IAAA;AACF,EAAA;EAEA4C,OAAAA,GAAU;AACR,IAAA,IAAA,CAAK5C,QAAQ4C,OAAAA,EAAO;AACpB,IAAA,IAAA,CAAKvB,UAAAA,CAAWO,KAAAA,CAAM,IAAIiB,yBAAAA,CAAe,oBAAA,CAAA,CAAA;AAC3C,EAAA;EAEA,OAAOC,KAAAA,CACLhB,QAAAA,EACAC,KAAAA,EACAjC,EAAAA,EACA;AACA,IAAA,MAAMkC,MAAAA,GAASd,UAAAA,CAAW,QAAA,CAAS6B,QAAAA,EAAQ;AAC3C,IAAA,MAAMrD,UAAAA,GAAa,IAAIwB,UAAAA,CAAWY,QAAAA,EAAUC,OAAOC,MAAAA,CAAAA;AAEnD,IAAA,OAAO,IAAI1C,IAAI,YAAA;AACb,MAAA,MAAMU,OAAAA,GAAUN,UAAAA,CAAWM,OAAAA,CAAQuC,KAAAA,CAA2B;QAC5DS,SAAAA,EAAW;AAAC,UAAA;;QACZC,OAAAA,EAASvD,UAAAA;QACTO,OAAAA,EAAS;UAAEiD,QAAAA,EAAU;AAAK;OAC5B,CAAA;AAEA,MAAA,MAAMC,cAAc,EAAC;AACrB,MAAA,IAAI;AACF,QAAA,MAAMC,UAAAA,GAA+D;AACnErB,UAAAA,KAAAA,EAAOrC,UAAAA,CAAW+B,SAAAA;UAClB4B,MAAAA,EAAQC,KAAAA;AACV,SAAA;AACA,QAAA,MAAMtD,OAAAA,CAAQuD,IAAAA,CAAK,OAAA,EAASH,UAAAA,CAAAA;AAG5B1D,QAAAA,UAAAA,CAAW+B,YAAY2B,UAAAA,CAAWrB,KAAAA;AAClCoB,QAAAA,WAAAA,CAAYpB,QAAQqB,UAAAA,CAAWrB,KAAAA;AAE/B,QAAA,MAAMyB,MAAAA,GAAa,MAAMC,OAAAA,CAAQC,IAAAA,CAAK;UACpCxC,UAAAA,CAAW,QAAA,CAASyC,GAAAA,CAClBjE,UAAAA,EACA0D,UAAAA,CAAWC,MAAAA,KAAWC,SAAYxD,EAAAA,GAAK,YAAYsD,UAAAA,CAAWC,MAAAA,EAC9D3D,UAAAA,CAAAA;AAEF,UAAA,IAAI+D,QAAe,CAACG,CAAAA,EAAGC,WACrBnE,UAAAA,CAAWiC,MAAAA,CAAOmC,iBAAiB,OAAA,EAAS,MAC1CC,UAAAA,CAAW,MAAMF,OAAOnE,UAAAA,CAAWiC,MAAAA,CAAOE,MAAM,CAAA,EAAG,CAAA,CAAA,CAAA;AAGxD,SAAA,CAAA;AACDsB,QAAAA,WAAAA,CAAYE,MAAAA,GAASG,MAAAA;AACrB,QAAA,MAAMxD,OAAAA,CAAQuD,IAAAA,CAAK,SAAA,EAAWC,MAAAA,CAAAA;AAC9B,QAAA,OAAOA,MAAAA;AACT,MAAA,CAAA,CAAA,OAASQ,EAAAA,EAAI;AACX,QAAA,MAAMC,CAAAA,GAAIpB,yBAAAA,CAAeqB,MAAAA,CAAOF,EAAAA,CAAAA;AAChCb,QAAAA,WAAAA,CAAYE,MAAAA,GAASY,CAAAA;AACrB,QAAA,MAAMjE,OAAAA,CAAQuD,IAAAA,CAAK,OAAA,EAASU,CAAAA,CAAAA;AAC5B,QAAA,MAAMA,CAAAA;MACR,CAAA,SAAA;AACE,QAAA,MAAMjE,OAAAA,CAAQuD,IAAAA,CAAK,QAAA,EAAUJ,WAAAA,CAAAA;AAC7BzD,QAAAA,UAAAA,CAAWkD,OAAAA,EAAO;AACpB,MAAA;AACF,IAAA,CAAA,EAAGlD,UAAAA,CAAAA;AACL,EAAA;EAEA;AACE,IAAA,IAAA,CAAKyE,QAAAA,EAAQ;AACf;EAEAC,cAAAA,GAAiB;AACf,IAAA,OAAO;AACL/C,MAAAA,UAAAA,EAAY,IAAA,CAAKA,UAAAA;AACjBC,MAAAA,KAAAA,EAAO,IAAA,CAAKA,KAAAA;AACZC,MAAAA,OAAAA,EAAS,IAAA,CAAKA,OAAAA;AACdC,MAAAA,QAAAA,EAAU,IAAA,CAAKA,QAAAA;AACfxB,MAAAA,OAAAA,EAAS,IAAA,CAAKA,OAAAA;MACdC,OAAAA,EAASoE,qBAAAA,CAAY,KAAKpE,OAAO,CAAA;MACjCwB,SAAAA,EAAW4C,qBAAAA,CAAY,KAAK5C,SAAS,CAAA;MACrCC,SAAAA,EAAW,IAAIO,IAAAA,CAAK,IAAA,CAAKP,SAAS;AACpC,KAAA;AACF,EAAA;AAEA4C,EAAAA,YAAAA,CAAaC,QAAAA,EAAkD;AAC7DpE,IAAAA,MAAAA,CAAOC,MAAAA,CAAO,MAAMmE,QAAAA,CAAAA;AACtB,EAAA;AACF","file":"context.cjs","sourcesContent":["/**\n * Copyright 2025 © BeeAI a Series of LF Projects, LLC\n * SPDX-License-Identifier: Apache-2.0\n */\n\nimport { AsyncLocalStorage } from \"node:async_hooks\";\nimport { Emitter } from \"@/emitter/emitter.js\";\nimport { createRandomHash } from \"@/internals/helpers/hash.js\";\nimport { isFunction, omit } from \"remeda\";\nimport { Callback, InferCallbackValue } from \"@/emitter/types.js\";\nimport { registerSignals } from \"@/internals/helpers/cancellation.js\";\nimport { Serializable } from \"@/internals/serializable.js\";\nimport { executeSequentially, LazyPromise } from \"@/internals/helpers/promise.js\";\nimport { FrameworkError } from \"@/errors.js\";\nimport { shallowCopy } from \"@/serializer/utils.js\";\nimport { isAsyncIterable } from \"@/internals/helpers/stream.js\";\n\nexport interface RunInstance<T = any> {\n  emitter: Emitter<T>;\n}\n\nexport interface RunContextCallbacks {\n  start: Callback<{ input: any; output: any }>;\n  success: Callback<unknown>; // TODO: should be updated to Callback<{ input: any; output: any }>\n  error: Callback<Error>;\n  finish: Callback<{ input: any; output?: any; error?: FrameworkError }>;\n}\n\nexport abstract class Middleware<T extends RunInstance = any> {\n  abstract bind(ctx: RunContext<T>): void;\n}\n\nexport type GetRunContext<T, P = any> = T extends RunInstance ? RunContext<T, P> : never;\nexport type GetRunInstance<T> = T extends RunInstance<infer P> ? P : never;\nexport type MiddlewareFn<T> = (context: GetRunContext<T>) => void;\nexport type MiddlewareType<T extends RunInstance> = Middleware<T> | MiddlewareFn<T>;\n\nexport class Run<R, I extends RunInstance, P = any> extends LazyPromise<R> {\n  protected readonly tasks: (() => Promise<void>)[] = [];\n\n  constructor(\n    handler: () => Promise<R>,\n    protected readonly runContext: GetRunContext<I, P>,\n  ) {\n    super(handler);\n  }\n\n  readonly [Symbol.toStringTag] = \"Promise\";\n\n  observe(fn: (emitter: Emitter<GetRunInstance<I>>) => void) {\n    this.tasks.push(async () => fn(this.runContext.emitter));\n    return this;\n  }\n\n  context(value: object) {\n    this.tasks.push(async () => {\n      Object.assign(this.runContext.context, value);\n      Object.assign(this.runContext.emitter.context, value);\n    });\n    return this;\n  }\n\n  middleware(...fns: MiddlewareType<I>[]) {\n    for (const fn of fns) {\n      this.tasks.push(async () =>\n        isFunction(fn) ? fn(this.runContext) : fn.bind(this.runContext),\n      );\n    }\n    return this;\n  }\n\n  protected async before(): Promise<void> {\n    await super.before();\n    await executeSequentially(this.tasks.splice(0, Infinity));\n  }\n\n  // @ts-expect-error too complex\n  async *[Symbol.asyncIterator](): R extends AsyncIterable<infer X> ? AsyncIterator<X> : never {\n    const response = await this.then();\n    if (isAsyncIterable(response)) {\n      yield* response;\n    } else {\n      throw new Error(\"Result is not iterable!\");\n    }\n  }\n}\n\nexport interface RunContextInput<P> {\n  params: P;\n  signal?: AbortSignal;\n}\n\nexport class RunContext<T extends RunInstance, P = any> extends Serializable {\n  static #storage = new AsyncLocalStorage<RunContext<any>>();\n\n  protected readonly controller: AbortController;\n  public readonly runId: string;\n  public readonly groupId: string;\n  public readonly parentId?: string;\n  public readonly emitter;\n  public readonly context: object;\n  public runParams: P;\n  public readonly createdAt: Date;\n\n  get signal() {\n    return this.controller.signal;\n  }\n\n  abort(reason?: Error) {\n    this.controller.abort(reason);\n  }\n\n  constructor(\n    public readonly instance: T,\n    protected readonly input: RunContextInput<P>,\n    parent?: RunContext<any>,\n  ) {\n    super();\n    this.createdAt = new Date();\n    this.runParams = input.params;\n    this.runId = createRandomHash(5);\n    this.parentId = parent?.runId;\n    this.groupId = parent?.groupId ?? createRandomHash();\n    this.context = omit((parent?.context ?? {}) as any, [\"id\", \"parentId\"]);\n\n    this.controller = new AbortController();\n    registerSignals(this.controller, [input.signal, parent?.signal]);\n\n    this.emitter = instance.emitter.child<GetRunInstance<T>>({\n      context: this.context,\n      trace: {\n        id: this.groupId,\n        runId: this.runId,\n        parentRunId: parent?.runId,\n      },\n    });\n    if (parent) {\n      this.emitter.pipe(parent.emitter);\n    }\n  }\n\n  destroy() {\n    this.emitter.destroy();\n    this.controller.abort(new FrameworkError(\"Context destroyed.\"));\n  }\n\n  static enter<C2 extends RunInstance, R2, P2>(\n    instance: C2,\n    input: RunContextInput<P2>,\n    fn: (context: GetRunContext<C2, P2>) => Promise<R2>,\n  ) {\n    const parent = RunContext.#storage.getStore();\n    const runContext = new RunContext(instance, input, parent) as GetRunContext<C2, P2>;\n\n    return new Run(async () => {\n      const emitter = runContext.emitter.child<RunContextCallbacks>({\n        namespace: [\"run\"],\n        creator: runContext,\n        context: { internal: true },\n      });\n\n      const finishEvent = {} as InferCallbackValue<RunContextCallbacks[\"finish\"]>;\n      try {\n        const startEvent: InferCallbackValue<RunContextCallbacks[\"start\"]> = {\n          input: runContext.runParams,\n          output: undefined,\n        };\n        await emitter.emit(\"start\", startEvent);\n\n        // Copy back any modifications made by middleware to run_params\n        runContext.runParams = startEvent.input;\n        finishEvent.input = startEvent.input;\n\n        const result: R2 = await Promise.race([\n          RunContext.#storage.run(\n            runContext,\n            startEvent.output === undefined ? fn : async () => startEvent.output,\n            runContext,\n          ),\n          new Promise<never>((_, reject) =>\n            runContext.signal.addEventListener(\"abort\", () =>\n              setTimeout(() => reject(runContext.signal.reason), 0),\n            ),\n          ),\n        ]);\n        finishEvent.output = result;\n        await emitter.emit(\"success\", result);\n        return result;\n      } catch (_e) {\n        const e = FrameworkError.ensure(_e);\n        finishEvent.output = e;\n        await emitter.emit(\"error\", e);\n        throw e;\n      } finally {\n        await emitter.emit(\"finish\", finishEvent);\n        runContext.destroy();\n      }\n    }, runContext);\n  }\n\n  static {\n    this.register();\n  }\n\n  createSnapshot() {\n    return {\n      controller: this.controller,\n      runId: this.runId,\n      groupId: this.groupId,\n      parentId: this.parentId,\n      emitter: this.emitter,\n      context: shallowCopy(this.context),\n      runParams: shallowCopy(this.runParams),\n      createdAt: new Date(this.createdAt),\n    };\n  }\n\n  loadSnapshot(snapshot: ReturnType<typeof this.createSnapshot>) {\n    Object.assign(this, snapshot);\n  }\n}\n"]}