{"version":3,"file":"index.cjs","names":["SYSTEM_CONTEXT_KEY"],"sources":["../../src/runtime/strategies/local-execution-strategy.ts","../../src/runtime/strategies/web-worker-strategy.ts","../../src/runtime/worker-runtime.ts"],"sourcesContent":["import { ILogger, PerformanceTimeEntry } from '@awesome-ecs/abstract/utils';\nimport { WorkerExecutionMode } from '../../abstract/worker-execution-mode';\nimport {\n  IWorkerExecutionStrategy,\n  WorkerResponseCallback\n} from '../../abstract/worker-execution-strategy';\nimport { IWorkerHandlerRegistry } from '../../abstract/worker-handler-registry';\nimport { WorkerRequestMessage, WorkerResponseMessage } from '../../abstract/worker-message';\nimport { WorkerStatus } from '../../components/worker-instance';\nimport { WorkerThreadEntity } from '../../entities/worker-thread';\n\nexport class LocalExecutionStrategy implements IWorkerExecutionStrategy {\n  readonly mode = WorkerExecutionMode.local;\n\n  constructor(\n    private readonly handlerRegistry: IWorkerHandlerRegistry,\n    private readonly logger: ILogger\n  ) {}\n\n  dispatch(\n    message: WorkerRequestMessage<any>,\n    entity: WorkerThreadEntity,\n    onResponse: WorkerResponseCallback\n  ): void {\n    entity.instance.status = WorkerStatus.busy;\n\n    try {\n      const handler = this.handlerRegistry.resolve(message.request.type);\n\n      if (!handler) {\n        this.logger.warn(\n          `[LOCAL-WORKER] No handler registered for request type: ${message.request.type}`\n        );\n        entity.instance.status = WorkerStatus.available;\n        return;\n      }\n\n      this.logger.debug(\n        '[LOCAL-WORKER] Executing handler locally',\n        message.request.type,\n        message.messageUid\n      );\n\n      const startTime = performance.now();\n      const result = handler.handle(message.request);\n      const duration = performance.now() - startTime;\n\n      const metric: PerformanceTimeEntry = {\n        name: message.request.type,\n        startedAt: startTime,\n        endedAt: startTime + duration,\n        msPassed: duration\n      };\n\n      const response: WorkerResponseMessage<any> = {\n        data: result.data,\n        transfer: result.transfer,\n        metrics: result.metrics ? [...result.metrics, metric] : [metric],\n        messageType: message.messageType,\n        messageUid: message.messageUid,\n        workerType: message.workerType\n      };\n\n      onResponse(entity, response);\n    } catch (error: any) {\n      this.logger.warn(`[LOCAL-WORKER] Handler error: ${error?.message}`, message.request.type, error);\n      entity.instance.status = WorkerStatus.error;\n    }\n  }\n}\n","import { ILogger } from '@awesome-ecs/abstract/utils';\nimport { WorkerExecutionMode } from '../../abstract/worker-execution-mode';\nimport { IWorkerExecutionStrategy } from '../../abstract/worker-execution-strategy';\nimport { WorkerRequestMessage } from '../../abstract/worker-message';\nimport { WorkerStatus } from '../../components/worker-instance';\nimport { WorkerThreadEntity } from '../../entities/worker-thread';\n\nexport class WebWorkerExecutionStrategy implements IWorkerExecutionStrategy {\n  readonly mode = WorkerExecutionMode.worker;\n\n  constructor(private readonly logger: ILogger) {}\n\n  dispatch(message: WorkerRequestMessage<any>, entity: WorkerThreadEntity): void {\n    this.logger.debug('Sending message to Worker', entity.identity.model.uid, message);\n\n    entity.instance.worker.postMessage(message);\n    entity.instance.status = WorkerStatus.busy;\n  }\n}\n","import {\n  IContextRepository,\n  IEntityProxy,\n  IEntityRepository,\n  SYSTEM_CONTEXT_KEY\n} from '@awesome-ecs/abstract/entities';\nimport { IMutableSystemContext } from '@awesome-ecs/abstract/systems';\nimport { ILogger } from '@awesome-ecs/abstract/utils';\nimport { WorkerThreadEventType } from '../abstract/events/event-type';\nimport { WorkerThreadEventMessageResponseData } from '../abstract/events/message';\nimport { WorkerEntityType } from '../abstract/types/entity-type';\nimport {\n  IWorkerExecutionStrategy,\n  WorkerResponseCallback\n} from '../abstract/worker-execution-strategy';\nimport { WorkerMessageType, WorkerResponseMessage } from '../abstract/worker-message';\nimport { IWorkerMessageQueue } from '../abstract/worker-queue';\nimport { WorkerStatus } from '../components/worker-instance';\nimport { WorkerThreadEntity } from '../entities/worker-thread';\nimport { WorkerPerformanceTracker } from '../utils/worker-performance-tracker';\n\ntype InstrumentedWorkerMessageQueue = IWorkerMessageQueue & {\n  flushWaitTimes?: () => number[];\n  markAllAsStalled?: () => void;\n};\n\n// @injectable\nexport class WorkerSystemRuntime {\n  private readonly onResponse: WorkerResponseCallback;\n\n  constructor(\n    private readonly workerMessageQueue: IWorkerMessageQueue,\n    private readonly entityRepository: IEntityRepository,\n    private readonly logger: ILogger,\n    private readonly strategy: IWorkerExecutionStrategy,\n    private readonly contextRepository?: IContextRepository,\n    private readonly performanceTracker?: WorkerPerformanceTracker\n  ) {\n    this.onResponse = (\n      entity: Readonly<WorkerThreadEntity>,\n      response: WorkerResponseMessage<any>\n    ) => {\n      if (!this.contextRepository) {\n        this.logger.warn(\n          '[WORKER-RUNTIME] No context repository configured for local response handling'\n        );\n        return;\n      }\n\n      const systemContext = this.contextRepository.get<IMutableSystemContext<WorkerThreadEntity>>(\n        entity.identity.proxy,\n        SYSTEM_CONTEXT_KEY\n      );\n      if (!systemContext) {\n        this.logger.warn('[WORKER-RUNTIME] No context found for entity', entity.identity.model.uid);\n        return;\n      }\n\n      const eventData: WorkerThreadEventMessageResponseData<any> = {\n        uid:\n          response.messageType === WorkerMessageType.health\n            ? WorkerThreadEventType.messageReceivedHealth\n            : WorkerThreadEventType.messageReceivedData,\n        message: response\n      };\n\n      systemContext.events.dispatchEvent(eventData, entity.identity.proxy as IEntityProxy);\n\n      this.logger.debug(\n        '[WORKER-LOCAL-RESPONSE] Dispatched response event',\n        response.messageType,\n        response.messageUid\n      );\n    };\n  }\n\n  runTick(): void {\n    while (this.workerMessageQueue.size > 0) {\n      const message = this.workerMessageQueue.peek();\n\n      if (!message) {\n        return;\n      }\n\n      const workers = this.entityRepository.getEntities<WorkerThreadEntity>(\n        WorkerEntityType.workerThread\n      );\n\n      if (!workers) {\n        this.logger.warn('No WorkerThreadEntity was registered for type: ', message.workerType);\n        return;\n      }\n\n      const availableWorker = this.findAvailableWorker(\n        workers.filter((worker): worker is WorkerThreadEntity => !!worker)\n      );\n\n      if (!availableWorker) {\n        this.logger.debug(\n          'No available worker found. Skipping remaining messages.',\n          message.workerType,\n          message.messageUid\n        );\n\n        if (this.performanceTracker) {\n          const queue = this.workerMessageQueue as InstrumentedWorkerMessageQueue;\n          if (typeof queue.markAllAsStalled === 'function') {\n            queue.markAllAsStalled();\n          }\n          this.performanceTracker.recordStalledCount(this.workerMessageQueue.size);\n        }\n        return;\n      }\n\n      this.strategy.dispatch(message, availableWorker, this.onResponse);\n      this.workerMessageQueue.dequeue();\n    }\n\n    this.flushQueueWaitTimes();\n  }\n\n  private flushQueueWaitTimes(): void {\n    if (!this.performanceTracker) return;\n    const queue = this.workerMessageQueue as InstrumentedWorkerMessageQueue;\n    if (typeof queue.flushWaitTimes === 'function') {\n      const waitTimes: number[] = queue.flushWaitTimes();\n      if (waitTimes.length) {\n        this.performanceTracker.recordWaitTimes(waitTimes);\n      }\n    }\n  }\n\n  private findAvailableWorker(\n    workers: ReadonlyArray<WorkerThreadEntity>\n  ): WorkerThreadEntity | undefined {\n    for (const worker of workers) {\n      if (worker.instance.status === WorkerStatus.available) {\n        return worker;\n      }\n    }\n\n    return undefined;\n  }\n}\n"],"mappings":";;;;;;AAWA,IAAa,yBAAb,MAAwE;CAInD;CACA;CAJnB,OAAS;CAET,YACE,iBACA,QACA;EAFiB,KAAA,kBAAA;EACA,KAAA,SAAA;CAChB;CAEH,SACE,SACA,QACA,YACM;EACN,OAAO,SAAS,SAAA;EAEhB,IAAI;GACF,MAAM,UAAU,KAAK,gBAAgB,QAAQ,QAAQ,QAAQ,IAAI;GAEjE,IAAI,CAAC,SAAS;IACZ,KAAK,OAAO,KACV,0DAA0D,QAAQ,QAAQ,MAC5E;IACA,OAAO,SAAS,SAAA;IAChB;GACF;GAEA,KAAK,OAAO,MACV,4CACA,QAAQ,QAAQ,MAChB,QAAQ,UACV;GAEA,MAAM,YAAY,YAAY,IAAI;GAClC,MAAM,SAAS,QAAQ,OAAO,QAAQ,OAAO;GAC7C,MAAM,WAAW,YAAY,IAAI,IAAI;GAErC,MAAM,SAA+B;IACnC,MAAM,QAAQ,QAAQ;IACtB,WAAW;IACX,SAAS,YAAY;IACrB,UAAU;GACZ;GAWA,WAAW,QAAQ;IARjB,MAAM,OAAO;IACb,UAAU,OAAO;IACjB,SAAS,OAAO,UAAU,CAAC,GAAG,OAAO,SAAS,MAAM,IAAI,CAAC,MAAM;IAC/D,aAAa,QAAQ;IACrB,YAAY,QAAQ;IACpB,YAAY,QAAQ;GAGI,CAAC;EAC7B,SAAS,OAAY;GACnB,KAAK,OAAO,KAAK,iCAAiC,OAAO,WAAW,QAAQ,QAAQ,MAAM,KAAK;GAC/F,OAAO,SAAS,SAAA;EAClB;CACF;AACF;;;AC9DA,IAAa,6BAAb,MAA4E;CAG7C;CAF7B,OAAS;CAET,YAAY,QAAkC;EAAjB,KAAA,SAAA;CAAkB;CAE/C,SAAS,SAAoC,QAAkC;EAC7E,KAAK,OAAO,MAAM,6BAA6B,OAAO,SAAS,MAAM,KAAK,OAAO;EAEjF,OAAO,SAAS,OAAO,YAAY,OAAO;EAC1C,OAAO,SAAS,SAAA;CAClB;AACF;;;ACSA,IAAa,sBAAb,MAAiC;CAIZ;CACA;CACA;CACA;CACA;CACA;CARnB;CAEA,YACE,oBACA,kBACA,QACA,UACA,mBACA,oBACA;EANiB,KAAA,qBAAA;EACA,KAAA,mBAAA;EACA,KAAA,SAAA;EACA,KAAA,WAAA;EACA,KAAA,oBAAA;EACA,KAAA,qBAAA;EAEjB,KAAK,cACH,QACA,aACG;GACH,IAAI,CAAC,KAAK,mBAAmB;IAC3B,KAAK,OAAO,KACV,+EACF;IACA;GACF;GAEA,MAAM,gBAAgB,KAAK,kBAAkB,IAC3C,OAAO,SAAS,OAChBA,+BAAAA,kBACF;GACA,IAAI,CAAC,eAAe;IAClB,KAAK,OAAO,KAAK,gDAAgD,OAAO,SAAS,MAAM,GAAG;IAC1F;GACF;GAEA,MAAM,YAAuD;IAC3D,KACE,SAAS,gBAAA,WAAA,4BAAA;IAGX,SAAS;GACX;GAEA,cAAc,OAAO,cAAc,WAAW,OAAO,SAAS,KAAqB;GAEnF,KAAK,OAAO,MACV,qDACA,SAAS,aACT,SAAS,UACX;EACF;CACF;CAEA,UAAgB;EACd,OAAO,KAAK,mBAAmB,OAAO,GAAG;GACvC,MAAM,UAAU,KAAK,mBAAmB,KAAK;GAE7C,IAAI,CAAC,SACH;GAGF,MAAM,UAAU,KAAK,iBAAiB,YAAA,eAEtC;GAEA,IAAI,CAAC,SAAS;IACZ,KAAK,OAAO,KAAK,mDAAmD,QAAQ,UAAU;IACtF;GACF;GAEA,MAAM,kBAAkB,KAAK,oBAC3B,QAAQ,QAAQ,WAAyC,CAAC,CAAC,MAAM,CACnE;GAEA,IAAI,CAAC,iBAAiB;IACpB,KAAK,OAAO,MACV,2DACA,QAAQ,YACR,QAAQ,UACV;IAEA,IAAI,KAAK,oBAAoB;KAC3B,MAAM,QAAQ,KAAK;KACnB,IAAI,OAAO,MAAM,qBAAqB,YACpC,MAAM,iBAAiB;KAEzB,KAAK,mBAAmB,mBAAmB,KAAK,mBAAmB,IAAI;IACzE;IACA;GACF;GAEA,KAAK,SAAS,SAAS,SAAS,iBAAiB,KAAK,UAAU;GAChE,KAAK,mBAAmB,QAAQ;EAClC;EAEA,KAAK,oBAAoB;CAC3B;CAEA,sBAAoC;EAClC,IAAI,CAAC,KAAK,oBAAoB;EAC9B,MAAM,QAAQ,KAAK;EACnB,IAAI,OAAO,MAAM,mBAAmB,YAAY;GAC9C,MAAM,YAAsB,MAAM,eAAe;GACjD,IAAI,UAAU,QACZ,KAAK,mBAAmB,gBAAgB,SAAS;EAErD;CACF;CAEA,oBACE,SACgC;EAChC,KAAK,MAAM,UAAU,SACnB,IAAI,OAAO,SAAS,WAAA,GAClB,OAAO;CAKb;AACF"}