Coverage Report

Created: 2026-09-28 06:55

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/postgres/src/backend/libpq/pqcomm.c
Line
Count
Source
1
/*-------------------------------------------------------------------------
2
 *
3
 * pqcomm.c
4
 *    Communication functions between the Frontend and the Backend
5
 *
6
 * These routines handle the low-level details of communication between
7
 * frontend and backend.  They just shove data across the communication
8
 * channel, and are ignorant of the semantics of the data.
9
 *
10
 * To emit an outgoing message, use the routines in pqformat.c to construct
11
 * the message in a buffer and then emit it in one call to pq_putmessage.
12
 * There are no functions to send raw bytes or partial messages; this
13
 * ensures that the channel will not be clogged by an incomplete message if
14
 * execution is aborted by ereport(ERROR) partway through the message.
15
 *
16
 * At one time, libpq was shared between frontend and backend, but now
17
 * the backend's "backend/libpq" is quite separate from "interfaces/libpq".
18
 * All that remains is similarities of names to trap the unwary...
19
 *
20
 * Portions Copyright (c) 1996-2026, PostgreSQL Global Development Group
21
 * Portions Copyright (c) 1994, Regents of the University of California
22
 *
23
 *  src/backend/libpq/pqcomm.c
24
 *
25
 *-------------------------------------------------------------------------
26
 */
27
28
/*------------------------
29
 * INTERFACE ROUTINES
30
 *
31
 * setup/teardown:
32
 *    ListenServerPort  - Open postmaster's server port
33
 *    AcceptConnection  - Accept new connection with client
34
 *    TouchSocketFiles  - Protect socket files against /tmp cleaners
35
 *    pq_init       - initialize libpq at backend startup
36
 *    socket_comm_reset - reset libpq during error recovery
37
 *    socket_close    - shutdown libpq at backend exit
38
 *
39
 * low-level I/O:
40
 *    pq_getbytes   - get a known number of bytes from connection
41
 *    pq_getmessage - get a message with length word from connection
42
 *    pq_getbyte    - get next byte from connection
43
 *    pq_peekbyte   - peek at next byte from connection
44
 *    pq_flush    - flush pending output
45
 *    pq_flush_if_writable - flush pending output if writable without blocking
46
 *    pq_getbyte_if_available - get a byte if available without blocking
47
 *
48
 * message-level I/O
49
 *    pq_putmessage - send a normal message (suppressed in COPY OUT mode)
50
 *    pq_putmessage_noblock - buffer a normal message (suppressed in COPY OUT)
51
 *
52
 *------------------------
53
 */
54
#include "postgres.h"
55
56
#include <signal.h>
57
#include <fcntl.h>
58
#include <grp.h>
59
#include <unistd.h>
60
#include <sys/file.h>
61
#include <sys/socket.h>
62
#include <sys/stat.h>
63
#include <sys/time.h>
64
#include <netdb.h>
65
#include <netinet/in.h>
66
#include <netinet/tcp.h>
67
#include <utime.h>
68
#ifdef WIN32
69
#include <mstcpip.h>
70
#endif
71
72
#include "common/ip.h"
73
#include "libpq/libpq.h"
74
#include "miscadmin.h"
75
#include "port/pg_bswap.h"
76
#include "postmaster/postmaster.h"
77
#include "storage/ipc.h"
78
#include "storage/latch.h"
79
#include "utils/guc_hooks.h"
80
#include "utils/memutils.h"
81
82
/*
83
 * Cope with the various platform-specific ways to spell TCP keepalive socket
84
 * options.  This doesn't cover Windows, which as usual does its own thing.
85
 */
86
#if defined(TCP_KEEPIDLE)
87
/* TCP_KEEPIDLE is the name of this option on Linux and *BSD */
88
0
#define PG_TCP_KEEPALIVE_IDLE TCP_KEEPIDLE
89
#define PG_TCP_KEEPALIVE_IDLE_STR "TCP_KEEPIDLE"
90
#elif defined(TCP_KEEPALIVE_THRESHOLD)
91
/* TCP_KEEPALIVE_THRESHOLD is the name of this option on Solaris >= 11 */
92
#define PG_TCP_KEEPALIVE_IDLE TCP_KEEPALIVE_THRESHOLD
93
#define PG_TCP_KEEPALIVE_IDLE_STR "TCP_KEEPALIVE_THRESHOLD"
94
#elif defined(TCP_KEEPALIVE) && defined(__darwin__)
95
/* TCP_KEEPALIVE is the name of this option on macOS */
96
/* Caution: Solaris has this symbol but it means something different */
97
#define PG_TCP_KEEPALIVE_IDLE TCP_KEEPALIVE
98
#define PG_TCP_KEEPALIVE_IDLE_STR "TCP_KEEPALIVE"
99
#endif
100
101
/*
102
 * Configuration options
103
 */
104
int     Unix_socket_permissions;
105
char     *Unix_socket_group;
106
107
/* Where the Unix socket files are (list of palloc'd strings) */
108
static List *sock_paths = NIL;
109
110
/*
111
 * Buffers for low-level I/O.
112
 *
113
 * The receive buffer is fixed size. Send buffer is usually 8k, but can be
114
 * enlarged by pq_putmessage_noblock() if the message doesn't fit otherwise.
115
 */
116
117
0
#define PQ_SEND_BUFFER_SIZE 8192
118
0
#define PQ_RECV_BUFFER_SIZE 8192
119
120
static char *PqSendBuffer;
121
static int  PqSendBufferSize; /* Size send buffer */
122
static size_t PqSendPointer;  /* Next index to store a byte in PqSendBuffer */
123
static size_t PqSendStart;    /* Next index to send a byte in PqSendBuffer */
124
125
static char PqRecvBuffer[PQ_RECV_BUFFER_SIZE];
126
static int  PqRecvPointer;    /* Next index to read a byte from PqRecvBuffer */
127
static int  PqRecvLength;   /* End of data available in PqRecvBuffer */
128
129
/*
130
 * Message status
131
 */
132
static bool PqCommBusy;     /* busy sending data to the client */
133
static bool PqCommReadingMsg; /* in the middle of reading a message */
134
135
136
/* Internal functions */
137
static void socket_comm_reset(void);
138
static void socket_close(int code, Datum arg);
139
static void socket_set_nonblocking(bool nonblocking);
140
static int  socket_flush(void);
141
static int  socket_flush_if_writable(void);
142
static bool socket_is_send_pending(void);
143
static int  socket_putmessage(char msgtype, const char *s, size_t len);
144
static void socket_putmessage_noblock(char msgtype, const char *s, size_t len);
145
static inline int internal_putbytes(const void *b, size_t len);
146
static inline int internal_flush(void);
147
static pg_noinline int internal_flush_buffer(const char *buf, size_t *start,
148
                       size_t *end);
149
150
static int  Lock_AF_UNIX(const char *unixSocketDir, const char *unixSocketPath);
151
static int  Setup_AF_UNIX(const char *sock_path);
152
153
static const PQcommMethods PqCommSocketMethods = {
154
  .comm_reset = socket_comm_reset,
155
  .flush = socket_flush,
156
  .flush_if_writable = socket_flush_if_writable,
157
  .is_send_pending = socket_is_send_pending,
158
  .putmessage = socket_putmessage,
159
  .putmessage_noblock = socket_putmessage_noblock
160
};
161
162
const PQcommMethods *PqCommMethods = &PqCommSocketMethods;
163
164
WaitEventSet *FeBeWaitSet;
165
166
167
/* --------------------------------
168
 *    pq_init - initialize libpq at backend startup
169
 * --------------------------------
170
 */
171
Port *
172
pq_init(ClientSocket *client_sock)
173
0
{
174
0
  Port     *port;
175
0
  int     socket_pos PG_USED_FOR_ASSERTS_ONLY;
176
0
  int     latch_pos PG_USED_FOR_ASSERTS_ONLY;
177
178
  /* allocate the Port struct and copy the ClientSocket contents to it */
179
0
  port = palloc0_object(Port);
180
0
  port->sock = client_sock->sock;
181
0
  memcpy(&port->raddr.addr, &client_sock->raddr.addr, client_sock->raddr.salen);
182
0
  port->raddr.salen = client_sock->raddr.salen;
183
184
  /* fill in the server (local) address */
185
0
  port->laddr.salen = sizeof(port->laddr.addr);
186
0
  if (getsockname(port->sock,
187
0
          (struct sockaddr *) &port->laddr.addr,
188
0
          &port->laddr.salen) < 0)
189
0
  {
190
0
    ereport(FATAL,
191
0
        (errmsg("%s() failed: %m", "getsockname")));
192
0
  }
193
194
  /* select NODELAY and KEEPALIVE options if it's a TCP connection */
195
0
  if (port->laddr.addr.ss_family != AF_UNIX)
196
0
  {
197
0
    int     on;
198
#ifdef WIN32
199
    int     oldopt;
200
    int     optlen;
201
    int     newopt;
202
#endif
203
204
0
#ifdef  TCP_NODELAY
205
0
    on = 1;
206
0
    if (setsockopt(port->sock, IPPROTO_TCP, TCP_NODELAY,
207
0
             (char *) &on, sizeof(on)) < 0)
208
0
    {
209
0
      ereport(FATAL,
210
0
          (errmsg("%s(%s) failed: %m", "setsockopt", "TCP_NODELAY")));
211
0
    }
212
0
#endif
213
0
    on = 1;
214
0
    if (setsockopt(port->sock, SOL_SOCKET, SO_KEEPALIVE,
215
0
             (char *) &on, sizeof(on)) < 0)
216
0
    {
217
0
      ereport(FATAL,
218
0
          (errmsg("%s(%s) failed: %m", "setsockopt", "SO_KEEPALIVE")));
219
0
    }
220
221
#ifdef WIN32
222
223
    /*
224
     * This is a Win32 socket optimization.  The OS send buffer should be
225
     * large enough to send the whole Postgres send buffer in one go, or
226
     * performance suffers.  The Postgres send buffer can be enlarged if a
227
     * very large message needs to be sent, but we won't attempt to
228
     * enlarge the OS buffer if that happens, so somewhat arbitrarily
229
     * ensure that the OS buffer is at least PQ_SEND_BUFFER_SIZE * 4.
230
     * (That's 32kB with the current default).
231
     *
232
     * The default OS buffer size used to be 8kB in earlier Windows
233
     * versions, but was raised to 64kB in Windows 2012.  So it shouldn't
234
     * be necessary to change it in later versions anymore.  Changing it
235
     * unnecessarily can even reduce performance, because setting
236
     * SO_SNDBUF in the application disables the "dynamic send buffering"
237
     * feature that was introduced in Windows 7.  So before fiddling with
238
     * SO_SNDBUF, check if the current buffer size is already large enough
239
     * and only increase it if necessary.
240
     *
241
     * See https://support.microsoft.com/kb/823764/EN-US/ and
242
     * https://msdn.microsoft.com/en-us/library/bb736549%28v=vs.85%29.aspx
243
     */
244
    optlen = sizeof(oldopt);
245
    if (getsockopt(port->sock, SOL_SOCKET, SO_SNDBUF, (char *) &oldopt,
246
             &optlen) < 0)
247
    {
248
      ereport(FATAL,
249
          (errmsg("%s(%s) failed: %m", "getsockopt", "SO_SNDBUF")));
250
    }
251
    newopt = PQ_SEND_BUFFER_SIZE * 4;
252
    if (oldopt < newopt)
253
    {
254
      if (setsockopt(port->sock, SOL_SOCKET, SO_SNDBUF, (char *) &newopt,
255
               sizeof(newopt)) < 0)
256
      {
257
        ereport(FATAL,
258
            (errmsg("%s(%s) failed: %m", "setsockopt", "SO_SNDBUF")));
259
      }
260
    }
261
#endif
262
263
    /*
264
     * Also apply the current keepalive parameters.  If we fail to set a
265
     * parameter, don't error out, because these aren't universally
266
     * supported.  (Note: you might think we need to reset the GUC
267
     * variables to 0 in such a case, but it's not necessary because the
268
     * show hooks for these variables report the truth anyway.)
269
     */
270
0
    (void) pq_setkeepalivesidle(tcp_keepalives_idle, port);
271
0
    (void) pq_setkeepalivesinterval(tcp_keepalives_interval, port);
272
0
    (void) pq_setkeepalivescount(tcp_keepalives_count, port);
273
0
    (void) pq_settcpusertimeout(tcp_user_timeout, port);
274
0
  }
275
276
  /* initialize state variables */
277
0
  PqSendBufferSize = PQ_SEND_BUFFER_SIZE;
278
0
  PqSendBuffer = MemoryContextAlloc(TopMemoryContext, PqSendBufferSize);
279
0
  PqSendPointer = PqSendStart = PqRecvPointer = PqRecvLength = 0;
280
0
  PqCommBusy = false;
281
0
  PqCommReadingMsg = false;
282
283
  /* set up process-exit hook to close the socket */
284
0
  on_proc_exit(socket_close, 0);
285
286
  /*
287
   * In backends (as soon as forked) we operate the underlying socket in
288
   * nonblocking mode and use latches to implement blocking semantics if
289
   * needed. That allows us to provide safely interruptible reads and
290
   * writes.
291
   */
292
0
#ifndef WIN32
293
0
  if (!pg_set_noblock(port->sock))
294
0
    ereport(FATAL,
295
0
        (errmsg("could not set socket to nonblocking mode: %m")));
296
0
#endif
297
298
0
#ifndef WIN32
299
300
  /* Don't give the socket to any subprograms we execute. */
301
0
  if (fcntl(port->sock, F_SETFD, FD_CLOEXEC) < 0)
302
0
    elog(FATAL, "fcntl(F_SETFD) failed on socket: %m");
303
0
#endif
304
305
0
  FeBeWaitSet = CreateWaitEventSet(NULL, FeBeWaitSetNEvents);
306
0
  socket_pos = AddWaitEventToSet(FeBeWaitSet, WL_SOCKET_WRITEABLE,
307
0
                   port->sock, NULL, NULL);
308
0
  latch_pos = AddWaitEventToSet(FeBeWaitSet, WL_LATCH_SET, PGINVALID_SOCKET,
309
0
                  MyLatch, NULL);
310
0
  AddWaitEventToSet(FeBeWaitSet, WL_POSTMASTER_DEATH, PGINVALID_SOCKET,
311
0
            NULL, NULL);
312
313
  /*
314
   * The event positions match the order we added them, but let's sanity
315
   * check them to be sure.
316
   */
317
0
  Assert(socket_pos == FeBeWaitSetSocketPos);
318
0
  Assert(latch_pos == FeBeWaitSetLatchPos);
319
320
0
  return port;
321
0
}
322
323
/* --------------------------------
324
 *    socket_comm_reset - reset libpq during error recovery
325
 *
326
 * This is called from error recovery at the outer idle loop.  It's
327
 * just to get us out of trouble if we somehow manage to elog() from
328
 * inside a pqcomm.c routine (which ideally will never happen, but...)
329
 * --------------------------------
330
 */
331
static void
332
socket_comm_reset(void)
333
5.70k
{
334
  /* Do not throw away pending data, but do reset the busy flag */
335
5.70k
  PqCommBusy = false;
336
5.70k
}
337
338
/* --------------------------------
339
 *    socket_close - shutdown libpq at backend exit
340
 *
341
 * This is the one pg_on_exit_callback in place during BackendInitialize().
342
 * That function's unusual signal handling constrains that this callback be
343
 * safe to run at any instant.
344
 * --------------------------------
345
 */
346
static void
347
socket_close(int code, Datum arg)
348
0
{
349
  /* Nothing to do in a standalone backend, where MyProcPort is NULL. */
350
0
  if (MyProcPort != NULL)
351
0
  {
352
#ifdef ENABLE_GSS
353
    /*
354
     * Shutdown GSSAPI layer.  This section does nothing when interrupting
355
     * BackendInitialize(), because pg_GSS_recvauth() makes first use of
356
     * "ctx" and "cred".
357
     *
358
     * Note that we don't bother to free MyProcPort->gss, since we're
359
     * about to exit anyway.
360
     */
361
    if (MyProcPort->gss)
362
    {
363
      OM_uint32 min_s;
364
365
      if (MyProcPort->gss->ctx != GSS_C_NO_CONTEXT)
366
        gss_delete_sec_context(&min_s, &MyProcPort->gss->ctx, NULL);
367
368
      if (MyProcPort->gss->cred != GSS_C_NO_CREDENTIAL)
369
        gss_release_cred(&min_s, &MyProcPort->gss->cred);
370
    }
371
#endif              /* ENABLE_GSS */
372
373
    /*
374
     * Cleanly shut down SSL layer.  Nowhere else does a postmaster child
375
     * call this, so this is safe when interrupting BackendInitialize().
376
     */
377
0
    secure_close(MyProcPort);
378
379
    /*
380
     * Formerly we did an explicit close() here, but it seems better to
381
     * leave the socket open until the process dies.  This allows clients
382
     * to perform a "synchronous close" if they care --- wait till the
383
     * transport layer reports connection closure, and you can be sure the
384
     * backend has exited.
385
     *
386
     * We do set sock to PGINVALID_SOCKET to prevent any further I/O,
387
     * though.
388
     */
389
0
    MyProcPort->sock = PGINVALID_SOCKET;
390
0
  }
391
0
}
392
393
394
395
/* --------------------------------
396
 * Postmaster functions to handle sockets.
397
 * --------------------------------
398
 */
399
400
/*
401
 * ListenServerPort -- open a "listening" port to accept connections.
402
 *
403
 * family should be AF_UNIX or AF_UNSPEC; portNumber is the port number.
404
 * For AF_UNIX ports, hostName should be NULL and unixSocketDir must be
405
 * specified.  For TCP ports, hostName is either NULL for all interfaces or
406
 * the interface to listen on, and unixSocketDir is ignored (can be NULL).
407
 *
408
 * Successfully opened sockets are appended to the ListenSockets[] array.  On
409
 * entry, *NumListenSockets holds the number of elements currently in the
410
 * array, and it is updated to reflect the opened sockets.  MaxListen is the
411
 * allocated size of the array.
412
 *
413
 * RETURNS: STATUS_OK or STATUS_ERROR
414
 */
415
int
416
ListenServerPort(int family, const char *hostName, unsigned short portNumber,
417
         const char *unixSocketDir,
418
         pgsocket ListenSockets[], int *NumListenSockets, int MaxListen)
419
0
{
420
0
  pgsocket  fd;
421
0
  int     err;
422
0
  int     maxconn;
423
0
  int     ret;
424
0
  char    portNumberStr[32];
425
0
  const char *familyDesc;
426
0
  char    familyDescBuf[64];
427
0
  const char *addrDesc;
428
0
  char    addrBuf[NI_MAXHOST];
429
0
  char     *service;
430
0
  struct addrinfo *addrs = NULL,
431
0
         *addr;
432
0
  struct addrinfo hint;
433
0
  int     added = 0;
434
0
  char    unixSocketPath[MAXPGPATH];
435
0
#if !defined(WIN32) || defined(IPV6_V6ONLY)
436
0
  int     one = 1;
437
0
#endif
438
439
  /* Initialize hint structure */
440
0
  MemSet(&hint, 0, sizeof(hint));
441
0
  hint.ai_family = family;
442
0
  hint.ai_flags = AI_PASSIVE;
443
0
  hint.ai_socktype = SOCK_STREAM;
444
445
0
  if (family == AF_UNIX)
446
0
  {
447
    /*
448
     * Create unixSocketPath from portNumber and unixSocketDir and lock
449
     * that file path
450
     */
451
0
    UNIXSOCK_PATH(unixSocketPath, portNumber, unixSocketDir);
452
0
    if (strlen(unixSocketPath) >= UNIXSOCK_PATH_BUFLEN)
453
0
    {
454
0
      ereport(LOG,
455
0
          (errmsg("Unix-domain socket path \"%s\" is too long (maximum %zu bytes)",
456
0
              unixSocketPath,
457
0
              (UNIXSOCK_PATH_BUFLEN - 1))));
458
0
      return STATUS_ERROR;
459
0
    }
460
0
    if (Lock_AF_UNIX(unixSocketDir, unixSocketPath) != STATUS_OK)
461
0
      return STATUS_ERROR;
462
0
    service = unixSocketPath;
463
0
  }
464
0
  else
465
0
  {
466
0
    snprintf(portNumberStr, sizeof(portNumberStr), "%d", portNumber);
467
0
    service = portNumberStr;
468
0
  }
469
470
0
  ret = pg_getaddrinfo_all(hostName, service, &hint, &addrs);
471
0
  if (ret || !addrs)
472
0
  {
473
0
    if (hostName)
474
0
      ereport(LOG,
475
0
          (errmsg("could not translate host name \"%s\", service \"%s\" to address: %s",
476
0
              hostName, service, gai_strerror(ret))));
477
0
    else
478
0
      ereport(LOG,
479
0
          (errmsg("could not translate service \"%s\" to address: %s",
480
0
              service, gai_strerror(ret))));
481
0
    if (addrs)
482
0
      pg_freeaddrinfo_all(hint.ai_family, addrs);
483
0
    return STATUS_ERROR;
484
0
  }
485
486
0
  for (addr = addrs; addr; addr = addr->ai_next)
487
0
  {
488
0
    if (family != AF_UNIX && addr->ai_family == AF_UNIX)
489
0
    {
490
      /*
491
       * Only set up a unix domain socket when they really asked for it.
492
       * The service/port is different in that case.
493
       */
494
0
      continue;
495
0
    }
496
497
    /* See if there is still room to add 1 more socket. */
498
0
    if (*NumListenSockets == MaxListen)
499
0
    {
500
0
      ereport(LOG,
501
0
          (errmsg("could not bind to all requested addresses: MAXLISTEN (%d) exceeded",
502
0
              MaxListen)));
503
0
      break;
504
0
    }
505
506
    /* set up address family name for log messages */
507
0
    switch (addr->ai_family)
508
0
    {
509
0
      case AF_INET:
510
0
        familyDesc = _("IPv4");
511
0
        break;
512
0
      case AF_INET6:
513
0
        familyDesc = _("IPv6");
514
0
        break;
515
0
      case AF_UNIX:
516
0
        familyDesc = _("Unix");
517
0
        break;
518
0
      default:
519
0
        snprintf(familyDescBuf, sizeof(familyDescBuf),
520
0
             _("unrecognized address family %d"),
521
0
             addr->ai_family);
522
0
        familyDesc = familyDescBuf;
523
0
        break;
524
0
    }
525
526
    /* set up text form of address for log messages */
527
0
    if (addr->ai_family == AF_UNIX)
528
0
      addrDesc = unixSocketPath;
529
0
    else
530
0
    {
531
0
      pg_getnameinfo_all((const struct sockaddr_storage *) addr->ai_addr,
532
0
                 addr->ai_addrlen,
533
0
                 addrBuf, sizeof(addrBuf),
534
0
                 NULL, 0,
535
0
                 NI_NUMERICHOST);
536
0
      addrDesc = addrBuf;
537
0
    }
538
539
0
    if ((fd = socket(addr->ai_family, SOCK_STREAM, 0)) == PGINVALID_SOCKET)
540
0
    {
541
0
      ereport(LOG,
542
0
          (errcode_for_socket_access(),
543
      /* translator: first %s is IPv4, IPv6, or Unix */
544
0
           errmsg("could not create %s socket for address \"%s\": %m",
545
0
              familyDesc, addrDesc)));
546
0
      continue;
547
0
    }
548
549
0
#ifndef WIN32
550
    /* Don't give the listen socket to any subprograms we execute. */
551
0
    if (fcntl(fd, F_SETFD, FD_CLOEXEC) < 0)
552
0
      elog(FATAL, "fcntl(F_SETFD) failed on socket: %m");
553
554
    /*
555
     * Without the SO_REUSEADDR flag, a new postmaster can't be started
556
     * right away after a stop or crash, giving "address already in use"
557
     * error on TCP ports.
558
     *
559
     * On win32, however, this behavior only happens if the
560
     * SO_EXCLUSIVEADDRUSE is set. With SO_REUSEADDR, win32 allows
561
     * multiple servers to listen on the same address, resulting in
562
     * unpredictable behavior. With no flags at all, win32 behaves as Unix
563
     * with SO_REUSEADDR.
564
     */
565
0
    if (addr->ai_family != AF_UNIX)
566
0
    {
567
0
      if ((setsockopt(fd, SOL_SOCKET, SO_REUSEADDR,
568
0
              (char *) &one, sizeof(one))) == -1)
569
0
      {
570
0
        ereport(LOG,
571
0
            (errcode_for_socket_access(),
572
        /* translator: third %s is IPv4 or IPv6 */
573
0
             errmsg("%s(%s) failed for %s address \"%s\": %m",
574
0
                "setsockopt", "SO_REUSEADDR",
575
0
                familyDesc, addrDesc)));
576
0
        closesocket(fd);
577
0
        continue;
578
0
      }
579
0
    }
580
0
#endif
581
582
0
#ifdef IPV6_V6ONLY
583
0
    if (addr->ai_family == AF_INET6)
584
0
    {
585
0
      if (setsockopt(fd, IPPROTO_IPV6, IPV6_V6ONLY,
586
0
               (char *) &one, sizeof(one)) == -1)
587
0
      {
588
0
        ereport(LOG,
589
0
            (errcode_for_socket_access(),
590
        /* translator: third %s is IPv6 */
591
0
             errmsg("%s(%s) failed for %s address \"%s\": %m",
592
0
                "setsockopt", "IPV6_V6ONLY",
593
0
                familyDesc, addrDesc)));
594
0
        closesocket(fd);
595
0
        continue;
596
0
      }
597
0
    }
598
0
#endif
599
600
    /*
601
     * Note: This might fail on some OS's, like Linux older than
602
     * 2.4.21-pre3, that don't have the IPV6_V6ONLY socket option, and map
603
     * ipv4 addresses to ipv6.  It will show ::ffff:ipv4 for all ipv4
604
     * connections.
605
     */
606
0
    err = bind(fd, addr->ai_addr, addr->ai_addrlen);
607
0
    if (err < 0)
608
0
    {
609
0
      int     saved_errno = errno;
610
611
0
      ereport(LOG,
612
0
          (errcode_for_socket_access(),
613
      /* translator: first %s is IPv4, IPv6, or Unix */
614
0
           errmsg("could not bind %s address \"%s\": %m",
615
0
              familyDesc, addrDesc),
616
0
           saved_errno == EADDRINUSE ?
617
0
           (addr->ai_family == AF_UNIX ?
618
0
            errhint("Is another postmaster already running on port %d?",
619
0
                portNumber) :
620
0
            errhint("Is another postmaster already running on port %d?"
621
0
                " If not, wait a few seconds and retry.",
622
0
                portNumber)) : 0));
623
0
      closesocket(fd);
624
0
      continue;
625
0
    }
626
627
0
    if (addr->ai_family == AF_UNIX)
628
0
    {
629
0
      if (Setup_AF_UNIX(service) != STATUS_OK)
630
0
      {
631
0
        closesocket(fd);
632
0
        break;
633
0
      }
634
0
    }
635
636
    /*
637
     * Select appropriate accept-queue length limit.  It seems reasonable
638
     * to use a value similar to the maximum number of child processes
639
     * that the postmaster will permit.
640
     */
641
0
    maxconn = MaxConnections * 2;
642
643
0
    err = listen(fd, maxconn);
644
0
    if (err < 0)
645
0
    {
646
0
      ereport(LOG,
647
0
          (errcode_for_socket_access(),
648
      /* translator: first %s is IPv4, IPv6, or Unix */
649
0
           errmsg("could not listen on %s address \"%s\": %m",
650
0
              familyDesc, addrDesc)));
651
0
      closesocket(fd);
652
0
      continue;
653
0
    }
654
655
0
    if (addr->ai_family == AF_UNIX)
656
0
      ereport(LOG,
657
0
          (errmsg("listening on Unix socket \"%s\"",
658
0
              addrDesc)));
659
0
    else
660
0
      ereport(LOG,
661
      /* translator: first %s is IPv4 or IPv6 */
662
0
          (errmsg("listening on %s address \"%s\", port %d",
663
0
              familyDesc, addrDesc, portNumber)));
664
665
0
    ListenSockets[*NumListenSockets] = fd;
666
0
    (*NumListenSockets)++;
667
0
    added++;
668
0
  }
669
670
0
  pg_freeaddrinfo_all(hint.ai_family, addrs);
671
672
0
  if (!added)
673
0
    return STATUS_ERROR;
674
675
0
  return STATUS_OK;
676
0
}
677
678
679
/*
680
 * Lock_AF_UNIX -- configure unix socket file path
681
 */
682
static int
683
Lock_AF_UNIX(const char *unixSocketDir, const char *unixSocketPath)
684
0
{
685
  /* no lock file for abstract sockets */
686
0
  if (unixSocketPath[0] == '@')
687
0
    return STATUS_OK;
688
689
  /*
690
   * Grab an interlock file associated with the socket file.
691
   *
692
   * Note: there are two reasons for using a socket lock file, rather than
693
   * trying to interlock directly on the socket itself.  First, it's a lot
694
   * more portable, and second, it lets us remove any pre-existing socket
695
   * file without race conditions.
696
   */
697
0
  CreateSocketLockFile(unixSocketPath, true, unixSocketDir);
698
699
  /*
700
   * Once we have the interlock, we can safely delete any pre-existing
701
   * socket file to avoid failure at bind() time.
702
   */
703
0
  (void) unlink(unixSocketPath);
704
705
  /*
706
   * Remember socket file pathnames for later maintenance.
707
   */
708
0
  sock_paths = lappend(sock_paths, pstrdup(unixSocketPath));
709
710
0
  return STATUS_OK;
711
0
}
712
713
714
/*
715
 * Setup_AF_UNIX -- configure unix socket permissions
716
 */
717
static int
718
Setup_AF_UNIX(const char *sock_path)
719
{
720
  /* no file system permissions for abstract sockets */
721
  if (sock_path[0] == '@')
722
    return STATUS_OK;
723
724
  /*
725
   * Fix socket ownership/permission if requested.  Note we must do this
726
   * before we listen() to avoid a window where unwanted connections could
727
   * get accepted.
728
   */
729
  Assert(Unix_socket_group);
730
  if (Unix_socket_group[0] != '\0')
731
  {
732
#ifdef WIN32
733
    elog(WARNING, "configuration item \"unix_socket_group\" is not supported on this platform");
734
#else
735
    char     *endptr;
736
    unsigned long val;
737
    gid_t   gid;
738
739
    val = strtoul(Unix_socket_group, &endptr, 10);
740
    if (*endptr == '\0')
741
    {           /* numeric group id */
742
      gid = val;
743
    }
744
    else
745
    {           /* convert group name to id */
746
      struct group *gr;
747
748
      gr = getgrnam(Unix_socket_group);
749
      if (!gr)
750
      {
751
        ereport(LOG,
752
            (errmsg("group \"%s\" does not exist",
753
                Unix_socket_group)));
754
        return STATUS_ERROR;
755
      }
756
      gid = gr->gr_gid;
757
    }
758
    if (chown(sock_path, -1, gid) == -1)
759
    {
760
      ereport(LOG,
761
          (errcode_for_file_access(),
762
           errmsg("could not set group of file \"%s\": %m",
763
              sock_path)));
764
      return STATUS_ERROR;
765
    }
766
#endif
767
  }
768
769
  if (chmod(sock_path, Unix_socket_permissions) == -1)
770
  {
771
    ereport(LOG,
772
        (errcode_for_file_access(),
773
         errmsg("could not set permissions of file \"%s\": %m",
774
            sock_path)));
775
    return STATUS_ERROR;
776
  }
777
  return STATUS_OK;
778
}
779
780
781
/*
782
 * AcceptConnection -- accept a new connection with client using
783
 *    server port.  Fills *client_sock with the FD and endpoint info
784
 *    of the new connection.
785
 *
786
 * ASSUME: that this doesn't need to be non-blocking because
787
 *    the Postmaster waits for the socket to be ready to accept().
788
 *
789
 * RETURNS: STATUS_OK or STATUS_ERROR
790
 */
791
int
792
AcceptConnection(pgsocket server_fd, ClientSocket *client_sock)
793
{
794
  /* accept connection and fill in the client (remote) address */
795
  client_sock->raddr.salen = sizeof(client_sock->raddr.addr);
796
  if ((client_sock->sock = accept(server_fd,
797
                  (struct sockaddr *) &client_sock->raddr.addr,
798
                  &client_sock->raddr.salen)) == PGINVALID_SOCKET)
799
  {
800
    ereport(LOG,
801
        (errcode_for_socket_access(),
802
         errmsg("could not accept new connection: %m")));
803
804
    /*
805
     * If accept() fails then postmaster.c will still see the server
806
     * socket as read-ready, and will immediately try again.  To avoid
807
     * uselessly sucking lots of CPU, delay a bit before trying again.
808
     * (The most likely reason for failure is being out of kernel file
809
     * table slots; we can do little except hope some will get freed up.)
810
     */
811
    pg_usleep(100000L);   /* wait 0.1 sec */
812
    return STATUS_ERROR;
813
  }
814
815
  return STATUS_OK;
816
}
817
818
/*
819
 * TouchSocketFiles -- mark socket files as recently accessed
820
 *
821
 * This routine should be called every so often to ensure that the socket
822
 * files have a recent mod date (ordinary operations on sockets usually won't
823
 * change the mod date).  That saves them from being removed by
824
 * overenthusiastic /tmp-directory-cleaner daemons.  (Another reason we should
825
 * never have put the socket file in /tmp...)
826
 */
827
void
828
TouchSocketFiles(void)
829
0
{
830
0
  ListCell   *l;
831
832
  /* Loop through all created sockets... */
833
0
  foreach(l, sock_paths)
834
0
  {
835
0
    char     *sock_path = (char *) lfirst(l);
836
837
    /* Ignore errors; there's no point in complaining */
838
0
    (void) utime(sock_path, NULL);
839
0
  }
840
0
}
841
842
/*
843
 * RemoveSocketFiles -- unlink socket files at postmaster shutdown
844
 */
845
void
846
RemoveSocketFiles(void)
847
0
{
848
0
  ListCell   *l;
849
850
  /* Loop through all created sockets... */
851
0
  foreach(l, sock_paths)
852
0
  {
853
0
    char     *sock_path = (char *) lfirst(l);
854
855
    /* Ignore any error. */
856
0
    (void) unlink(sock_path);
857
0
  }
858
  /* Since we're about to exit, no need to reclaim storage */
859
0
}
860
861
862
/* --------------------------------
863
 * Low-level I/O routines begin here.
864
 *
865
 * These routines communicate with a frontend client across a connection
866
 * already established by the preceding routines.
867
 * --------------------------------
868
 */
869
870
/* --------------------------------
871
 *        socket_set_nonblocking - set socket blocking/non-blocking
872
 *
873
 * Sets the socket non-blocking if nonblocking is true, or sets it
874
 * blocking otherwise.
875
 * --------------------------------
876
 */
877
static void
878
socket_set_nonblocking(bool nonblocking)
879
0
{
880
0
  if (MyProcPort == NULL)
881
0
    ereport(ERROR,
882
0
        (errcode(ERRCODE_CONNECTION_DOES_NOT_EXIST),
883
0
         errmsg("there is no client connection")));
884
885
0
  MyProcPort->noblock = nonblocking;
886
0
}
887
888
/* --------------------------------
889
 *    pq_recvbuf - load some bytes into the input buffer
890
 *
891
 *    returns 0 if OK, EOF if trouble
892
 * --------------------------------
893
 */
894
static int
895
pq_recvbuf(void)
896
0
{
897
0
  if (PqRecvPointer > 0)
898
0
  {
899
0
    if (PqRecvLength > PqRecvPointer)
900
0
    {
901
      /* still some unread data, left-justify it in the buffer */
902
0
      memmove(PqRecvBuffer, PqRecvBuffer + PqRecvPointer,
903
0
          PqRecvLength - PqRecvPointer);
904
0
      PqRecvLength -= PqRecvPointer;
905
0
      PqRecvPointer = 0;
906
0
    }
907
0
    else
908
0
      PqRecvLength = PqRecvPointer = 0;
909
0
  }
910
911
  /* Ensure that we're in blocking mode */
912
0
  socket_set_nonblocking(false);
913
914
  /* Can fill buffer from PqRecvLength and upwards */
915
0
  for (;;)
916
0
  {
917
0
    ssize_t   r;
918
919
0
    errno = 0;
920
921
0
    r = secure_read(MyProcPort, PqRecvBuffer + PqRecvLength,
922
0
            PQ_RECV_BUFFER_SIZE - PqRecvLength);
923
924
0
    if (r < 0)
925
0
    {
926
0
      if (errno == EINTR)
927
0
        continue;   /* Ok if interrupted */
928
929
      /*
930
       * Careful: an ereport() that tries to write to the client would
931
       * cause recursion to here, leading to stack overflow and core
932
       * dump!  This message must go *only* to the postmaster log.
933
       *
934
       * If errno is zero, assume it's EOF and let the caller complain.
935
       */
936
0
      if (errno != 0)
937
0
        ereport(COMMERROR,
938
0
            (errcode_for_socket_access(),
939
0
             errmsg("could not receive data from client: %m")));
940
0
      return EOF;
941
0
    }
942
0
    if (r == 0)
943
0
    {
944
      /*
945
       * EOF detected.  We used to write a log message here, but it's
946
       * better to expect the ultimate caller to do that.
947
       */
948
0
      return EOF;
949
0
    }
950
    /* r contains number of bytes read, so just incr length */
951
0
    PqRecvLength += r;
952
0
    return 0;
953
0
  }
954
0
}
955
956
/* --------------------------------
957
 *    pq_getbyte  - get a single byte from connection, or return EOF
958
 * --------------------------------
959
 */
960
int
961
pq_getbyte(void)
962
0
{
963
0
  Assert(PqCommReadingMsg);
964
965
0
  while (PqRecvPointer >= PqRecvLength)
966
0
  {
967
0
    if (pq_recvbuf())   /* If nothing in buffer, then recv some */
968
0
      return EOF;     /* Failed to recv data */
969
0
  }
970
0
  return (unsigned char) PqRecvBuffer[PqRecvPointer++];
971
0
}
972
973
/* --------------------------------
974
 *    pq_peekbyte   - peek at next byte from connection
975
 *
976
 *   Same as pq_getbyte() except we don't advance the pointer.
977
 * --------------------------------
978
 */
979
int
980
pq_peekbyte(void)
981
0
{
982
0
  Assert(PqCommReadingMsg);
983
984
0
  while (PqRecvPointer >= PqRecvLength)
985
0
  {
986
0
    if (pq_recvbuf())   /* If nothing in buffer, then recv some */
987
0
      return EOF;     /* Failed to recv data */
988
0
  }
989
0
  return (unsigned char) PqRecvBuffer[PqRecvPointer];
990
0
}
991
992
/* --------------------------------
993
 *    pq_getbyte_if_available - get a single byte from connection,
994
 *      if available
995
 *
996
 * The received byte is stored in *c. Returns 1 if a byte was read,
997
 * 0 if no data was available, or EOF if trouble.
998
 * --------------------------------
999
 */
1000
int
1001
pq_getbyte_if_available(unsigned char *c)
1002
0
{
1003
0
  ssize_t   r;
1004
1005
0
  Assert(PqCommReadingMsg);
1006
1007
0
  if (PqRecvPointer < PqRecvLength)
1008
0
  {
1009
0
    *c = PqRecvBuffer[PqRecvPointer++];
1010
0
    return 1;
1011
0
  }
1012
1013
  /* Put the socket into non-blocking mode */
1014
0
  socket_set_nonblocking(true);
1015
1016
0
  errno = 0;
1017
1018
0
  r = secure_read(MyProcPort, c, 1);
1019
0
  if (r < 0)
1020
0
  {
1021
    /*
1022
     * Ok if no data available without blocking or interrupted (though
1023
     * EINTR really shouldn't happen with a non-blocking socket). Report
1024
     * other errors.
1025
     */
1026
0
    if (errno == EAGAIN || errno == EWOULDBLOCK || errno == EINTR)
1027
0
      r = 0;
1028
0
    else
1029
0
    {
1030
      /*
1031
       * Careful: an ereport() that tries to write to the client would
1032
       * cause recursion to here, leading to stack overflow and core
1033
       * dump!  This message must go *only* to the postmaster log.
1034
       *
1035
       * If errno is zero, assume it's EOF and let the caller complain.
1036
       */
1037
0
      if (errno != 0)
1038
0
        ereport(COMMERROR,
1039
0
            (errcode_for_socket_access(),
1040
0
             errmsg("could not receive data from client: %m")));
1041
0
      r = EOF;
1042
0
    }
1043
0
  }
1044
0
  else if (r == 0)
1045
0
  {
1046
    /* EOF detected */
1047
0
    r = EOF;
1048
0
  }
1049
1050
0
  return r;
1051
0
}
1052
1053
/* --------------------------------
1054
 *    pq_getbytes   - get a known number of bytes from connection
1055
 *
1056
 *    returns 0 if OK, EOF if trouble
1057
 * --------------------------------
1058
 */
1059
int
1060
pq_getbytes(void *b, size_t len)
1061
0
{
1062
0
  char     *s = b;
1063
0
  size_t    amount;
1064
1065
0
  Assert(PqCommReadingMsg);
1066
1067
0
  while (len > 0)
1068
0
  {
1069
0
    while (PqRecvPointer >= PqRecvLength)
1070
0
    {
1071
0
      if (pq_recvbuf()) /* If nothing in buffer, then recv some */
1072
0
        return EOF;   /* Failed to recv data */
1073
0
    }
1074
0
    amount = PqRecvLength - PqRecvPointer;
1075
0
    if (amount > len)
1076
0
      amount = len;
1077
0
    memcpy(s, PqRecvBuffer + PqRecvPointer, amount);
1078
0
    PqRecvPointer += amount;
1079
0
    s += amount;
1080
0
    len -= amount;
1081
0
  }
1082
0
  return 0;
1083
0
}
1084
1085
/* --------------------------------
1086
 *    pq_discardbytes   - throw away a known number of bytes
1087
 *
1088
 *    same as pq_getbytes except we do not copy the data to anyplace.
1089
 *    this is used for resynchronizing after read errors.
1090
 *
1091
 *    returns 0 if OK, EOF if trouble
1092
 * --------------------------------
1093
 */
1094
static int
1095
pq_discardbytes(size_t len)
1096
0
{
1097
0
  size_t    amount;
1098
1099
0
  Assert(PqCommReadingMsg);
1100
1101
0
  while (len > 0)
1102
0
  {
1103
0
    while (PqRecvPointer >= PqRecvLength)
1104
0
    {
1105
0
      if (pq_recvbuf()) /* If nothing in buffer, then recv some */
1106
0
        return EOF;   /* Failed to recv data */
1107
0
    }
1108
0
    amount = PqRecvLength - PqRecvPointer;
1109
0
    if (amount > len)
1110
0
      amount = len;
1111
0
    PqRecvPointer += amount;
1112
0
    len -= amount;
1113
0
  }
1114
0
  return 0;
1115
0
}
1116
1117
/* --------------------------------
1118
 *    pq_buffer_remaining_data  - return number of bytes in receive buffer
1119
 *
1120
 * This will *not* attempt to read more data. And reading up to that number of
1121
 * bytes should not cause reading any more data either.
1122
 * --------------------------------
1123
 */
1124
ssize_t
1125
pq_buffer_remaining_data(void)
1126
0
{
1127
0
  Assert(PqRecvLength >= PqRecvPointer);
1128
0
  return (PqRecvLength - PqRecvPointer);
1129
0
}
1130
1131
1132
/* --------------------------------
1133
 *    pq_startmsgread - begin reading a message from the client.
1134
 *
1135
 *    This must be called before any of the pq_get* functions.
1136
 * --------------------------------
1137
 */
1138
void
1139
pq_startmsgread(void)
1140
0
{
1141
  /*
1142
   * There shouldn't be a read active already, but let's check just to be
1143
   * sure.
1144
   */
1145
0
  if (PqCommReadingMsg)
1146
0
    ereport(FATAL,
1147
0
        (errcode(ERRCODE_PROTOCOL_VIOLATION),
1148
0
         errmsg("terminating connection because protocol synchronization was lost")));
1149
1150
0
  PqCommReadingMsg = true;
1151
0
}
1152
1153
1154
/* --------------------------------
1155
 *    pq_endmsgread - finish reading message.
1156
 *
1157
 *    This must be called after reading a message with pq_getbytes()
1158
 *    and friends, to indicate that we have read the whole message.
1159
 *    pq_getmessage() does this implicitly.
1160
 * --------------------------------
1161
 */
1162
void
1163
pq_endmsgread(void)
1164
0
{
1165
0
  Assert(PqCommReadingMsg);
1166
1167
0
  PqCommReadingMsg = false;
1168
0
}
1169
1170
/* --------------------------------
1171
 *    pq_is_reading_msg - are we currently reading a message?
1172
 *
1173
 * This is used in error recovery at the outer idle loop to detect if we have
1174
 * lost protocol sync, and need to terminate the connection. pq_startmsgread()
1175
 * will check for that too, but it's nicer to detect it earlier.
1176
 * --------------------------------
1177
 */
1178
bool
1179
pq_is_reading_msg(void)
1180
0
{
1181
0
  return PqCommReadingMsg;
1182
0
}
1183
1184
/* --------------------------------
1185
 *    pq_getmessage - get a message with length word from connection
1186
 *
1187
 *    The return value is placed in an expansible StringInfo, which has
1188
 *    already been initialized by the caller.
1189
 *    Only the message body is placed in the StringInfo; the length word
1190
 *    is removed.  Also, s->cursor is initialized to zero for convenience
1191
 *    in scanning the message contents.
1192
 *
1193
 *    maxlen is the upper limit on the length of the
1194
 *    message we are willing to accept.  We abort the connection (by
1195
 *    returning EOF) if client tries to send more than that.
1196
 *
1197
 *    returns 0 if OK, EOF if trouble
1198
 * --------------------------------
1199
 */
1200
int
1201
pq_getmessage(StringInfo s, int maxlen)
1202
{
1203
  int32   len;
1204
1205
  Assert(PqCommReadingMsg);
1206
1207
  resetStringInfo(s);
1208
1209
  /* Read message length word */
1210
  if (pq_getbytes(&len, 4) == EOF)
1211
  {
1212
    ereport(COMMERROR,
1213
        (errcode(ERRCODE_PROTOCOL_VIOLATION),
1214
         errmsg("unexpected EOF within message length word")));
1215
    return EOF;
1216
  }
1217
1218
  len = pg_ntoh32(len);
1219
1220
  if (len < 4 || len > maxlen)
1221
  {
1222
    ereport(COMMERROR,
1223
        (errcode(ERRCODE_PROTOCOL_VIOLATION),
1224
         errmsg("invalid message length")));
1225
    return EOF;
1226
  }
1227
1228
  len -= 4;         /* discount length itself */
1229
1230
  if (len > 0)
1231
  {
1232
    /*
1233
     * Allocate space for message.  If we run out of room (ridiculously
1234
     * large message), we will elog(ERROR), but we want to discard the
1235
     * message body so as not to lose communication sync.
1236
     */
1237
    PG_TRY();
1238
    {
1239
      enlargeStringInfo(s, len);
1240
    }
1241
    PG_CATCH();
1242
    {
1243
      if (pq_discardbytes(len) == EOF)
1244
        ereport(COMMERROR,
1245
            (errcode(ERRCODE_PROTOCOL_VIOLATION),
1246
             errmsg("incomplete message from client")));
1247
1248
      /* we discarded the rest of the message so we're back in sync. */
1249
      PqCommReadingMsg = false;
1250
      PG_RE_THROW();
1251
    }
1252
    PG_END_TRY();
1253
1254
    /* And grab the message */
1255
    if (pq_getbytes(s->data, len) == EOF)
1256
    {
1257
      ereport(COMMERROR,
1258
          (errcode(ERRCODE_PROTOCOL_VIOLATION),
1259
           errmsg("incomplete message from client")));
1260
      return EOF;
1261
    }
1262
    s->len = len;
1263
    /* Place a trailing null per StringInfo convention */
1264
    s->data[len] = '\0';
1265
  }
1266
1267
  /* finished reading the message. */
1268
  PqCommReadingMsg = false;
1269
1270
  return 0;
1271
}
1272
1273
1274
static inline int
1275
internal_putbytes(const void *b, size_t len)
1276
0
{
1277
0
  const char *s = b;
1278
1279
0
  while (len > 0)
1280
0
  {
1281
    /* If buffer is full, then flush it out */
1282
0
    if (PqSendPointer >= PqSendBufferSize)
1283
0
    {
1284
0
      socket_set_nonblocking(false);
1285
0
      if (internal_flush())
1286
0
        return EOF;
1287
0
    }
1288
1289
    /*
1290
     * If the buffer is empty and data length is larger than the buffer
1291
     * size, send it without buffering.  Otherwise, copy as much data as
1292
     * possible into the buffer.
1293
     */
1294
0
    if (len >= PqSendBufferSize && PqSendStart == PqSendPointer)
1295
0
    {
1296
0
      size_t    start = 0;
1297
1298
0
      socket_set_nonblocking(false);
1299
0
      if (internal_flush_buffer(s, &start, &len))
1300
0
        return EOF;
1301
0
    }
1302
0
    else
1303
0
    {
1304
0
      size_t    amount = PqSendBufferSize - PqSendPointer;
1305
1306
0
      if (amount > len)
1307
0
        amount = len;
1308
0
      memcpy(PqSendBuffer + PqSendPointer, s, amount);
1309
0
      PqSendPointer += amount;
1310
0
      s += amount;
1311
0
      len -= amount;
1312
0
    }
1313
0
  }
1314
1315
0
  return 0;
1316
0
}
1317
1318
/* --------------------------------
1319
 *    socket_flush    - flush pending output
1320
 *
1321
 *    returns 0 if OK, EOF if trouble
1322
 * --------------------------------
1323
 */
1324
static int
1325
socket_flush(void)
1326
0
{
1327
0
  int     res;
1328
1329
  /* No-op if reentrant call */
1330
0
  if (PqCommBusy)
1331
0
    return 0;
1332
0
  PqCommBusy = true;
1333
0
  socket_set_nonblocking(false);
1334
0
  res = internal_flush();
1335
0
  PqCommBusy = false;
1336
0
  return res;
1337
0
}
1338
1339
/* --------------------------------
1340
 *    internal_flush - flush pending output
1341
 *
1342
 * Returns 0 if OK (meaning everything was sent, or operation would block
1343
 * and the socket is in non-blocking mode), or EOF if trouble.
1344
 * --------------------------------
1345
 */
1346
static inline int
1347
internal_flush(void)
1348
0
{
1349
0
  return internal_flush_buffer(PqSendBuffer, &PqSendStart, &PqSendPointer);
1350
0
}
1351
1352
/* --------------------------------
1353
 *    internal_flush_buffer - flush the given buffer content
1354
 *
1355
 * Returns 0 if OK (meaning everything was sent, or operation would block
1356
 * and the socket is in non-blocking mode), or EOF if trouble.
1357
 * --------------------------------
1358
 */
1359
static pg_noinline int
1360
internal_flush_buffer(const char *buf, size_t *start, size_t *end)
1361
{
1362
  static int  last_reported_send_errno = 0;
1363
1364
  const char *bufptr = buf + *start;
1365
  const char *bufend = buf + *end;
1366
1367
  while (bufptr < bufend)
1368
  {
1369
    ssize_t   r;
1370
1371
    r = secure_write(MyProcPort, bufptr, bufend - bufptr);
1372
1373
    if (r <= 0)
1374
    {
1375
      if (errno == EINTR)
1376
        continue;   /* Ok if we were interrupted */
1377
1378
      /*
1379
       * Ok if no data writable without blocking, and the socket is in
1380
       * non-blocking mode.
1381
       */
1382
      if (errno == EAGAIN ||
1383
        errno == EWOULDBLOCK)
1384
      {
1385
        return 0;
1386
      }
1387
1388
      /*
1389
       * Careful: an ereport() that tries to write to the client would
1390
       * cause recursion to here, leading to stack overflow and core
1391
       * dump!  This message must go *only* to the postmaster log.
1392
       *
1393
       * If a client disconnects while we're in the midst of output, we
1394
       * might write quite a bit of data before we get to a safe query
1395
       * abort point.  So, suppress duplicate log messages.
1396
       */
1397
      if (errno != last_reported_send_errno)
1398
      {
1399
        last_reported_send_errno = errno;
1400
        ereport(COMMERROR,
1401
            (errcode_for_socket_access(),
1402
             errmsg("could not send data to client: %m")));
1403
      }
1404
1405
      /*
1406
       * We drop the buffered data anyway so that processing can
1407
       * continue, even though we'll probably quit soon. We also set a
1408
       * flag that'll cause the next CHECK_FOR_INTERRUPTS to terminate
1409
       * the connection.
1410
       */
1411
      *start = *end = 0;
1412
      ClientConnectionLost = 1;
1413
      InterruptPending = 1;
1414
      return EOF;
1415
    }
1416
1417
    last_reported_send_errno = 0; /* reset after any successful send */
1418
    bufptr += r;
1419
    *start += r;
1420
  }
1421
1422
  *start = *end = 0;
1423
  return 0;
1424
}
1425
1426
/* --------------------------------
1427
 *    pq_flush_if_writable - flush pending output if writable without blocking
1428
 *
1429
 * Returns 0 if OK, or EOF if trouble.
1430
 * --------------------------------
1431
 */
1432
static int
1433
socket_flush_if_writable(void)
1434
0
{
1435
0
  int     res;
1436
1437
  /* Quick exit if nothing to do */
1438
0
  if (PqSendPointer == PqSendStart)
1439
0
    return 0;
1440
1441
  /* No-op if reentrant call */
1442
0
  if (PqCommBusy)
1443
0
    return 0;
1444
1445
  /* Temporarily put the socket into non-blocking mode */
1446
0
  socket_set_nonblocking(true);
1447
1448
0
  PqCommBusy = true;
1449
0
  res = internal_flush();
1450
0
  PqCommBusy = false;
1451
0
  return res;
1452
0
}
1453
1454
/* --------------------------------
1455
 *  socket_is_send_pending  - is there any pending data in the output buffer?
1456
 * --------------------------------
1457
 */
1458
static bool
1459
socket_is_send_pending(void)
1460
0
{
1461
0
  return (PqSendStart < PqSendPointer);
1462
0
}
1463
1464
/* --------------------------------
1465
 * Message-level I/O routines begin here.
1466
 * --------------------------------
1467
 */
1468
1469
1470
/* --------------------------------
1471
 *    socket_putmessage - send a normal message (suppressed in COPY OUT mode)
1472
 *
1473
 *    msgtype is a message type code to place before the message body.
1474
 *
1475
 *    len is the length of the message body data at *s.  A message length
1476
 *    word (equal to len+4 because it counts itself too) is inserted by this
1477
 *    routine.
1478
 *
1479
 *    We suppress messages generated while pqcomm.c is busy.  This
1480
 *    avoids any possibility of messages being inserted within other
1481
 *    messages.  The only known trouble case arises if SIGQUIT occurs
1482
 *    during a pqcomm.c routine --- quickdie() will try to send a warning
1483
 *    message, and the most reasonable approach seems to be to drop it.
1484
 *
1485
 *    returns 0 if OK, EOF if trouble
1486
 * --------------------------------
1487
 */
1488
static int
1489
socket_putmessage(char msgtype, const char *s, size_t len)
1490
0
{
1491
0
  uint32    n32;
1492
1493
0
  Assert(msgtype != 0);
1494
1495
0
  if (PqCommBusy)
1496
0
    return 0;
1497
0
  PqCommBusy = true;
1498
0
  if (internal_putbytes(&msgtype, 1))
1499
0
    goto fail;
1500
1501
0
  n32 = pg_hton32((uint32) (len + 4));
1502
0
  if (internal_putbytes(&n32, 4))
1503
0
    goto fail;
1504
1505
0
  if (internal_putbytes(s, len))
1506
0
    goto fail;
1507
0
  PqCommBusy = false;
1508
0
  return 0;
1509
1510
0
fail:
1511
0
  PqCommBusy = false;
1512
0
  return EOF;
1513
0
}
1514
1515
/* --------------------------------
1516
 *    pq_putmessage_noblock - like pq_putmessage, but never blocks
1517
 *
1518
 *    If the output buffer is too small to hold the message, the buffer
1519
 *    is enlarged.
1520
 */
1521
static void
1522
socket_putmessage_noblock(char msgtype, const char *s, size_t len)
1523
0
{
1524
0
  int     res PG_USED_FOR_ASSERTS_ONLY;
1525
0
  int     required;
1526
1527
  /*
1528
   * Ensure we have enough space in the output buffer for the message header
1529
   * as well as the message itself.
1530
   */
1531
0
  required = PqSendPointer + 1 + 4 + len;
1532
0
  if (required > PqSendBufferSize)
1533
0
  {
1534
0
    PqSendBuffer = repalloc(PqSendBuffer, required);
1535
0
    PqSendBufferSize = required;
1536
0
  }
1537
0
  res = socket_putmessage(msgtype, s, len);
1538
0
  Assert(res == 0);     /* should not fail when the message fits in
1539
                 * buffer */
1540
0
}
1541
1542
/* --------------------------------
1543
 *    pq_putmessage_v2 - send a message in protocol version 2
1544
 *
1545
 *    msgtype is a message type code to place before the message body.
1546
 *
1547
 *    We no longer support protocol version 2, but we have kept this
1548
 *    function so that if a client tries to connect with protocol version 2,
1549
 *    as a courtesy we can still send the "unsupported protocol version"
1550
 *    error to the client in the old format.
1551
 *
1552
 *    Like in pq_putmessage(), we suppress messages generated while
1553
 *    pqcomm.c is busy.
1554
 *
1555
 *    returns 0 if OK, EOF if trouble
1556
 * --------------------------------
1557
 */
1558
int
1559
pq_putmessage_v2(char msgtype, const char *s, size_t len)
1560
0
{
1561
0
  Assert(msgtype != 0);
1562
1563
0
  if (PqCommBusy)
1564
0
    return 0;
1565
0
  PqCommBusy = true;
1566
0
  if (internal_putbytes(&msgtype, 1))
1567
0
    goto fail;
1568
1569
0
  if (internal_putbytes(s, len))
1570
0
    goto fail;
1571
0
  PqCommBusy = false;
1572
0
  return 0;
1573
1574
0
fail:
1575
0
  PqCommBusy = false;
1576
0
  return EOF;
1577
0
}
1578
1579
/*
1580
 * Support for TCP Keepalive parameters
1581
 */
1582
1583
/*
1584
 * On Windows, we need to set both idle and interval at the same time.
1585
 * We also cannot reset them to the default (setting to zero will
1586
 * actually set them to zero, not default), therefore we fallback to
1587
 * the out-of-the-box default instead.
1588
 */
1589
#if defined(WIN32) && defined(SIO_KEEPALIVE_VALS)
1590
static int
1591
pq_setkeepaliveswin32(Port *port, int idle, int interval)
1592
{
1593
  struct tcp_keepalive ka;
1594
  DWORD   retsize;
1595
1596
  if (idle <= 0)
1597
    idle = 2 * 60 * 60;   /* default = 2 hours */
1598
  if (interval <= 0)
1599
    interval = 1;     /* default = 1 second */
1600
1601
  ka.onoff = 1;
1602
  ka.keepalivetime = idle * 1000;
1603
  ka.keepaliveinterval = interval * 1000;
1604
1605
  if (WSAIoctl(port->sock,
1606
         SIO_KEEPALIVE_VALS,
1607
         (LPVOID) &ka,
1608
         sizeof(ka),
1609
         NULL,
1610
         0,
1611
         &retsize,
1612
         NULL,
1613
         NULL)
1614
    != 0)
1615
  {
1616
    ereport(LOG,
1617
        (errmsg("%s(%s) failed: error code %d",
1618
            "WSAIoctl", "SIO_KEEPALIVE_VALS", WSAGetLastError())));
1619
    return STATUS_ERROR;
1620
  }
1621
  if (port->keepalives_idle != idle)
1622
    port->keepalives_idle = idle;
1623
  if (port->keepalives_interval != interval)
1624
    port->keepalives_interval = interval;
1625
  return STATUS_OK;
1626
}
1627
#endif
1628
1629
int
1630
pq_getkeepalivesidle(Port *port)
1631
{
1632
#if defined(PG_TCP_KEEPALIVE_IDLE) || defined(SIO_KEEPALIVE_VALS)
1633
  if (port == NULL || port->laddr.addr.ss_family == AF_UNIX)
1634
    return 0;
1635
1636
  if (port->keepalives_idle != 0)
1637
    return port->keepalives_idle;
1638
1639
  if (port->default_keepalives_idle == 0)
1640
  {
1641
#ifndef WIN32
1642
    socklen_t size = sizeof(port->default_keepalives_idle);
1643
1644
    if (getsockopt(port->sock, IPPROTO_TCP, PG_TCP_KEEPALIVE_IDLE,
1645
             (char *) &port->default_keepalives_idle,
1646
             &size) < 0)
1647
    {
1648
      ereport(LOG,
1649
          (errmsg("%s(%s) failed: %m", "getsockopt", PG_TCP_KEEPALIVE_IDLE_STR)));
1650
      port->default_keepalives_idle = -1; /* don't know */
1651
    }
1652
#else             /* WIN32 */
1653
    /* We can't get the defaults on Windows, so return "don't know" */
1654
    port->default_keepalives_idle = -1;
1655
#endif              /* WIN32 */
1656
  }
1657
1658
  return port->default_keepalives_idle;
1659
#else
1660
  return 0;
1661
#endif
1662
}
1663
1664
int
1665
pq_setkeepalivesidle(int idle, Port *port)
1666
2
{
1667
2
  if (port == NULL || port->laddr.addr.ss_family == AF_UNIX)
1668
2
    return STATUS_OK;
1669
1670
/* check SIO_KEEPALIVE_VALS here, not just WIN32, as some toolchains lack it */
1671
0
#if defined(PG_TCP_KEEPALIVE_IDLE) || defined(SIO_KEEPALIVE_VALS)
1672
0
  if (idle == port->keepalives_idle)
1673
0
    return STATUS_OK;
1674
1675
0
#ifndef WIN32
1676
0
  if (port->default_keepalives_idle <= 0)
1677
0
  {
1678
0
    if (pq_getkeepalivesidle(port) < 0)
1679
0
    {
1680
0
      if (idle == 0)
1681
0
        return STATUS_OK; /* default is set but unknown */
1682
0
      else
1683
0
        return STATUS_ERROR;
1684
0
    }
1685
0
  }
1686
1687
0
  if (idle == 0)
1688
0
    idle = port->default_keepalives_idle;
1689
1690
0
  if (setsockopt(port->sock, IPPROTO_TCP, PG_TCP_KEEPALIVE_IDLE,
1691
0
           (char *) &idle, sizeof(idle)) < 0)
1692
0
  {
1693
0
    ereport(LOG,
1694
0
        (errmsg("%s(%s) failed: %m", "setsockopt", PG_TCP_KEEPALIVE_IDLE_STR)));
1695
0
    return STATUS_ERROR;
1696
0
  }
1697
1698
0
  port->keepalives_idle = idle;
1699
#else             /* WIN32 */
1700
  return pq_setkeepaliveswin32(port, idle, port->keepalives_interval);
1701
#endif
1702
#else
1703
  if (idle != 0)
1704
  {
1705
    ereport(LOG,
1706
        (errmsg("setting the keepalive idle time is not supported")));
1707
    return STATUS_ERROR;
1708
  }
1709
#endif
1710
1711
0
  return STATUS_OK;
1712
0
}
1713
1714
int
1715
pq_getkeepalivesinterval(Port *port)
1716
{
1717
#if defined(TCP_KEEPINTVL) || defined(SIO_KEEPALIVE_VALS)
1718
  if (port == NULL || port->laddr.addr.ss_family == AF_UNIX)
1719
    return 0;
1720
1721
  if (port->keepalives_interval != 0)
1722
    return port->keepalives_interval;
1723
1724
  if (port->default_keepalives_interval == 0)
1725
  {
1726
#ifndef WIN32
1727
    socklen_t size = sizeof(port->default_keepalives_interval);
1728
1729
    if (getsockopt(port->sock, IPPROTO_TCP, TCP_KEEPINTVL,
1730
             (char *) &port->default_keepalives_interval,
1731
             &size) < 0)
1732
    {
1733
      ereport(LOG,
1734
          (errmsg("%s(%s) failed: %m", "getsockopt", "TCP_KEEPINTVL")));
1735
      port->default_keepalives_interval = -1; /* don't know */
1736
    }
1737
#else
1738
    /* We can't get the defaults on Windows, so return "don't know" */
1739
    port->default_keepalives_interval = -1;
1740
#endif              /* WIN32 */
1741
  }
1742
1743
  return port->default_keepalives_interval;
1744
#else
1745
  return 0;
1746
#endif
1747
}
1748
1749
int
1750
pq_setkeepalivesinterval(int interval, Port *port)
1751
2
{
1752
2
  if (port == NULL || port->laddr.addr.ss_family == AF_UNIX)
1753
2
    return STATUS_OK;
1754
1755
0
#if defined(TCP_KEEPINTVL) || defined(SIO_KEEPALIVE_VALS)
1756
0
  if (interval == port->keepalives_interval)
1757
0
    return STATUS_OK;
1758
1759
0
#ifndef WIN32
1760
0
  if (port->default_keepalives_interval <= 0)
1761
0
  {
1762
0
    if (pq_getkeepalivesinterval(port) < 0)
1763
0
    {
1764
0
      if (interval == 0)
1765
0
        return STATUS_OK; /* default is set but unknown */
1766
0
      else
1767
0
        return STATUS_ERROR;
1768
0
    }
1769
0
  }
1770
1771
0
  if (interval == 0)
1772
0
    interval = port->default_keepalives_interval;
1773
1774
0
  if (setsockopt(port->sock, IPPROTO_TCP, TCP_KEEPINTVL,
1775
0
           (char *) &interval, sizeof(interval)) < 0)
1776
0
  {
1777
0
    ereport(LOG,
1778
0
        (errmsg("%s(%s) failed: %m", "setsockopt", "TCP_KEEPINTVL")));
1779
0
    return STATUS_ERROR;
1780
0
  }
1781
1782
0
  port->keepalives_interval = interval;
1783
#else             /* WIN32 */
1784
  return pq_setkeepaliveswin32(port, port->keepalives_idle, interval);
1785
#endif
1786
#else
1787
  if (interval != 0)
1788
  {
1789
    ereport(LOG,
1790
        (errmsg("%s(%s) not supported", "setsockopt", "TCP_KEEPINTVL")));
1791
    return STATUS_ERROR;
1792
  }
1793
#endif
1794
1795
0
  return STATUS_OK;
1796
0
}
1797
1798
int
1799
pq_getkeepalivescount(Port *port)
1800
{
1801
#ifdef TCP_KEEPCNT
1802
  if (port == NULL || port->laddr.addr.ss_family == AF_UNIX)
1803
    return 0;
1804
1805
  if (port->keepalives_count != 0)
1806
    return port->keepalives_count;
1807
1808
  if (port->default_keepalives_count == 0)
1809
  {
1810
    socklen_t size = sizeof(port->default_keepalives_count);
1811
1812
    if (getsockopt(port->sock, IPPROTO_TCP, TCP_KEEPCNT,
1813
             (char *) &port->default_keepalives_count,
1814
             &size) < 0)
1815
    {
1816
      ereport(LOG,
1817
          (errmsg("%s(%s) failed: %m", "getsockopt", "TCP_KEEPCNT")));
1818
      port->default_keepalives_count = -1;  /* don't know */
1819
    }
1820
  }
1821
1822
  return port->default_keepalives_count;
1823
#else
1824
  return 0;
1825
#endif
1826
}
1827
1828
int
1829
pq_setkeepalivescount(int count, Port *port)
1830
2
{
1831
2
  if (port == NULL || port->laddr.addr.ss_family == AF_UNIX)
1832
2
    return STATUS_OK;
1833
1834
0
#ifdef TCP_KEEPCNT
1835
0
  if (count == port->keepalives_count)
1836
0
    return STATUS_OK;
1837
1838
0
  if (port->default_keepalives_count <= 0)
1839
0
  {
1840
0
    if (pq_getkeepalivescount(port) < 0)
1841
0
    {
1842
0
      if (count == 0)
1843
0
        return STATUS_OK; /* default is set but unknown */
1844
0
      else
1845
0
        return STATUS_ERROR;
1846
0
    }
1847
0
  }
1848
1849
0
  if (count == 0)
1850
0
    count = port->default_keepalives_count;
1851
1852
0
  if (setsockopt(port->sock, IPPROTO_TCP, TCP_KEEPCNT,
1853
0
           (char *) &count, sizeof(count)) < 0)
1854
0
  {
1855
0
    ereport(LOG,
1856
0
        (errmsg("%s(%s) failed: %m", "setsockopt", "TCP_KEEPCNT")));
1857
0
    return STATUS_ERROR;
1858
0
  }
1859
1860
0
  port->keepalives_count = count;
1861
#else
1862
  if (count != 0)
1863
  {
1864
    ereport(LOG,
1865
        (errmsg("%s(%s) not supported", "setsockopt", "TCP_KEEPCNT")));
1866
    return STATUS_ERROR;
1867
  }
1868
#endif
1869
1870
0
  return STATUS_OK;
1871
0
}
1872
1873
int
1874
pq_gettcpusertimeout(Port *port)
1875
{
1876
#ifdef TCP_USER_TIMEOUT
1877
  if (port == NULL || port->laddr.addr.ss_family == AF_UNIX)
1878
    return 0;
1879
1880
  if (port->tcp_user_timeout != 0)
1881
    return port->tcp_user_timeout;
1882
1883
  if (port->default_tcp_user_timeout == 0)
1884
  {
1885
    socklen_t size = sizeof(port->default_tcp_user_timeout);
1886
1887
    if (getsockopt(port->sock, IPPROTO_TCP, TCP_USER_TIMEOUT,
1888
             (char *) &port->default_tcp_user_timeout,
1889
             &size) < 0)
1890
    {
1891
      ereport(LOG,
1892
          (errmsg("%s(%s) failed: %m", "getsockopt", "TCP_USER_TIMEOUT")));
1893
      port->default_tcp_user_timeout = -1;  /* don't know */
1894
    }
1895
  }
1896
1897
  return port->default_tcp_user_timeout;
1898
#else
1899
  return 0;
1900
#endif
1901
}
1902
1903
int
1904
pq_settcpusertimeout(int timeout, Port *port)
1905
2
{
1906
2
  if (port == NULL || port->laddr.addr.ss_family == AF_UNIX)
1907
2
    return STATUS_OK;
1908
1909
0
#ifdef TCP_USER_TIMEOUT
1910
0
  if (timeout == port->tcp_user_timeout)
1911
0
    return STATUS_OK;
1912
1913
0
  if (port->default_tcp_user_timeout <= 0)
1914
0
  {
1915
0
    if (pq_gettcpusertimeout(port) < 0)
1916
0
    {
1917
0
      if (timeout == 0)
1918
0
        return STATUS_OK; /* default is set but unknown */
1919
0
      else
1920
0
        return STATUS_ERROR;
1921
0
    }
1922
0
  }
1923
1924
0
  if (timeout == 0)
1925
0
    timeout = port->default_tcp_user_timeout;
1926
1927
0
  if (setsockopt(port->sock, IPPROTO_TCP, TCP_USER_TIMEOUT,
1928
0
           (char *) &timeout, sizeof(timeout)) < 0)
1929
0
  {
1930
0
    ereport(LOG,
1931
0
        (errmsg("%s(%s) failed: %m", "setsockopt", "TCP_USER_TIMEOUT")));
1932
0
    return STATUS_ERROR;
1933
0
  }
1934
1935
0
  port->tcp_user_timeout = timeout;
1936
#else
1937
  if (timeout != 0)
1938
  {
1939
    ereport(LOG,
1940
        (errmsg("%s(%s) not supported", "setsockopt", "TCP_USER_TIMEOUT")));
1941
    return STATUS_ERROR;
1942
  }
1943
#endif
1944
1945
0
  return STATUS_OK;
1946
0
}
1947
1948
/*
1949
 * GUC assign_hook for tcp_keepalives_idle
1950
 */
1951
void
1952
assign_tcp_keepalives_idle(int newval, void *extra)
1953
2
{
1954
  /*
1955
   * The kernel API provides no way to test a value without setting it; and
1956
   * once we set it we might fail to unset it.  So there seems little point
1957
   * in fully implementing the check-then-assign GUC API for these
1958
   * variables.  Instead we just do the assignment on demand.
1959
   * pq_setkeepalivesidle reports any problems via ereport(LOG).
1960
   *
1961
   * This approach means that the GUC value might have little to do with the
1962
   * actual kernel value, so we use a show_hook that retrieves the kernel
1963
   * value rather than trusting GUC's copy.
1964
   */
1965
2
  (void) pq_setkeepalivesidle(newval, MyProcPort);
1966
2
}
1967
1968
/*
1969
 * GUC show_hook for tcp_keepalives_idle
1970
 */
1971
const char *
1972
show_tcp_keepalives_idle(void)
1973
0
{
1974
  /* See comments in assign_tcp_keepalives_idle */
1975
0
  static char nbuf[16];
1976
1977
0
  snprintf(nbuf, sizeof(nbuf), "%d", pq_getkeepalivesidle(MyProcPort));
1978
0
  return nbuf;
1979
0
}
1980
1981
/*
1982
 * GUC assign_hook for tcp_keepalives_interval
1983
 */
1984
void
1985
assign_tcp_keepalives_interval(int newval, void *extra)
1986
2
{
1987
  /* See comments in assign_tcp_keepalives_idle */
1988
2
  (void) pq_setkeepalivesinterval(newval, MyProcPort);
1989
2
}
1990
1991
/*
1992
 * GUC show_hook for tcp_keepalives_interval
1993
 */
1994
const char *
1995
show_tcp_keepalives_interval(void)
1996
0
{
1997
  /* See comments in assign_tcp_keepalives_idle */
1998
0
  static char nbuf[16];
1999
2000
0
  snprintf(nbuf, sizeof(nbuf), "%d", pq_getkeepalivesinterval(MyProcPort));
2001
0
  return nbuf;
2002
0
}
2003
2004
/*
2005
 * GUC assign_hook for tcp_keepalives_count
2006
 */
2007
void
2008
assign_tcp_keepalives_count(int newval, void *extra)
2009
2
{
2010
  /* See comments in assign_tcp_keepalives_idle */
2011
2
  (void) pq_setkeepalivescount(newval, MyProcPort);
2012
2
}
2013
2014
/*
2015
 * GUC show_hook for tcp_keepalives_count
2016
 */
2017
const char *
2018
show_tcp_keepalives_count(void)
2019
0
{
2020
  /* See comments in assign_tcp_keepalives_idle */
2021
0
  static char nbuf[16];
2022
2023
0
  snprintf(nbuf, sizeof(nbuf), "%d", pq_getkeepalivescount(MyProcPort));
2024
0
  return nbuf;
2025
0
}
2026
2027
/*
2028
 * GUC assign_hook for tcp_user_timeout
2029
 */
2030
void
2031
assign_tcp_user_timeout(int newval, void *extra)
2032
2
{
2033
  /* See comments in assign_tcp_keepalives_idle */
2034
2
  (void) pq_settcpusertimeout(newval, MyProcPort);
2035
2
}
2036
2037
/*
2038
 * GUC show_hook for tcp_user_timeout
2039
 */
2040
const char *
2041
show_tcp_user_timeout(void)
2042
0
{
2043
  /* See comments in assign_tcp_keepalives_idle */
2044
0
  static char nbuf[16];
2045
2046
0
  snprintf(nbuf, sizeof(nbuf), "%d", pq_gettcpusertimeout(MyProcPort));
2047
0
  return nbuf;
2048
0
}
2049
2050
/*
2051
 * Check if the client is still connected.
2052
 */
2053
bool
2054
pq_check_connection(void)
2055
0
{
2056
0
  WaitEvent events[FeBeWaitSetNEvents];
2057
0
  int     rc;
2058
2059
  /*
2060
   * It's OK to modify the socket event filter without restoring, because
2061
   * all FeBeWaitSet socket wait sites do the same.
2062
   */
2063
0
  ModifyWaitEvent(FeBeWaitSet, FeBeWaitSetSocketPos, WL_SOCKET_CLOSED, NULL);
2064
2065
0
retry:
2066
0
  rc = WaitEventSetWait(FeBeWaitSet, 0, events, lengthof(events), 0);
2067
0
  for (int i = 0; i < rc; ++i)
2068
0
  {
2069
0
    if (events[i].events & WL_SOCKET_CLOSED)
2070
0
      return false;
2071
0
    if (events[i].events & WL_LATCH_SET)
2072
0
    {
2073
      /*
2074
       * A latch event might be preventing other events from being
2075
       * reported.  Reset it and poll again.  No need to restore it
2076
       * because no code should expect latches to survive across
2077
       * CHECK_FOR_INTERRUPTS().
2078
       */
2079
0
      ResetLatch(MyLatch);
2080
0
      goto retry;
2081
0
    }
2082
0
  }
2083
2084
0
  return true;
2085
0
}