codekingpro/portable-devtools
114k
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 