Line | Count | Source |
1 | | /* |
2 | | * Queue management functions. |
3 | | * |
4 | | * Copyright 2000-2009 Willy Tarreau <w@1wt.eu> |
5 | | * |
6 | | * This program is free software; you can redistribute it and/or |
7 | | * modify it under the terms of the GNU General Public License |
8 | | * as published by the Free Software Foundation; either version |
9 | | * 2 of the License, or (at your option) any later version. |
10 | | * |
11 | | */ |
12 | | |
13 | | /* Short explanation on the locking, which is far from being trivial : a |
14 | | * pendconn is a list element which necessarily is associated with an existing |
15 | | * stream. It has pendconn->strm always valid. A pendconn may only be in one of |
16 | | * these three states : |
17 | | * - unlinked : in this case it is an empty list head ; |
18 | | * - linked into the server's queue ; |
19 | | * - linked into the proxy's queue. |
20 | | * |
21 | | * A stream does not necessarily have such a pendconn. Thus the pendconn is |
22 | | * designated by the stream->pend_pos pointer. This results in some properties : |
23 | | * - pendconn->strm->pend_pos is never NULL for any valid pendconn |
24 | | * - if p->node.node.leaf_p is NULL, the element is unlinked, |
25 | | * otherwise it necessarily belongs to one of the other lists ; this may |
26 | | * not be atomically checked under threads though ; |
27 | | * - pendconn->queue->px is never NULL if pendconn->node.node.leaf_p is not NULL |
28 | | * - pendconn->queue->sv is never NULL if pendconn is in the server's queue, |
29 | | * and is always NULL if it is in the proxy's queue or unlinked. |
30 | | * - pendconn->target is NULL while the element is queued, and points to the |
31 | | * assigned server when the pendconn is picked. |
32 | | * |
33 | | * Threads complicate the design a little bit but rules remain simple : |
34 | | * - the server's queue lock must be held at least when manipulating the |
35 | | * server's queue, which is when adding a pendconn to the queue and when |
36 | | * removing a pendconn from the queue. It protects the queue's integrity. |
37 | | * |
38 | | * - the proxy's queue lock must be held at least when manipulating the |
39 | | * proxy's queue, which is when adding a pendconn to the queue and when |
40 | | * removing a pendconn from the queue. It protects the queue's integrity. |
41 | | * |
42 | | * - both locks are compatible and may be held at the same time. |
43 | | * |
44 | | * - a pendconn_add() is only performed by the stream which will own the |
45 | | * pendconn ; the pendconn is allocated at this moment and returned ; it is |
46 | | * added to either the server or the proxy's queue while holding this |
47 | | * queue's lock. |
48 | | * |
49 | | * - the pendconn is then met by a thread walking over the proxy or server's |
50 | | * queue with the respective lock held. This lock is exclusive and the |
51 | | * pendconn can only appear in one queue so by definition a single thread |
52 | | * may find this pendconn at a time. |
53 | | * |
54 | | * - the pendconn is unlinked either by its own stream upon success/abort/ |
55 | | * free, or by another one offering it its server slot. This is achieved by |
56 | | * pendconn_process_next_strm() under either the server or proxy's lock, |
57 | | * pendconn_redistribute() under the server's lock, or pendconn_unlink() |
58 | | * under either the proxy's or the server's lock depending |
59 | | * on the queue the pendconn is attached to. |
60 | | * |
61 | | * - no single operation except the pendconn initialisation prior to the |
62 | | * insertion are performed without either a queue lock held or the element |
63 | | * being unlinked and visible exclusively to its stream. |
64 | | * |
65 | | * - pendconn_process_next_strm() assign ->target so that the stream knows |
66 | | * what server to work with (via pendconn_dequeue() which sets it on |
67 | | * strm->target). |
68 | | * |
69 | | * - a pendconn doesn't switch between queues, it stays where it is. |
70 | | */ |
71 | | |
72 | | #include <import/eb32tree.h> |
73 | | #include <haproxy/api.h> |
74 | | #include <haproxy/backend.h> |
75 | | #include <haproxy/counters.h> |
76 | | #include <haproxy/http_rules.h> |
77 | | #include <haproxy/pool.h> |
78 | | #include <haproxy/queue.h> |
79 | | #include <haproxy/sample.h> |
80 | | #include <haproxy/server-t.h> |
81 | | #include <haproxy/stream.h> |
82 | | #include <haproxy/task.h> |
83 | | #include <haproxy/tcp_rules.h> |
84 | | #include <haproxy/thread.h> |
85 | | #include <haproxy/time.h> |
86 | | #include <haproxy/tools.h> |
87 | | |
88 | | |
89 | 0 | #define NOW_OFFSET_BOUNDARY() ((now_ms - (TIMER_LOOK_BACK >> 12)) & 0xfffff) |
90 | 0 | #define KEY_CLASS(key) ((u32)key & 0xfff00000) |
91 | 0 | #define KEY_OFFSET(key) ((u32)key & 0x000fffff) |
92 | 0 | #define KEY_CLASS_OFFSET_BOUNDARY(key) (KEY_CLASS(key) | NOW_OFFSET_BOUNDARY()) |
93 | 0 | #define MAKE_KEY(class, offset) (((u32)(class + 0x7ff) << 20) | ((u32)(now_ms + offset) & 0xfffff)) |
94 | | |
95 | | DECLARE_TYPED_POOL(pool_head_pendconn, "pendconn", struct pendconn, 0, 64); |
96 | | |
97 | | /* returns the effective dynamic maxconn for a server, considering the minconn |
98 | | * and the proxy's usage relative to its dynamic connections limit. It is |
99 | | * expected that 0 < s->minconn <= s->maxconn when this is called. If the |
100 | | * server is currently warming up, the slowstart is also applied to the |
101 | | * resulting value, which can be lower than minconn in this case, but never |
102 | | * less than 1. |
103 | | */ |
104 | | unsigned int srv_dynamic_maxconn(const struct server *s) |
105 | 0 | { |
106 | 0 | unsigned int max; |
107 | |
|
108 | 0 | if (s->minconn == s->maxconn || s->proxy->beconn >= s->proxy->fullconn) |
109 | | /* static limit, or no fullconn or proxy is full */ |
110 | 0 | max = s->maxconn; |
111 | 0 | else max = MAX(s->minconn, |
112 | 0 | s->proxy->beconn * s->maxconn / s->proxy->fullconn); |
113 | |
|
114 | 0 | if ((s->cur_state == SRV_ST_STARTING) && |
115 | 0 | ns_to_sec(now_ns) < s->last_change + s->slowstart && |
116 | 0 | ns_to_sec(now_ns) >= s->last_change) { |
117 | 0 | unsigned int ratio; |
118 | 0 | ratio = 100 * (ns_to_sec(now_ns) - s->last_change) / s->slowstart; |
119 | 0 | max = MAX(1, max * ratio / 100); |
120 | 0 | } |
121 | 0 | return max; |
122 | 0 | } |
123 | | |
124 | | /* Remove the pendconn from the server's queue. At this stage, the connection |
125 | | * is not really dequeued. It will be done during the process_stream. It is |
126 | | * up to the caller to atomically decrement the pending counts. |
127 | | * |
128 | | * The caller must own the lock on the server queue. The pendconn must still be |
129 | | * queued (p->node.node.leaf_p != NULL) and must be in a server queue |
130 | | * (p->queue->sv != NULL). |
131 | | */ |
132 | | static void __pendconn_unlink_srv(struct pendconn *p) |
133 | 0 | { |
134 | 0 | p->strm->logs.srv_queue_pos += _HA_ATOMIC_LOAD(&p->queue->idx) - p->queue_idx; |
135 | 0 | eb32_delete(&p->node); |
136 | 0 | } |
137 | | |
138 | | /* Remove the pendconn from the proxy's queue. At this stage, the connection |
139 | | * is not really dequeued. It will be done during the process_stream. It is |
140 | | * up to the caller to atomically decrement the pending counts. |
141 | | * |
142 | | * The caller must own the lock on the proxy queue. The pendconn must still be |
143 | | * queued (p->node.node.leaf_p != NULL) and must be in the proxy queue |
144 | | * (p->queue->sv == NULL). |
145 | | */ |
146 | | static void __pendconn_unlink_prx(struct pendconn *p) |
147 | 0 | { |
148 | 0 | p->strm->logs.prx_queue_pos += _HA_ATOMIC_LOAD(&p->queue->idx) - p->queue_idx; |
149 | 0 | eb32_delete(&p->node); |
150 | 0 | } |
151 | | |
152 | | /* Locks the queue the pendconn element belongs to. This relies on p->queue |
153 | | * being properly initialized (which is always the case once the element |
154 | | * has been added). |
155 | | */ |
156 | | static inline void pendconn_queue_lock(struct pendconn *p) |
157 | 0 | { |
158 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &p->queue->lock); |
159 | 0 | } |
160 | | |
161 | | /* Unlocks the queue the pendconn element belongs to. This relies on p->queue |
162 | | * being properly initialized (which is always the case once the element |
163 | | * has been added). |
164 | | */ |
165 | | static inline void pendconn_queue_unlock(struct pendconn *p) |
166 | 0 | { |
167 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &p->queue->lock); |
168 | 0 | } |
169 | | |
170 | | /* Removes the pendconn from the server/proxy queue. At this stage, the |
171 | | * connection is not really dequeued. It will be done during process_stream(). |
172 | | * This function takes all the required locks for the operation. The pendconn |
173 | | * must be valid, though it doesn't matter if it was already unlinked. Prefer |
174 | | * pendconn_cond_unlink() to first check <p>. It also forces a serialization |
175 | | * on p->del_lock to make sure another thread currently waking it up finishes |
176 | | * first. |
177 | | */ |
178 | | void pendconn_unlink(struct pendconn *p) |
179 | 0 | { |
180 | 0 | struct queue *q = p->queue; |
181 | 0 | struct proxy *px = q->px; |
182 | 0 | struct server *sv = q->sv; |
183 | 0 | uint oldidx; |
184 | 0 | int done = 0; |
185 | |
|
186 | 0 | oldidx = _HA_ATOMIC_LOAD(&p->queue->idx); |
187 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &q->lock); |
188 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &p->del_lock); |
189 | |
|
190 | 0 | if (p->node.node.leaf_p) { |
191 | 0 | eb32_delete(&p->node); |
192 | 0 | done = 1; |
193 | 0 | } |
194 | |
|
195 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &p->del_lock); |
196 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &q->lock); |
197 | |
|
198 | 0 | if (done) { |
199 | 0 | oldidx -= p->queue_idx; |
200 | 0 | if (sv) { |
201 | 0 | p->strm->logs.srv_queue_pos += oldidx; |
202 | 0 | _HA_ATOMIC_DEC(&sv->queueslength); |
203 | 0 | } |
204 | 0 | else { |
205 | 0 | p->strm->logs.prx_queue_pos += oldidx; |
206 | 0 | _HA_ATOMIC_DEC(&px->queueslength); |
207 | 0 | } |
208 | |
|
209 | 0 | _HA_ATOMIC_DEC(&q->length); |
210 | 0 | _HA_ATOMIC_DEC(&px->totpend); |
211 | 0 | } |
212 | 0 | } |
213 | | |
214 | | /* Retrieve the first pendconn from tree <pendconns>. Classes are always |
215 | | * considered first, then the time offset. The time does wrap, so the |
216 | | * lookup is performed twice, one to retrieve the first class and a second |
217 | | * time to retrieve the earliest time in this class. |
218 | | */ |
219 | | static struct pendconn *pendconn_first(struct eb_root *pendconns) |
220 | 0 | { |
221 | 0 | struct eb32_node *node, *node2 = NULL; |
222 | 0 | u32 key; |
223 | |
|
224 | 0 | node = eb32_first(pendconns); |
225 | 0 | if (!node) |
226 | 0 | return NULL; |
227 | | |
228 | 0 | key = KEY_CLASS_OFFSET_BOUNDARY(node->key); |
229 | 0 | node2 = eb32_lookup_ge(pendconns, key); |
230 | |
|
231 | 0 | if (!node2 || |
232 | 0 | KEY_CLASS(node2->key) != KEY_CLASS(node->key)) { |
233 | | /* no other key in the tree, or in this class */ |
234 | 0 | return eb32_entry(node, struct pendconn, node); |
235 | 0 | } |
236 | | |
237 | | /* found a better key */ |
238 | 0 | return eb32_entry(node2, struct pendconn, node); |
239 | 0 | } |
240 | | |
241 | | /* Process the next pending connection from either a server or a proxy, and |
242 | | * returns a strictly positive value on success (see below). If no pending |
243 | | * connection is found, 0 is returned. Note that neither <srv> nor <px> may be |
244 | | * NULL. Priority is given to the oldest request in the queue if both <srv> and |
245 | | * <px> have pending requests. This ensures that no request will be left |
246 | | * unserved. The <px> queue is not considered if the server (or a tracked |
247 | | * server) is not RUNNING, is disabled, or has a null weight (server going |
248 | | * down). The <srv> queue is still considered in this case, because if some |
249 | | * connections remain there, it means that some requests have been forced there |
250 | | * after it was seen down (eg: due to option persist). The stream is |
251 | | * immediately marked as "assigned", and both its <srv> and <srv_conn> are set |
252 | | * to <srv>. |
253 | | * |
254 | | * The proxy's queue will be consulted only if px_ok is non-zero. |
255 | | * |
256 | | * This function must only be called if the server queue is locked _AND_ the |
257 | | * proxy queue is not. Today it is only called by process_srv_queue. |
258 | | * When a pending connection is dequeued, this function returns 1 if a pendconn |
259 | | * is dequeued, otherwise 0. |
260 | | */ |
261 | | static int pendconn_process_next_strm(struct server *srv, struct proxy *px, int px_ok, int tgrp) |
262 | 0 | { |
263 | 0 | struct pendconn *p = NULL; |
264 | 0 | struct pendconn *pp = NULL; |
265 | 0 | u32 pkey, ppkey; |
266 | 0 | int served; |
267 | 0 | int maxconn; |
268 | 0 | int got_it = 0; |
269 | |
|
270 | 0 | p = NULL; |
271 | 0 | if (srv->per_tgrp[tgrp - 1].queue.length) |
272 | 0 | p = pendconn_first(&srv->per_tgrp[tgrp - 1].queue.head); |
273 | |
|
274 | 0 | pp = NULL; |
275 | 0 | if (px_ok && px->per_tgrp[tgrp - 1].queue.length) { |
276 | | /* the lock only remains held as long as the pp is |
277 | | * in the proxy's queue. |
278 | | */ |
279 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &px->per_tgrp[tgrp - 1].queue.lock); |
280 | 0 | pp = pendconn_first(&px->per_tgrp[tgrp - 1].queue.head); |
281 | 0 | if (!pp) |
282 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &px->per_tgrp[tgrp - 1].queue.lock); |
283 | 0 | } |
284 | |
|
285 | 0 | if (!p && !pp) |
286 | 0 | return 0; |
287 | | |
288 | 0 | served = _HA_ATOMIC_LOAD(&srv->served); |
289 | 0 | maxconn = srv_dynamic_maxconn(srv); |
290 | |
|
291 | 0 | while (served < maxconn && !got_it) |
292 | 0 | got_it = _HA_ATOMIC_CAS(&srv->served, &served, served + 1); |
293 | | |
294 | | /* No more slot available, give up */ |
295 | 0 | if (!got_it) { |
296 | 0 | if (pp) |
297 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &px->per_tgrp[tgrp - 1].queue.lock); |
298 | 0 | return 0; |
299 | 0 | } |
300 | | |
301 | | /* |
302 | | * Now we know we'll have something available. |
303 | | * Let's try to allocate a slot on the server. |
304 | | */ |
305 | 0 | if (!pp) |
306 | 0 | goto use_p; /* p != NULL */ |
307 | 0 | else if (!p) |
308 | 0 | goto use_pp; /* pp != NULL */ |
309 | | |
310 | | /* p != NULL && pp != NULL*/ |
311 | | |
312 | 0 | if (KEY_CLASS(p->node.key) < KEY_CLASS(pp->node.key)) |
313 | 0 | goto use_p; |
314 | | |
315 | 0 | if (KEY_CLASS(pp->node.key) < KEY_CLASS(p->node.key)) |
316 | 0 | goto use_pp; |
317 | | |
318 | 0 | pkey = KEY_OFFSET(p->node.key); |
319 | 0 | ppkey = KEY_OFFSET(pp->node.key); |
320 | |
|
321 | 0 | if (pkey < NOW_OFFSET_BOUNDARY()) |
322 | 0 | pkey += 0x100000; // key in the future |
323 | |
|
324 | 0 | if (ppkey < NOW_OFFSET_BOUNDARY()) |
325 | 0 | ppkey += 0x100000; // key in the future |
326 | |
|
327 | 0 | if (pkey <= ppkey) |
328 | 0 | goto use_p; |
329 | | |
330 | 0 | use_pp: |
331 | | /* we'd like to release the proxy lock ASAP to let other threads |
332 | | * work with other servers. But for this we must first hold the |
333 | | * pendconn alive to prevent a removal from its owning stream. |
334 | | */ |
335 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &pp->del_lock); |
336 | | |
337 | | /* now the element won't go, we can release the proxy */ |
338 | 0 | __pendconn_unlink_prx(pp); |
339 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &px->per_tgrp[tgrp - 1].queue.lock); |
340 | |
|
341 | 0 | pp->strm_flags |= SF_ASSIGNED; |
342 | 0 | pp->target = srv; |
343 | 0 | stream_add_srv_conn(pp->strm, srv); |
344 | | |
345 | | /* we must wake the task up before releasing the lock as it's the only |
346 | | * way to make sure the task still exists. The pendconn cannot vanish |
347 | | * under us since the task will need to take the lock anyway and to wait |
348 | | * if it wakes up on a different thread. |
349 | | */ |
350 | 0 | task_wakeup(pp->strm->task, TASK_WOKEN_RES); |
351 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &pp->del_lock); |
352 | |
|
353 | 0 | _HA_ATOMIC_DEC(&px->per_tgrp[tgrp - 1].queue.length); |
354 | 0 | _HA_ATOMIC_INC(&px->per_tgrp[tgrp - 1].queue.idx); |
355 | 0 | _HA_ATOMIC_DEC(&px->queueslength); |
356 | 0 | return 1; |
357 | | |
358 | 0 | use_p: |
359 | | /* we don't need the px queue lock anymore, we have the server's lock */ |
360 | 0 | if (pp) |
361 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &px->per_tgrp[tgrp - 1].queue.lock); |
362 | |
|
363 | 0 | p->strm_flags |= SF_ASSIGNED; |
364 | 0 | p->target = srv; |
365 | 0 | stream_add_srv_conn(p->strm, srv); |
366 | | |
367 | | /* we must wake the task up before releasing the lock as it's the only |
368 | | * way to make sure the task still exists. The pendconn cannot vanish |
369 | | * under us since the task will need to take the lock anyway and to wait |
370 | | * if it wakes up on a different thread. |
371 | | */ |
372 | 0 | task_wakeup(p->strm->task, TASK_WOKEN_RES); |
373 | 0 | __pendconn_unlink_srv(p); |
374 | |
|
375 | 0 | _HA_ATOMIC_DEC(&srv->per_tgrp[tgrp - 1].queue.length); |
376 | 0 | _HA_ATOMIC_INC(&srv->per_tgrp[tgrp - 1].queue.idx); |
377 | 0 | _HA_ATOMIC_DEC(&srv->queueslength); |
378 | 0 | return 1; |
379 | 0 | } |
380 | | |
381 | | /* Manages a server's connection queue. This function will try to dequeue as |
382 | | * many pending streams as possible, and wake them up. |
383 | | */ |
384 | | int process_srv_queue(struct server *s) |
385 | 0 | { |
386 | 0 | struct server *ref = s->track ? s->track : s; |
387 | 0 | struct proxy *p = s->proxy; |
388 | 0 | long non_empty_tgids[(global.nbtgroups / LONGBITS) + 1]; |
389 | 0 | int maxconn; |
390 | 0 | int done = 0; |
391 | 0 | int px_ok; |
392 | 0 | int cur_tgrp; |
393 | 0 | int i = global.nbtgroups; |
394 | 0 | int curgrpnb = i; |
395 | | |
396 | |
|
397 | 0 | while (i >= LONGBITS) { |
398 | 0 | non_empty_tgids[(global.nbtgroups - i) / LONGBITS] = ULONG_MAX; |
399 | 0 | i -= LONGBITS; |
400 | 0 | } |
401 | 0 | while (i > 0) { |
402 | 0 | ha_bit_set(global.nbtgroups - i, non_empty_tgids); |
403 | 0 | i--; |
404 | 0 | } |
405 | | |
406 | | /* if a server is not usable or backup and must not be used |
407 | | * to dequeue backend requests. |
408 | | */ |
409 | 0 | px_ok = srv_currently_usable(ref) && |
410 | 0 | (!(s->flags & SRV_F_BACKUP) || |
411 | 0 | (!p->srv_act && |
412 | 0 | (s == p->lbprm.fbck || (p->options & PR_O_USE_ALL_BK)))); |
413 | | |
414 | | /* let's repeat that under the lock on each round. Threads competing |
415 | | * for the same server will give up, knowing that at least one of |
416 | | * them will check the conditions again before quitting. In order |
417 | | * to avoid the deadly situation where one thread spends its time |
418 | | * dequeueing for others, we limit the number of rounds it does. |
419 | | * However we still re-enter the loop for one pass if there's no |
420 | | * more served, otherwise we could end up with no other thread |
421 | | * trying to dequeue them. |
422 | | * |
423 | | * There's one racy part: we don't want to have more than one thread |
424 | | * in charge of dequeuing, hence the dequeung flag. We cannot rely |
425 | | * on a trylock here because it would compete against pendconn_add() |
426 | | * and would occasionally leave entries in the queue that are never |
427 | | * dequeued. Nobody else uses the dequeuing flag so when seeing it |
428 | | * non-null, we're certain that another thread is waiting on it. |
429 | | * |
430 | | * We'll dequeue MAX_SELF_USE_QUEUE items from the queue corresponding |
431 | | * to our thread group, then we'll get one from a different one, to |
432 | | * be sure those actually get processed too. |
433 | | */ |
434 | 0 | while (curgrpnb != 0 |
435 | 0 | && (done < global.tune.maxpollevents || !s->served) && |
436 | 0 | s->served < (maxconn = srv_dynamic_maxconn(s))) { |
437 | 0 | int self_served; |
438 | 0 | int to_dequeue; |
439 | | |
440 | | /* |
441 | | * self_served contains the number of times we dequeued items |
442 | | * from our own thread-group queue. |
443 | | */ |
444 | 0 | self_served = _HA_ATOMIC_LOAD(&s->per_tgrp[tgid - 1].self_served) % (MAX_SELF_USE_QUEUE + 1); |
445 | 0 | if ((self_served == MAX_SELF_USE_QUEUE && (curgrpnb > 1 || !ha_bit_test(tgid - 1, non_empty_tgids))) || |
446 | 0 | !ha_bit_test(tgid - 1, non_empty_tgids)) { |
447 | 0 | unsigned int old_served, new_served; |
448 | | |
449 | | /* |
450 | | * We want to dequeue from another queue. The last |
451 | | * one we used is stored in last_other_tgrp_served. |
452 | | */ |
453 | 0 | old_served = _HA_ATOMIC_LOAD(&s->per_tgrp[tgid - 1].last_other_tgrp_served); |
454 | 0 | do { |
455 | 0 | new_served = old_served + 1; |
456 | | |
457 | | /* |
458 | | * Find the next tgrp to dequeue from. |
459 | | * If we're here then we know there is |
460 | | * at least one tgrp that is not the current |
461 | | * tgrp that we can dequeue from, so that |
462 | | * loop will end eventually. |
463 | | */ |
464 | 0 | while (new_served == tgid || |
465 | 0 | new_served == global.nbtgroups + 1 || |
466 | 0 | !ha_bit_test(new_served - 1, non_empty_tgids)) { |
467 | 0 | if (new_served == global.nbtgroups + 1) |
468 | 0 | new_served = 1; |
469 | 0 | else |
470 | 0 | new_served++; |
471 | 0 | } |
472 | 0 | } while (!_HA_ATOMIC_CAS(&s->per_tgrp[tgid - 1].last_other_tgrp_served, &old_served, new_served) && __ha_cpu_relax()); |
473 | 0 | cur_tgrp = new_served; |
474 | 0 | to_dequeue = 1; |
475 | 0 | } else { |
476 | 0 | cur_tgrp = tgid; |
477 | 0 | if (self_served == MAX_SELF_USE_QUEUE) |
478 | 0 | self_served = 0; |
479 | 0 | to_dequeue = MAX_SELF_USE_QUEUE - self_served; |
480 | 0 | } |
481 | 0 | if (HA_ATOMIC_XCHG(&s->per_tgrp[cur_tgrp - 1].dequeuing, 1)) { |
482 | 0 | ha_bit_clr(cur_tgrp - 1, non_empty_tgids); |
483 | 0 | curgrpnb--; |
484 | 0 | continue; |
485 | 0 | } |
486 | | |
487 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &s->per_tgrp[cur_tgrp - 1].queue.lock); |
488 | 0 | while (to_dequeue > 0 && s->served < maxconn) { |
489 | | /* |
490 | | * pendconn_process_next_strm() will increment |
491 | | * the served field, only if it is < maxconn. |
492 | | */ |
493 | 0 | if (!pendconn_process_next_strm(s, p, px_ok, cur_tgrp)) { |
494 | 0 | ha_bit_clr(cur_tgrp - 1, non_empty_tgids); |
495 | 0 | curgrpnb--; |
496 | 0 | break; |
497 | 0 | } |
498 | 0 | to_dequeue--; |
499 | 0 | if (cur_tgrp == tgid) |
500 | 0 | _HA_ATOMIC_INC(&s->per_tgrp[tgid - 1].self_served); |
501 | 0 | done++; |
502 | 0 | if (done >= global.tune.maxpollevents) |
503 | 0 | break; |
504 | 0 | } |
505 | 0 | HA_ATOMIC_STORE(&s->per_tgrp[cur_tgrp - 1].dequeuing, 0); |
506 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &s->per_tgrp[cur_tgrp - 1].queue.lock); |
507 | 0 | } |
508 | |
|
509 | 0 | if (done) { |
510 | 0 | _HA_ATOMIC_SUB(&p->totpend, done); |
511 | 0 | _HA_ATOMIC_ADD(&p->served, done); |
512 | 0 | __ha_barrier_atomic_store(); |
513 | 0 | if (p->lbprm.ops && p->lbprm.ops->server_take_conn) |
514 | 0 | p->lbprm.ops->server_take_conn(s); |
515 | 0 | } |
516 | 0 | if (s->served == 0 && p->served == 0 && !HA_ATOMIC_LOAD(&p->ready_srv)) { |
517 | 0 | int i; |
518 | | |
519 | | /* |
520 | | * If there is no task running on the server, and the proxy, |
521 | | * let it known that we are ready, there is a small race |
522 | | * condition if a task was being added just before we checked |
523 | | * the proxy queue. It will look for that server, and use it |
524 | | * if nothing is currently running, as there would be nobody |
525 | | * to wake it up. |
526 | | */ |
527 | 0 | _HA_ATOMIC_STORE(&p->ready_srv, s); |
528 | | /* |
529 | | * Maybe a stream was added to the queue just after we |
530 | | * checked, but before we set ready_srv so it would not see it, |
531 | | * just in case try to run one more stream. |
532 | | */ |
533 | 0 | for (i = 0; i < global.nbtgroups; i++) { |
534 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &s->per_tgrp[i].queue.lock); |
535 | 0 | if (pendconn_process_next_strm(s, p, px_ok, i + 1)) { |
536 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &s->per_tgrp[i].queue.lock); |
537 | 0 | _HA_ATOMIC_SUB(&p->totpend, 1); |
538 | 0 | _HA_ATOMIC_ADD(&p->served, 1); |
539 | 0 | done++; |
540 | 0 | break; |
541 | 0 | } |
542 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &s->per_tgrp[i].queue.lock); |
543 | 0 | } |
544 | 0 | } |
545 | 0 | return done; |
546 | 0 | } |
547 | | |
548 | | /* Adds the stream <strm> to the pending connection queue of server <strm>->srv |
549 | | * or to the one of <strm>->proxy if srv is NULL. All counters and back pointers |
550 | | * are updated accordingly. Returns NULL if no memory is available, otherwise the |
551 | | * pendconn itself. If the stream was already marked as served, its flag is |
552 | | * cleared. It is illegal to call this function with a non-NULL strm->srv_conn. |
553 | | * The stream's queue position is counted with an offset of -1 because we want |
554 | | * to make sure that being at the first position in the queue reports 1. |
555 | | * |
556 | | * The queue is sorted by the composition of the priority_class, and the current |
557 | | * timestamp offset by strm->priority_offset. The timestamp is in milliseconds |
558 | | * and truncated to 20 bits, so will wrap every 17m28s575ms. |
559 | | * The offset can be positive or negative, and an offset of 0 puts it in the |
560 | | * middle of this range (~ 8 min). Note that this also means if the adjusted |
561 | | * timestamp wraps around, the request will be misinterpreted as being of |
562 | | * the highest priority for that priority class. |
563 | | * |
564 | | * This function must be called by the stream itself, so in the context of |
565 | | * process_stream. |
566 | | */ |
567 | | struct pendconn *pendconn_add(struct stream *strm) |
568 | 0 | { |
569 | 0 | struct pendconn *p; |
570 | 0 | struct proxy *px; |
571 | 0 | struct server *srv; |
572 | 0 | struct queue *q; |
573 | 0 | unsigned int *max_ptr; |
574 | 0 | unsigned int *queueslength; |
575 | 0 | unsigned int old_max, new_max; |
576 | |
|
577 | 0 | p = pool_alloc(pool_head_pendconn); |
578 | 0 | if (!p) |
579 | 0 | return NULL; |
580 | | |
581 | 0 | p->target = NULL; |
582 | 0 | p->node.key = MAKE_KEY(strm->priority_class, strm->priority_offset); |
583 | 0 | p->strm = strm; |
584 | 0 | p->strm_flags = strm->flags; |
585 | 0 | HA_SPIN_INIT(&p->del_lock); |
586 | 0 | strm->pend_pos = p; |
587 | |
|
588 | 0 | px = strm->be; |
589 | 0 | if (strm->flags & SF_ASSIGNED) |
590 | 0 | srv = objt_server(strm->target); |
591 | 0 | else |
592 | 0 | srv = NULL; |
593 | |
|
594 | 0 | if (srv) { |
595 | 0 | q = &srv->per_tgrp[tgid - 1].queue; |
596 | 0 | max_ptr = &srv->counters.nbpend_max; |
597 | 0 | queueslength = &srv->queueslength; |
598 | 0 | } |
599 | 0 | else { |
600 | 0 | q = &px->per_tgrp[tgid - 1].queue; |
601 | 0 | max_ptr = &px->be_counters.nbpend_max; |
602 | 0 | queueslength = &px->queueslength; |
603 | 0 | } |
604 | |
|
605 | 0 | p->queue = q; |
606 | 0 | p->queue_idx = _HA_ATOMIC_LOAD(&q->idx) - 1; // for logging only |
607 | 0 | new_max = _HA_ATOMIC_ADD_FETCH(queueslength, 1); |
608 | 0 | _HA_ATOMIC_INC(&q->length); |
609 | 0 | old_max = _HA_ATOMIC_LOAD(max_ptr); |
610 | 0 | while (new_max > old_max) { |
611 | 0 | if (likely(_HA_ATOMIC_CAS(max_ptr, &old_max, new_max))) |
612 | 0 | break; |
613 | 0 | } |
614 | 0 | __ha_barrier_atomic_store(); |
615 | |
|
616 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &q->lock); |
617 | 0 | eb32_insert(&q->head, &p->node); |
618 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &q->lock); |
619 | |
|
620 | 0 | _HA_ATOMIC_INC(&px->totpend); |
621 | 0 | return p; |
622 | 0 | } |
623 | | |
624 | | /* Redistribute pending connections when a server goes down. The number of |
625 | | * connections redistributed is returned. It will take the server queue lock |
626 | | * and does not use nor depend on other locks. |
627 | | */ |
628 | | int pendconn_redistribute(struct server *s) |
629 | 0 | { |
630 | 0 | struct pendconn *p; |
631 | 0 | struct eb32_node *node, *nodeb; |
632 | 0 | struct proxy *px = s->proxy; |
633 | 0 | int px_xferred = 0; |
634 | 0 | int xferred = 0; |
635 | 0 | int i; |
636 | | |
637 | | /* The REDISP option was specified. We will ignore cookie and force to |
638 | | * balance or use the dispatcher. |
639 | | */ |
640 | 0 | if (!(s->cur_admin & SRV_ADMF_MAINT) && |
641 | 0 | (s->proxy->options & (PR_O_REDISP|PR_O_PERSIST)) != PR_O_REDISP) |
642 | 0 | goto skip_srv_queue; |
643 | | |
644 | 0 | for (i = 0; i < global.nbtgroups; i++) { |
645 | 0 | struct queue *queue = &s->per_tgrp[i].queue; |
646 | 0 | int local_xferred = 0; |
647 | |
|
648 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &queue->lock); |
649 | 0 | for (node = eb32_first(&queue->head); node; node = nodeb) { |
650 | 0 | nodeb = eb32_next(node); |
651 | |
|
652 | 0 | p = eb32_entry(node, struct pendconn, node); |
653 | 0 | if (p->strm_flags & SF_FORCE_PRST) |
654 | 0 | continue; |
655 | | |
656 | | /* it's left to the dispatcher to choose a server */ |
657 | 0 | __pendconn_unlink_srv(p); |
658 | 0 | if (!(s->proxy->options & PR_O_REDISP)) |
659 | 0 | p->strm_flags &= ~(SF_DIRECT | SF_ASSIGNED); |
660 | |
|
661 | 0 | task_wakeup(p->strm->task, TASK_WOKEN_RES); |
662 | 0 | local_xferred++; |
663 | 0 | } |
664 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &queue->lock); |
665 | 0 | xferred += local_xferred; |
666 | 0 | if (local_xferred) |
667 | 0 | _HA_ATOMIC_SUB(&queue->length, local_xferred); |
668 | 0 | } |
669 | |
|
670 | 0 | if (xferred) { |
671 | 0 | _HA_ATOMIC_SUB(&s->queueslength, xferred); |
672 | 0 | _HA_ATOMIC_SUB(&s->proxy->totpend, xferred); |
673 | 0 | } |
674 | |
|
675 | 0 | skip_srv_queue: |
676 | 0 | if (px->lbprm.tot_wact || px->lbprm.tot_wbck) |
677 | 0 | goto done; |
678 | | |
679 | 0 | for (i = 0; i < global.nbtgroups; i++) { |
680 | 0 | struct queue *queue = &px->per_tgrp[i].queue; |
681 | 0 | int local_xferred = 0; |
682 | |
|
683 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &queue->lock); |
684 | 0 | for (node = eb32_first(&queue->head); node; node = nodeb) { |
685 | 0 | nodeb = eb32_next(node); |
686 | 0 | p = eb32_entry(node, struct pendconn, node); |
687 | | |
688 | | /* force-persist streams may occasionally appear in the |
689 | | * proxy's queue, and we certainly don't want them here! |
690 | | */ |
691 | 0 | p->strm_flags &= ~SF_FORCE_PRST; |
692 | 0 | __pendconn_unlink_prx(p); |
693 | |
|
694 | 0 | task_wakeup(p->strm->task, TASK_WOKEN_RES); |
695 | 0 | local_xferred++; |
696 | 0 | } |
697 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &queue->lock); |
698 | 0 | if (local_xferred) |
699 | 0 | _HA_ATOMIC_SUB(&queue->length, local_xferred); |
700 | 0 | px_xferred += local_xferred; |
701 | 0 | } |
702 | |
|
703 | 0 | if (px_xferred) { |
704 | 0 | _HA_ATOMIC_SUB(&px->queueslength, px_xferred); |
705 | 0 | _HA_ATOMIC_SUB(&px->totpend, px_xferred); |
706 | 0 | } |
707 | 0 | done: |
708 | 0 | return xferred + px_xferred; |
709 | 0 | } |
710 | | |
711 | | /* Try to dequeue pending connection attached to the stream <strm>. It must |
712 | | * always exists here. If the pendconn is still linked to the server or the |
713 | | * proxy queue, nothing is done and the function returns 1. Otherwise, |
714 | | * <strm>->flags and <strm>->target are updated, the pendconn is released and 0 |
715 | | * is returned. |
716 | | * |
717 | | * This function must be called by the stream itself, so in the context of |
718 | | * process_stream. |
719 | | */ |
720 | | int pendconn_dequeue(struct stream *strm) |
721 | 0 | { |
722 | 0 | struct pendconn *p; |
723 | 0 | int is_unlinked; |
724 | | |
725 | | /* unexpected case because it is called by the stream itself and |
726 | | * only the stream can release a pendconn. So it is only |
727 | | * possible if a pendconn is released by someone else or if the |
728 | | * stream is supposed to be queued but without its associated |
729 | | * pendconn. In both cases it is a bug! */ |
730 | 0 | BUG_ON(!strm->pend_pos); |
731 | |
|
732 | 0 | p = strm->pend_pos; |
733 | | |
734 | | /* note below : we need to grab the queue's lock to check for emptiness |
735 | | * because we don't want a partial process_srv_queue() or redistribute() |
736 | | * to be called in parallel and show an empty list without having the |
737 | | * time to finish. With this we know that if we see the element |
738 | | * unlinked, these functions were completely done. |
739 | | */ |
740 | 0 | pendconn_queue_lock(p); |
741 | 0 | is_unlinked = !p->node.node.leaf_p; |
742 | 0 | pendconn_queue_unlock(p); |
743 | | |
744 | | /* serialize to make sure the element was finished processing */ |
745 | 0 | HA_SPIN_LOCK(QUEUE_LOCK, &p->del_lock); |
746 | 0 | HA_SPIN_UNLOCK(QUEUE_LOCK, &p->del_lock); |
747 | |
|
748 | 0 | if (!is_unlinked) |
749 | 0 | return 1; |
750 | | |
751 | | /* the pendconn is not queued anymore and will not be so we're safe |
752 | | * to proceed. |
753 | | */ |
754 | 0 | strm->flags &= ~(SF_DIRECT | SF_ASSIGNED); |
755 | 0 | strm->flags |= p->strm_flags & (SF_DIRECT | SF_ASSIGNED); |
756 | | |
757 | | /* the entry might have been redistributed to another server */ |
758 | 0 | if (!(strm->flags & SF_ASSIGNED)) |
759 | 0 | sockaddr_free(&strm->scb->dst); |
760 | |
|
761 | 0 | if (p->target) { |
762 | | /* a server picked this pendconn, it must skip LB */ |
763 | 0 | stream_set_srv_target(strm, p->target); |
764 | 0 | strm->flags |= SF_ASSIGNED; |
765 | 0 | } |
766 | |
|
767 | 0 | strm->pend_pos = NULL; |
768 | 0 | pool_free(pool_head_pendconn, p); |
769 | 0 | return 0; |
770 | 0 | } |
771 | | |
772 | | static enum act_return action_set_priority_class(struct act_rule *rule, struct proxy *px, |
773 | | struct session *sess, struct stream *s, int flags) |
774 | 0 | { |
775 | 0 | struct sample *smp; |
776 | |
|
777 | 0 | smp = sample_fetch_as_type(px, sess, s, SMP_OPT_DIR_REQ|SMP_OPT_FINAL, rule->arg.expr, SMP_T_SINT); |
778 | 0 | if (!smp) |
779 | 0 | return ACT_RET_CONT; |
780 | | |
781 | 0 | s->priority_class = queue_limit_class(smp->data.u.sint); |
782 | 0 | return ACT_RET_CONT; |
783 | 0 | } |
784 | | |
785 | | static enum act_return action_set_priority_offset(struct act_rule *rule, struct proxy *px, |
786 | | struct session *sess, struct stream *s, int flags) |
787 | 0 | { |
788 | 0 | struct sample *smp; |
789 | |
|
790 | 0 | smp = sample_fetch_as_type(px, sess, s, SMP_OPT_DIR_REQ|SMP_OPT_FINAL, rule->arg.expr, SMP_T_SINT); |
791 | 0 | if (!smp) |
792 | 0 | return ACT_RET_CONT; |
793 | | |
794 | 0 | s->priority_offset = queue_limit_offset(smp->data.u.sint); |
795 | |
|
796 | 0 | return ACT_RET_CONT; |
797 | 0 | } |
798 | | |
799 | | static enum act_parse_ret parse_set_priority_class(const char **args, int *arg, struct proxy *px, |
800 | | struct act_rule *rule, char **err) |
801 | 0 | { |
802 | 0 | unsigned int where = 0; |
803 | |
|
804 | 0 | rule->arg.expr = sample_parse_expr((char **)args, arg, px->conf.args.file, |
805 | 0 | px->conf.args.line, err, &px->conf.args, NULL); |
806 | 0 | if (!rule->arg.expr) |
807 | 0 | return ACT_RET_PRS_ERR; |
808 | | |
809 | 0 | if (px->cap & PR_CAP_FE) |
810 | 0 | where |= SMP_VAL_FE_HRQ_HDR; |
811 | 0 | if (px->cap & PR_CAP_BE) |
812 | 0 | where |= SMP_VAL_BE_HRQ_HDR; |
813 | |
|
814 | 0 | if (!(rule->arg.expr->fetch->val & where)) { |
815 | 0 | memprintf(err, |
816 | 0 | "fetch method '%s' extracts information from '%s', none of which is available here", |
817 | 0 | args[0], sample_src_names(rule->arg.expr->fetch->use)); |
818 | 0 | free(rule->arg.expr); |
819 | 0 | return ACT_RET_PRS_ERR; |
820 | 0 | } |
821 | | |
822 | 0 | rule->action = ACT_CUSTOM; |
823 | 0 | rule->action_ptr = action_set_priority_class; |
824 | 0 | return ACT_RET_PRS_OK; |
825 | 0 | } |
826 | | |
827 | | static enum act_parse_ret parse_set_priority_offset(const char **args, int *arg, struct proxy *px, |
828 | | struct act_rule *rule, char **err) |
829 | 0 | { |
830 | 0 | unsigned int where = 0; |
831 | |
|
832 | 0 | rule->arg.expr = sample_parse_expr((char **)args, arg, px->conf.args.file, |
833 | 0 | px->conf.args.line, err, &px->conf.args, NULL); |
834 | 0 | if (!rule->arg.expr) |
835 | 0 | return ACT_RET_PRS_ERR; |
836 | | |
837 | 0 | if (px->cap & PR_CAP_FE) |
838 | 0 | where |= SMP_VAL_FE_HRQ_HDR; |
839 | 0 | if (px->cap & PR_CAP_BE) |
840 | 0 | where |= SMP_VAL_BE_HRQ_HDR; |
841 | |
|
842 | 0 | if (!(rule->arg.expr->fetch->val & where)) { |
843 | 0 | memprintf(err, |
844 | 0 | "fetch method '%s' extracts information from '%s', none of which is available here", |
845 | 0 | args[0], sample_src_names(rule->arg.expr->fetch->use)); |
846 | 0 | free(rule->arg.expr); |
847 | 0 | return ACT_RET_PRS_ERR; |
848 | 0 | } |
849 | | |
850 | 0 | rule->action = ACT_CUSTOM; |
851 | 0 | rule->action_ptr = action_set_priority_offset; |
852 | 0 | return ACT_RET_PRS_OK; |
853 | 0 | } |
854 | | |
855 | | static struct action_kw_list tcp_cont_kws = {ILH, { |
856 | | { "set-priority-class", parse_set_priority_class }, |
857 | | { "set-priority-offset", parse_set_priority_offset }, |
858 | | { /* END */ } |
859 | | }}; |
860 | | |
861 | | INITCALL1(STG_REGISTER, tcp_req_cont_keywords_register, &tcp_cont_kws); |
862 | | |
863 | | static struct action_kw_list http_req_kws = {ILH, { |
864 | | { "set-priority-class", parse_set_priority_class }, |
865 | | { "set-priority-offset", parse_set_priority_offset }, |
866 | | { /* END */ } |
867 | | }}; |
868 | | |
869 | | INITCALL1(STG_REGISTER, http_req_keywords_register, &http_req_kws); |
870 | | |
871 | | static int |
872 | | smp_fetch_priority_class(const struct arg *args, struct sample *smp, const char *kw, void *private) |
873 | 0 | { |
874 | 0 | if (!smp->strm) |
875 | 0 | return 0; |
876 | | |
877 | 0 | smp->data.type = SMP_T_SINT; |
878 | 0 | smp->data.u.sint = smp->strm->priority_class; |
879 | |
|
880 | 0 | return 1; |
881 | 0 | } |
882 | | |
883 | | static int |
884 | | smp_fetch_priority_offset(const struct arg *args, struct sample *smp, const char *kw, void *private) |
885 | 0 | { |
886 | 0 | if (!smp->strm) |
887 | 0 | return 0; |
888 | | |
889 | 0 | smp->data.type = SMP_T_SINT; |
890 | 0 | smp->data.u.sint = smp->strm->priority_offset; |
891 | |
|
892 | 0 | return 1; |
893 | 0 | } |
894 | | |
895 | | |
896 | | static struct sample_fetch_kw_list smp_kws = {ILH, { |
897 | | { "prio_class", smp_fetch_priority_class, 0, NULL, SMP_T_SINT, SMP_USE_INTRN, }, |
898 | | { "prio_offset", smp_fetch_priority_offset, 0, NULL, SMP_T_SINT, SMP_USE_INTRN, }, |
899 | | { /* END */}, |
900 | | }}; |
901 | | |
902 | | INITCALL1(STG_REGISTER, sample_register_fetches, &smp_kws); |
903 | | |
904 | | /* |
905 | | * Local variables: |
906 | | * c-indent-level: 8 |
907 | | * c-basic-offset: 8 |
908 | | * End: |
909 | | */ |