dispatcher-base.js 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184
  1. 'use strict'
  2. const Dispatcher = require('./dispatcher')
  3. const UnwrapHandler = require('../handler/unwrap-handler')
  4. const {
  5. ClientDestroyedError,
  6. ClientClosedError,
  7. InvalidArgumentError
  8. } = require('../core/errors')
  9. const { kDestroy, kClose, kClosed, kDestroyed, kDispatch } = require('../core/symbols')
  10. const kOnDestroyed = Symbol('onDestroyed')
  11. const kOnClosed = Symbol('onClosed')
  12. const kWebSocketOptions = Symbol('webSocketOptions')
  13. class DispatcherBase extends Dispatcher {
  14. /** @type {boolean} */
  15. [kDestroyed] = false;
  16. /** @type {Array<Function|null} */
  17. [kOnDestroyed] = null;
  18. /** @type {boolean} */
  19. [kClosed] = false;
  20. /** @type {Array<Function>|null} */
  21. [kOnClosed] = null
  22. /**
  23. * @param {import('../../types/dispatcher').DispatcherOptions} [opts]
  24. */
  25. constructor (opts) {
  26. super()
  27. this[kWebSocketOptions] = opts?.webSocket ?? {}
  28. }
  29. /**
  30. * @returns {import('../../types/dispatcher').WebSocketOptions}
  31. */
  32. get webSocketOptions () {
  33. return {
  34. maxFragments: this[kWebSocketOptions].maxFragments ?? 131072,
  35. maxPayloadSize: this[kWebSocketOptions].maxPayloadSize ?? 128 * 1024 * 1024 // 128 MB default
  36. }
  37. }
  38. /** @returns {boolean} */
  39. get destroyed () {
  40. return this[kDestroyed]
  41. }
  42. /** @returns {boolean} */
  43. get closed () {
  44. return this[kClosed]
  45. }
  46. close (callback) {
  47. if (callback === undefined) {
  48. return new Promise((resolve, reject) => {
  49. this.close((err, data) => {
  50. return err ? reject(err) : resolve(data)
  51. })
  52. })
  53. }
  54. if (typeof callback !== 'function') {
  55. throw new InvalidArgumentError('invalid callback')
  56. }
  57. if (this[kDestroyed]) {
  58. const err = new ClientDestroyedError()
  59. queueMicrotask(() => callback(err, null))
  60. return
  61. }
  62. if (this[kClosed]) {
  63. if (this[kOnClosed]) {
  64. this[kOnClosed].push(callback)
  65. } else {
  66. queueMicrotask(() => callback(null, null))
  67. }
  68. return
  69. }
  70. this[kClosed] = true
  71. this[kOnClosed] ??= []
  72. this[kOnClosed].push(callback)
  73. const onClosed = () => {
  74. const callbacks = this[kOnClosed]
  75. this[kOnClosed] = null
  76. for (let i = 0; i < callbacks.length; i++) {
  77. callbacks[i](null, null)
  78. }
  79. }
  80. // Should not error.
  81. this[kClose]()
  82. .then(() => this.destroy())
  83. .then(() => queueMicrotask(onClosed))
  84. }
  85. destroy (err, callback) {
  86. if (typeof err === 'function') {
  87. callback = err
  88. err = null
  89. }
  90. if (callback === undefined) {
  91. return new Promise((resolve, reject) => {
  92. this.destroy(err, (err, data) => {
  93. return err ? reject(err) : resolve(data)
  94. })
  95. })
  96. }
  97. if (typeof callback !== 'function') {
  98. throw new InvalidArgumentError('invalid callback')
  99. }
  100. if (this[kDestroyed]) {
  101. if (this[kOnDestroyed]) {
  102. this[kOnDestroyed].push(callback)
  103. } else {
  104. queueMicrotask(() => callback(null, null))
  105. }
  106. return
  107. }
  108. if (!err) {
  109. err = new ClientDestroyedError()
  110. }
  111. this[kDestroyed] = true
  112. this[kOnDestroyed] ??= []
  113. this[kOnDestroyed].push(callback)
  114. const onDestroyed = () => {
  115. const callbacks = this[kOnDestroyed]
  116. this[kOnDestroyed] = null
  117. for (let i = 0; i < callbacks.length; i++) {
  118. callbacks[i](null, null)
  119. }
  120. }
  121. // Should not error.
  122. this[kDestroy](err)
  123. .then(() => queueMicrotask(onDestroyed))
  124. }
  125. dispatch (opts, handler) {
  126. if (!handler || typeof handler !== 'object') {
  127. throw new InvalidArgumentError('handler must be an object')
  128. }
  129. handler = UnwrapHandler.unwrap(handler)
  130. try {
  131. if (!opts || typeof opts !== 'object') {
  132. throw new InvalidArgumentError('opts must be an object.')
  133. }
  134. if (this[kDestroyed] || this[kOnDestroyed]) {
  135. throw new ClientDestroyedError()
  136. }
  137. if (this[kClosed]) {
  138. throw new ClientClosedError()
  139. }
  140. return this[kDispatch](opts, handler)
  141. } catch (err) {
  142. if (typeof handler.onError !== 'function') {
  143. throw err
  144. }
  145. handler.onError(err)
  146. return false
  147. }
  148. }
  149. }
  150. module.exports = DispatcherBase