vinyl-cache/bin/vinyld/waiter/cache_waiter_poll.c
0
/*-
1
 * Copyright (c) 2006 Verdens Gang AS
2
 * Copyright (c) 2006-2011 Varnish Software AS
3
 * All rights reserved.
4
 *
5
 * Author: Poul-Henning Kamp <phk@phk.freebsd.dk>
6
 *
7
 * SPDX-License-Identifier: BSD-2-Clause
8
 *
9
 * Redistribution and use in source and binary forms, with or without
10
 * modification, are permitted provided that the following conditions
11
 * are met:
12
 * 1. Redistributions of source code must retain the above copyright
13
 *    notice, this list of conditions and the following disclaimer.
14
 * 2. Redistributions in binary form must reproduce the above copyright
15
 *    notice, this list of conditions and the following disclaimer in the
16
 *    documentation and/or other materials provided with the distribution.
17
 *
18
 * THIS SOFTWARE IS PROVIDED BY THE AUTHOR AND CONTRIBUTORS ``AS IS'' AND
19
 * ANY EXPRESS OR IMPLIED WARRANTIES, INCLUDING, BUT NOT LIMITED TO, THE
20
 * IMPLIED WARRANTIES OF MERCHANTABILITY AND FITNESS FOR A PARTICULAR PURPOSE
21
 * ARE DISCLAIMED.  IN NO EVENT SHALL AUTHOR OR CONTRIBUTORS BE LIABLE
22
 * FOR ANY DIRECT, INDIRECT, INCIDENTAL, SPECIAL, EXEMPLARY, OR CONSEQUENTIAL
23
 * DAMAGES (INCLUDING, BUT NOT LIMITED TO, PROCUREMENT OF SUBSTITUTE GOODS
24
 * OR SERVICES; LOSS OF USE, DATA, OR PROFITS; OR BUSINESS INTERRUPTION)
25
 * HOWEVER CAUSED AND ON ANY THEORY OF LIABILITY, WHETHER IN CONTRACT, STRICT
26
 * LIABILITY, OR TORT (INCLUDING NEGLIGENCE OR OTHERWISE) ARISING IN ANY WAY
27
 * OUT OF THE USE OF THIS SOFTWARE, EVEN IF ADVISED OF THE POSSIBILITY OF
28
 * SUCH DAMAGE.
29
 *
30
 */
31
32
#include "config.h"
33
34
#include <fcntl.h>
35
#include <poll.h>
36
#include <stdlib.h>
37
38
#include "cache/cache_int.h"
39
40
#include "waiter/waiter.h"
41
#include "waiter/waiter_priv.h"
42
#include "vtim.h"
43
44
struct vwp {
45
        unsigned                magic;
46
#define VWP_MAGIC               0x4b2cc735
47
        struct waiter           *waiter;
48
49
        int                     pipes[2];
50
51
        pthread_t               thread;
52
        struct pollfd           *pollfd;
53
        struct waited           **idx;
54
        size_t                  npoll;
55
        size_t                  hpoll;
56
};
57
58
/*--------------------------------------------------------------------
59
 * It would make much more sense to not use two large vectors, but
60
 * the poll(2) API forces us to use at least one, so ... KISS.
61
 */
62
63
static void
64 210
vwp_extend_pollspace(struct vwp *vwp)
65
{
66
        size_t inc;
67
68 210
        if (vwp->npoll < (1<<12))
69 210
                inc = (1<<10);
70 0
        else if (vwp->npoll < (1<<14))
71 0
                inc = (1<<12);
72 0
        else if (vwp->npoll < (1<<16))
73 0
                inc = (1<<14);
74
        else
75 0
                inc = (1<<16);
76
77 420
        VSL(SLT_Debug, NO_VXID, "Acceptor poll space increased by %zu to %zu",
78 210
            inc, vwp->npoll + inc);
79
80 420
        vwp->pollfd = realloc(vwp->pollfd,
81 210
            (vwp->npoll + inc) * sizeof(*vwp->pollfd));
82 210
        AN(vwp->pollfd);
83 210
        memset(vwp->pollfd + vwp->npoll, 0, inc * sizeof(*vwp->pollfd));
84
85 210
        vwp->idx = realloc(vwp->idx, (vwp->npoll + inc) * sizeof(*vwp->idx));
86 210
        AN(vwp->idx);
87 210
        memset(vwp->idx + vwp->npoll, 0, inc * sizeof(*vwp->idx));
88
89 215250
        for (; inc > 0; inc--)
90 215040
                vwp->pollfd[vwp->npoll++].fd = -1;
91 210
}
92
93
/*--------------------------------------------------------------------*/
94
95
static void
96 2450
vwp_add(struct vwp *vwp, struct waited *wp)
97
{
98
99 2450
        CHECK_OBJ_NOTNULL(wp, WAITED_MAGIC);
100 2450
        VSL(SLT_Debug, NO_VXID, "vwp: ADD %d", wp->fd);
101 2450
        CHECK_OBJ_NOTNULL(vwp, VWP_MAGIC);
102 2450
        if (vwp->hpoll == vwp->npoll)
103 0
                vwp_extend_pollspace(vwp);
104 2450
        assert(vwp->hpoll < vwp->npoll);
105 2450
        assert(vwp->pollfd[vwp->hpoll].fd == -1);
106 2450
        AZ(vwp->idx[vwp->hpoll]);
107 2450
        vwp->pollfd[vwp->hpoll].fd = wp->fd;
108 2450
        vwp->pollfd[vwp->hpoll].events = POLLIN;
109 2450
        vwp->idx[vwp->hpoll] = wp;
110 2450
        vwp->hpoll++;
111 2450
        Wait_HeapInsert(vwp->waiter, wp);
112 2450
}
113
114
static void
115 2325
vwp_del(struct vwp *vwp, int n)
116
{
117 2325
        vwp->hpoll--;
118 2325
        if (n != vwp->hpoll) {
119 324
                vwp->pollfd[n] = vwp->pollfd[vwp->hpoll];
120 324
                vwp->idx[n] = vwp->idx[vwp->hpoll];
121 324
        }
122 2325
        memset(&vwp->pollfd[vwp->hpoll], 0, sizeof(*vwp->pollfd));
123 2325
        vwp->pollfd[vwp->hpoll].fd = -1;
124 2325
        vwp->idx[vwp->hpoll] = NULL;
125 2325
}
126
127
/*--------------------------------------------------------------------*/
128
129
static void
130 2660
vwp_dopipe(struct vwp *vwp)
131
{
132
        struct waited *w[128];
133
        ssize_t ss;
134
        int i;
135
136 2660
        ss = read(vwp->pipes[0], w, sizeof w);
137 2660
        assert(ss > 0);
138 2660
        i = 0;
139 5110
        while (ss) {
140 2660
                if (w[i] == NULL) {
141 210
                        assert(ss == sizeof w[0]);
142 210
                        pthread_exit(NULL);
143
                }
144 2450
                CHECK_OBJ_NOTNULL(w[i], WAITED_MAGIC);
145 2450
                assert(w[i]->fd > 0);                   // no stdin
146 2450
                vwp_add(vwp, w[i++]);
147 2450
                ss -= sizeof w[0];
148
        }
149 2450
}
150
151
/*--------------------------------------------------------------------*/
152
153
static void *
154 0
vwp_main(void *priv)
155
{
156
        int t, v;
157
        struct vwp *vwp;
158
        struct waiter *w;
159
        struct waited *wp;
160
        double now, then;
161
        size_t z;
162
163 0
        THR_SetName("cache-poll");
164 0
        THR_Init();
165 0
        CAST_OBJ_NOTNULL(vwp, priv, VWP_MAGIC);
166 0
        w = vwp->waiter;
167
168 4680
        while (1) {
169 4680
                then = Wait_HeapDue(w, &wp);
170 4680
                if (wp == NULL)
171 1396
                        t = -1;
172
                else {
173 3284
                        t = (int)floor(1e3 * (then - VTIM_real()));
174 3284
                        if (t < 0)
175 0
                                t = 0;
176
                }
177 4680
                assert(vwp->hpoll > 0);
178 4680
                AN(vwp->pollfd);
179 4680
                v = poll(vwp->pollfd, vwp->hpoll, t);
180 4680
                assert(v >= 0);
181 4680
                now = VTIM_real();
182 4680
                if (vwp->pollfd[0].revents)
183 2660
                        v--;
184 7807
                for (z = 1; z < vwp->hpoll;) {
185 4410
                        assert(vwp->pollfd[z].fd != vwp->pipes[0]);
186 4410
                        wp = vwp->idx[z];
187 4410
                        CHECK_OBJ_NOTNULL(wp, WAITED_MAGIC);
188
189 4410
                        if (v == 0 && Wait_HeapDue(w, NULL) > now)
190 1283
                                break;
191 3127
                        if (vwp->pollfd[z].revents)
192 2324
                                v--;
193 3127
                        then = Wait_When(wp);
194 3127
                        if (then <= now) {
195 0
                                AN(Wait_HeapDelete(w, wp));
196 0
                                Wait_Call(w, wp, WAITER_TIMEOUT, now);
197 0
                                vwp_del(vwp, z);
198 3127
                        } else if (vwp->pollfd[z].revents & POLLIN) {
199 2325
                                assert(wp->fd > 0);
200 2325
                                assert(wp->fd == vwp->pollfd[z].fd);
201 2325
                                AN(Wait_HeapDelete(w, wp));
202 2325
                                Wait_Call(w, wp, WAITER_ACTION, now);
203 2325
                                vwp_del(vwp, z);
204 2325
                        } else {
205 802
                                z++;
206
                        }
207
                }
208
                // vwp_dopipe calls pthread_exit()
209 4680
                if (vwp->pollfd[0].revents)
210 2660
                        vwp_dopipe(vwp);
211
        }
212
        NEEDLESS(return (NULL));
213
}
214
215
/*--------------------------------------------------------------------*/
216
217
static int v_matchproto_(waiter_enter_f)
218 2450
vwp_enter(void *priv, struct waited *wp)
219
{
220
        struct vwp *vwp;
221
222 2450
        CAST_OBJ_NOTNULL(vwp, priv, VWP_MAGIC);
223
224 2450
        if (write(vwp->pipes[1], &wp, sizeof wp) != sizeof wp)
225 0
                return (-1);
226 2450
        return (0);
227 2450
}
228
229
/*--------------------------------------------------------------------*/
230
231
static void v_matchproto_(waiter_init_f)
232 210
vwp_init(struct waiter *w)
233
{
234
        struct vwp *vwp;
235
236 210
        CHECK_OBJ_NOTNULL(w, WAITER_MAGIC);
237 210
        vwp = w->priv;
238 210
        INIT_OBJ(vwp, VWP_MAGIC);
239 210
        vwp->waiter = w;
240 210
        AZ(pipe(vwp->pipes));
241
        // XXX: set write pipe non-blocking
242
243 210
        vwp->hpoll = 1;
244 210
        vwp_extend_pollspace(vwp);
245 210
        vwp->pollfd[0].fd = vwp->pipes[0];
246 210
        vwp->pollfd[0].events = POLLIN;
247 210
        PTOK(pthread_create(&vwp->thread, NULL, vwp_main, vwp));
248 210
}
249
250
/*--------------------------------------------------------------------
251
 * It is the callers responsibility to trigger all fd's waited on to
252
 * fail somehow.
253
 */
254
255
static void v_matchproto_(waiter_fini_f)
256 210
vwp_fini(struct waiter *w)
257
{
258
        struct vwp *vwp;
259
        void *vp;
260
261 210
        CAST_OBJ_NOTNULL(vwp, w->priv, VWP_MAGIC);
262 210
        vp = NULL;
263
        // XXX: set write pipe blocking
264 210
        assert(write(vwp->pipes[1], &vp, sizeof vp) == sizeof vp);
265 210
        PTOK(pthread_join(vwp->thread, &vp));
266 210
        closefd(&vwp->pipes[0]);
267 210
        closefd(&vwp->pipes[1]);
268 210
        free(vwp->pollfd);
269 210
        free(vwp->idx);
270 210
}
271
272
/*--------------------------------------------------------------------*/
273
274
#include "waiter/mgt_waiter.h"
275
276
const struct waiter_impl waiter_poll = {
277
        .name =         "poll",
278
        .init =         vwp_init,
279
        .fini =         vwp_fini,
280
        .enter =        vwp_enter,
281
        .size =         sizeof(struct vwp),
282
};