balanced-pool.js 5.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221
  1. 'use strict'
  2. const {
  3. BalancedPoolMissingUpstreamError,
  4. InvalidArgumentError
  5. } = require('../core/errors')
  6. const {
  7. PoolBase,
  8. kClients,
  9. kNeedDrain,
  10. kAddClient,
  11. kRemoveClient,
  12. kGetDispatcher
  13. } = require('./pool-base')
  14. const Pool = require('./pool')
  15. const { kUrl } = require('../core/symbols')
  16. const util = require('../core/util')
  17. const kFactory = Symbol('factory')
  18. const kOptions = Symbol('options')
  19. const kGreatestCommonDivisor = Symbol('kGreatestCommonDivisor')
  20. const kCurrentWeight = Symbol('kCurrentWeight')
  21. const kIndex = Symbol('kIndex')
  22. const kWeight = Symbol('kWeight')
  23. const kMaxWeightPerServer = Symbol('kMaxWeightPerServer')
  24. const kErrorPenalty = Symbol('kErrorPenalty')
  25. /**
  26. * Calculate the greatest common divisor of two numbers by
  27. * using the Euclidean algorithm.
  28. *
  29. * @param {number} a
  30. * @param {number} b
  31. * @returns {number}
  32. */
  33. function getGreatestCommonDivisor (a, b) {
  34. if (a === 0) return b
  35. while (b !== 0) {
  36. const t = b
  37. b = a % b
  38. a = t
  39. }
  40. return a
  41. }
  42. function defaultFactory (origin, opts) {
  43. return new Pool(origin, opts)
  44. }
  45. class BalancedPool extends PoolBase {
  46. constructor (upstreams = [], { factory = defaultFactory, connect, tls, ...opts } = {}) {
  47. if (typeof factory !== 'function') {
  48. throw new InvalidArgumentError('factory must be a function.')
  49. }
  50. super(opts)
  51. if (connect && typeof connect !== 'function') connect = { ...connect }
  52. if (tls && typeof tls !== 'function') tls = { ...tls }
  53. this[kOptions] = { ...util.deepClone(opts), connect, tls }
  54. this[kOptions].interceptors = opts.interceptors
  55. ? { ...opts.interceptors }
  56. : undefined
  57. this[kIndex] = -1
  58. this[kCurrentWeight] = 0
  59. this[kMaxWeightPerServer] = this[kOptions].maxWeightPerServer || 100
  60. this[kErrorPenalty] = this[kOptions].errorPenalty || 15
  61. if (!Array.isArray(upstreams)) {
  62. upstreams = [upstreams]
  63. }
  64. this[kFactory] = factory
  65. for (const upstream of upstreams) {
  66. this.addUpstream(upstream)
  67. }
  68. this._updateBalancedPoolStats()
  69. }
  70. addUpstream (upstream) {
  71. const upstreamOrigin = util.parseOrigin(upstream).origin
  72. if (this[kClients].find((pool) => (
  73. pool[kUrl].origin === upstreamOrigin &&
  74. pool.closed !== true &&
  75. pool.destroyed !== true
  76. ))) {
  77. return this
  78. }
  79. const pool = this[kFactory](upstreamOrigin, this[kOptions])
  80. this[kAddClient](pool)
  81. pool.on('connect', () => {
  82. pool[kWeight] = Math.min(this[kMaxWeightPerServer], pool[kWeight] + this[kErrorPenalty])
  83. })
  84. pool.on('connectionError', () => {
  85. pool[kWeight] = Math.max(1, pool[kWeight] - this[kErrorPenalty])
  86. this._updateBalancedPoolStats()
  87. })
  88. pool.on('disconnect', (...args) => {
  89. const err = args[2]
  90. if (err && err.code === 'UND_ERR_SOCKET') {
  91. // decrease the weight of the pool.
  92. pool[kWeight] = Math.max(1, pool[kWeight] - this[kErrorPenalty])
  93. this._updateBalancedPoolStats()
  94. }
  95. })
  96. for (const client of this[kClients]) {
  97. client[kWeight] = this[kMaxWeightPerServer]
  98. }
  99. this._updateBalancedPoolStats()
  100. return this
  101. }
  102. _updateBalancedPoolStats () {
  103. let result = 0
  104. for (let i = 0; i < this[kClients].length; i++) {
  105. result = getGreatestCommonDivisor(this[kClients][i][kWeight], result)
  106. }
  107. this[kGreatestCommonDivisor] = result
  108. }
  109. removeUpstream (upstream) {
  110. const upstreamOrigin = util.parseOrigin(upstream).origin
  111. const pool = this[kClients].find((pool) => (
  112. pool[kUrl].origin === upstreamOrigin &&
  113. pool.closed !== true &&
  114. pool.destroyed !== true
  115. ))
  116. if (pool) {
  117. this[kRemoveClient](pool)
  118. }
  119. return this
  120. }
  121. getUpstream (upstream) {
  122. const upstreamOrigin = util.parseOrigin(upstream).origin
  123. return this[kClients].find((pool) => (
  124. pool[kUrl].origin === upstreamOrigin &&
  125. pool.closed !== true &&
  126. pool.destroyed !== true
  127. ))
  128. }
  129. get upstreams () {
  130. return this[kClients]
  131. .filter(dispatcher => dispatcher.closed !== true && dispatcher.destroyed !== true)
  132. .map((p) => p[kUrl].origin)
  133. }
  134. [kGetDispatcher] () {
  135. // We validate that pools is greater than 0,
  136. // otherwise we would have to wait until an upstream
  137. // is added, which might never happen.
  138. if (this[kClients].length === 0) {
  139. throw new BalancedPoolMissingUpstreamError()
  140. }
  141. const dispatcher = this[kClients].find(dispatcher => (
  142. !dispatcher[kNeedDrain] &&
  143. dispatcher.closed !== true &&
  144. dispatcher.destroyed !== true
  145. ))
  146. if (!dispatcher) {
  147. return
  148. }
  149. const allClientsBusy = this[kClients].map(pool => pool[kNeedDrain]).reduce((a, b) => a && b, true)
  150. if (allClientsBusy) {
  151. return
  152. }
  153. let counter = 0
  154. let maxWeightIndex = this[kClients].findIndex(pool => !pool[kNeedDrain])
  155. while (counter++ < this[kClients].length) {
  156. this[kIndex] = (this[kIndex] + 1) % this[kClients].length
  157. const pool = this[kClients][this[kIndex]]
  158. // find pool index with the largest weight
  159. if (pool[kWeight] > this[kClients][maxWeightIndex][kWeight] && !pool[kNeedDrain]) {
  160. maxWeightIndex = this[kIndex]
  161. }
  162. // decrease the current weight every `this[kClients].length`.
  163. if (this[kIndex] === 0) {
  164. // Set the current weight to the next lower weight.
  165. this[kCurrentWeight] = this[kCurrentWeight] - this[kGreatestCommonDivisor]
  166. if (this[kCurrentWeight] <= 0) {
  167. this[kCurrentWeight] = this[kMaxWeightPerServer]
  168. }
  169. }
  170. if (pool[kWeight] >= this[kCurrentWeight] && (!pool[kNeedDrain])) {
  171. return pool
  172. }
  173. }
  174. this[kCurrentWeight] = this[kClients][maxWeightIndex][kWeight]
  175. this[kIndex] = maxWeightIndex
  176. return this[kClients][maxWeightIndex]
  177. }
  178. }
  179. module.exports = BalancedPool