{"version":3,"file":"SocketServerTransport.cjs","names":[],"sources":["../../src/SocketServerTransport.ts"],"sourcesContent":["import { Socket as IoSocket } from 'socket.io'\nimport { Transport, AuthMessage } from '@bsv/sdk'\n\nconst DEFAULT_MAX_PENDING_MESSAGES = 32\n\nexport interface SocketServerTransportOptions {\n  /** Maximum authentication messages that may be processed concurrently per socket. */\n  maxPendingMessages?: number\n  /** Receives contained authentication failures. The hook is never allowed to throw outward. */\n  onError?: (error: unknown) => void | Promise<void>\n}\n\n/**\n * Implements the Transport interface for a specific client socket.\n *\n * This transport simply relays AuthMessages over 'authMessage'\n * in the underlying Socket.IO connection.\n */\nexport class SocketServerTransport implements Transport {\n  private readonly maxPendingMessages: number\n  private readonly onError?: (error: unknown) => void | Promise<void>\n  private pendingMessages = 0\n  private failed = false\n\n  constructor(\n    private readonly socket: IoSocket,\n    options: SocketServerTransportOptions = {}\n  ) {\n    const maxPendingMessages = options.maxPendingMessages ?? DEFAULT_MAX_PENDING_MESSAGES\n    if (!Number.isSafeInteger(maxPendingMessages) || maxPendingMessages < 1) {\n      throw new RangeError('maxPendingMessages must be a positive safe integer')\n    }\n    this.maxPendingMessages = maxPendingMessages\n    this.onError = options.onError\n  }\n\n  async send(message: AuthMessage): Promise<void> {\n    // We'll emit with a special low-level event named: 'authMessage'\n    this.socket.emit('authMessage', message)\n  }\n\n  async onData(callback: (message: AuthMessage) => Promise<void>): Promise<void> {\n    // Listen for 'authMessage' from the client\n    this.socket.on('authMessage', (msg: AuthMessage) => {\n      return this.processMessage(msg, callback)\n    })\n  }\n\n  private async processMessage(\n    message: AuthMessage,\n    callback: (message: AuthMessage) => Promise<void>\n  ): Promise<void> {\n    if (this.failed) return\n    if (this.pendingMessages >= this.maxPendingMessages) {\n      this.fail(new Error('Authentication message concurrency limit exceeded'))\n      return\n    }\n\n    this.pendingMessages += 1\n    try {\n      await callback(message)\n    } catch (error) {\n      this.fail(error)\n    } finally {\n      this.pendingMessages -= 1\n    }\n  }\n\n  private fail(error: unknown): void {\n    if (this.failed) return\n    this.failed = true\n    void Promise.resolve()\n      .then(async () => await this.onError?.(error))\n      .catch(() => {})\n    try {\n      this.socket.disconnect(true)\n    } catch {\n      // A transport failure is already contained; disconnect errors are non-actionable here.\n    }\n  }\n}\n"],"mappings":";AAGA,MAAM,+BAA+B;;;;;;;AAerC,IAAa,wBAAb,MAAwD;CAOnC;CANnB;CACA;CACA,kBAA0B;CAC1B,SAAiB;CAEjB,YACE,QACA,UAAwC,CAAC,GACzC;EAFiB,KAAA,SAAA;EAGjB,MAAM,qBAAqB,QAAQ,sBAAsB;EACzD,IAAI,CAAC,OAAO,cAAc,kBAAkB,KAAK,qBAAqB,GACpE,MAAM,IAAI,WAAW,oDAAoD;EAE3E,KAAK,qBAAqB;EAC1B,KAAK,UAAU,QAAQ;CACzB;CAEA,MAAM,KAAK,SAAqC;EAE9C,KAAK,OAAO,KAAK,eAAe,OAAO;CACzC;CAEA,MAAM,OAAO,UAAkE;EAE7E,KAAK,OAAO,GAAG,gBAAgB,QAAqB;GAClD,OAAO,KAAK,eAAe,KAAK,QAAQ;EAC1C,CAAC;CACH;CAEA,MAAc,eACZ,SACA,UACe;EACf,IAAI,KAAK,QAAQ;EACjB,IAAI,KAAK,mBAAmB,KAAK,oBAAoB;GACnD,KAAK,qBAAK,IAAI,MAAM,mDAAmD,CAAC;GACxE;EACF;EAEA,KAAK,mBAAmB;EACxB,IAAI;GACF,MAAM,SAAS,OAAO;EACxB,SAAS,OAAO;GACd,KAAK,KAAK,KAAK;EACjB,UAAU;GACR,KAAK,mBAAmB;EAC1B;CACF;CAEA,KAAa,OAAsB;EACjC,IAAI,KAAK,QAAQ;EACjB,KAAK,SAAS;EACd,QAAa,QAAQ,CAAC,CACnB,KAAK,YAAY,MAAM,KAAK,UAAU,KAAK,CAAC,CAAC,CAC7C,YAAY,CAAC,CAAC;EACjB,IAAI;GACF,KAAK,OAAO,WAAW,IAAI;EAC7B,QAAQ,CAER;CACF;AACF"}