Team Ai
Datasetpublic

codekingpro/portable-devtools

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