PostgreSQL Source Code git master
Loading...
Searching...
No Matches
parallel.c
Go to the documentation of this file.
1/*-------------------------------------------------------------------------
2 *
3 * parallel.c
4 *
5 * Parallel support for pg_dump and pg_restore
6 *
7 * Portions Copyright (c) 1996-2026, PostgreSQL Global Development Group
8 * Portions Copyright (c) 1994, Regents of the University of California
9 *
10 * IDENTIFICATION
11 * src/bin/pg_dump/parallel.c
12 *
13 *-------------------------------------------------------------------------
14 */
15
16/*
17 * Parallel operation works like this:
18 *
19 * The original, leader process calls ParallelBackupStart(), which forks off
20 * the desired number of worker processes, which each enter WaitForCommands().
21 *
22 * The leader process dispatches an individual work item to one of the worker
23 * processes in DispatchJobForTocEntry(). We send a command string such as
24 * "DUMP 1234" or "RESTORE 1234", where 1234 is the TocEntry ID.
25 * The worker process receives and decodes the command and passes it to the
26 * routine pointed to by AH->WorkerJobDumpPtr or AH->WorkerJobRestorePtr,
27 * which are routines of the current archive format. That routine performs
28 * the required action (dump or restore) and returns an integer status code.
29 * This is passed back to the leader where we pass it to the
30 * ParallelCompletionPtr callback function that was passed to
31 * DispatchJobForTocEntry(). The callback function does state updating
32 * for the leader control logic in pg_backup_archiver.c.
33 *
34 * In principle additional archive-format-specific information might be needed
35 * in commands or worker status responses, but so far that hasn't proved
36 * necessary, since workers have full copies of the ArchiveHandle/TocEntry
37 * data structures. Remember that we have forked off the workers only after
38 * we have read in the catalog. That's why our worker processes can also
39 * access the catalog information. (In the Windows case, the workers are
40 * threads in the same process. To avoid problems, they work with cloned
41 * copies of the Archive data structure; see RunWorker().)
42 *
43 * In the leader process, the workerStatus field for each worker has one of
44 * the following values:
45 * WRKR_NOT_STARTED: we've not yet forked this worker
46 * WRKR_IDLE: it's waiting for a command
47 * WRKR_WORKING: it's working on a command
48 * WRKR_TERMINATED: process ended
49 * The pstate->te[] entry for each worker is valid when it's in WRKR_WORKING
50 * state, and must be NULL in other states.
51 */
52
53#include "postgres_fe.h"
54
55#ifndef WIN32
56#include <sys/select.h>
57#include <sys/wait.h>
58#include <signal.h>
59#include <unistd.h>
60#include <fcntl.h>
61#endif
62
64#include "parallel.h"
65#include "pg_backup_utils.h"
66#ifdef WIN32
67#include "port/pg_bswap.h"
68#endif
69
70/* Mnemonic macros for indexing the fd array returned by pipe(2) */
71#define PIPE_READ 0
72#define PIPE_WRITE 1
73
74#define NO_SLOT (-1) /* Failure result for GetIdleWorker() */
75
76/* Worker process statuses */
84
85#define WORKER_IS_RUNNING(workerStatus) \
86 ((workerStatus) == WRKR_IDLE || (workerStatus) == WRKR_WORKING)
87
88/*
89 * Private per-parallel-worker state (typedef for this is in parallel.h).
90 *
91 * Much of this is valid only in the leader process (or, on Windows, should
92 * be touched only by the leader thread). But the AH field should be touched
93 * only by workers. The pipe descriptors are valid everywhere.
94 */
96{
97 T_WorkerStatus workerStatus; /* see enum above */
98
99 /* These fields are valid if workerStatus == WRKR_WORKING: */
100 ParallelCompletionPtr callback; /* function to call on completion */
101 void *callback_data; /* passthrough data for it */
102
103 ArchiveHandle *AH; /* Archive data worker is using */
104
105 int pipeRead; /* leader's end of the pipes */
107 int pipeRevRead; /* child's end of the pipes */
109
110 /* Child process/thread identity info: */
111#ifdef WIN32
113 unsigned int threadId;
114#else
116#endif
117};
118
119#ifdef WIN32
120
121/*
122 * Structure to hold info passed by _beginthreadex() to the function it calls
123 * via its single allowed argument.
124 */
125typedef struct
126{
127 ArchiveHandle *AH; /* leader database connection */
128 ParallelSlot *slot; /* this worker's parallel slot */
129} WorkerInfo;
130
131/* Windows implementation of pipe access */
132static int pgpipe(int handles[2]);
133#define piperead(a,b,c) recv(a,b,c,0)
134#define pipewrite(a,b,c) send(a,b,c,0)
135
136#else /* !WIN32 */
137
138/* Non-Windows implementation of pipe access */
139#define pgpipe(a) pipe(a)
140#define piperead(a,b,c) read(a,b,c)
141#define pipewrite(a,b,c) write(a,b,c)
142
143#endif /* WIN32 */
144
145/*
146 * State info for archive_close_connection() shutdown callback.
147 */
153
155
156/*
157 * State info for signal handling.
158 * We assume signal_info initializes to zeroes.
159 *
160 * On Unix, myAH is the leader DB connection in the leader process, and the
161 * worker's own connection in worker processes. On Windows, we have only one
162 * instance of signal_info, so myAH is the leader connection and the worker
163 * connections must be dug out of pstate->parallelSlot[].
164 */
166{
167 ArchiveHandle *myAH; /* database connection to issue cancel for */
168 ParallelState *pstate; /* parallel state, if any */
169 bool handler_set; /* signal handler set up in this process? */
170#ifndef WIN32
171 bool am_worker; /* am I a worker process? */
172#endif
174
176
177#ifdef WIN32
179#endif
180
181/*
182 * Write a simple string to stderr --- must be safe in a signal handler.
183 * We ignore the write() result since there's not much we could do about it.
184 * Certain compilers make that harder than it ought to be.
185 */
186#define write_stderr(str) \
187 do { \
188 const char *str_ = (str); \
189 ssize_t rc_; \
190 rc_ = write(fileno(stderr), str_, strlen(str_)); \
191 (void) rc_; \
192 } while (0)
193
194
195#ifdef WIN32
196/* file-scope variables */
197static DWORD tls_index;
198
199/* globally visible variables (needed by exit_nicely) */
200bool parallel_init_done = false;
202#endif /* WIN32 */
203
204/* Local function prototypes */
205static ParallelSlot *GetMyPSlot(ParallelState *pstate);
206static void archive_close_connection(int code, void *arg);
207static void ShutdownWorkersHard(ParallelState *pstate);
208static void WaitForTerminatingWorkers(ParallelState *pstate);
209static void set_cancel_handler(void);
210static void set_cancel_pstate(ParallelState *pstate);
212static void RunWorker(ArchiveHandle *AH, ParallelSlot *slot);
213static int GetIdleWorker(ParallelState *pstate);
214static bool HasEveryWorkerTerminated(ParallelState *pstate);
215static void lockTableForWorker(ArchiveHandle *AH, TocEntry *te);
216static void WaitForCommands(ArchiveHandle *AH, int pipefd[2]);
217static bool ListenToWorkers(ArchiveHandle *AH, ParallelState *pstate,
218 bool do_wait);
219static char *getMessageFromLeader(int pipefd[2]);
220static void sendMessageToLeader(int pipefd[2], const char *str);
221static int select_loop(int maxFd, fd_set *workerset);
222static char *getMessageFromWorker(ParallelState *pstate,
223 bool do_wait, int *worker);
224static void sendMessageToWorker(ParallelState *pstate,
225 int worker, const char *str);
226static char *readMessageFromPipe(int fd);
227
228#define messageStartsWith(msg, prefix) \
229 (strncmp(msg, prefix, strlen(prefix)) == 0)
230
231
232/*
233 * Initialize parallel dump support --- should be called early in process
234 * startup. (Currently, this is called whether or not we intend parallel
235 * activity.)
236 */
237void
239{
240#ifdef WIN32
242 {
244 int err;
245
246 /* Prepare for threaded operation */
249
250 /* Initialize socket access */
251 err = WSAStartup(MAKEWORD(2, 2), &wsaData);
252 if (err != 0)
253 pg_fatal("%s() failed: error code %d", "WSAStartup", err);
254
255 parallel_init_done = true;
256 }
257#endif
258}
259
260/*
261 * Find the ParallelSlot for the current worker process or thread.
262 *
263 * Returns NULL if no matching slot is found (this implies we're the leader).
264 */
265static ParallelSlot *
267{
268 int i;
269
270 for (i = 0; i < pstate->numWorkers; i++)
271 {
272#ifdef WIN32
273 if (pstate->parallelSlot[i].threadId == GetCurrentThreadId())
274#else
275 if (pstate->parallelSlot[i].pid == getpid())
276#endif
277 return &(pstate->parallelSlot[i]);
278 }
279
280 return NULL;
281}
282
283/*
284 * A thread-local version of getLocalPQExpBuffer().
285 *
286 * Non-reentrant but reduces memory leakage: we'll consume one buffer per
287 * thread, which is much better than one per fmtId/fmtQualifiedId call.
288 */
289#ifdef WIN32
290static PQExpBuffer
292{
293 /*
294 * The Tls code goes awry if we use a static var, so we provide for both
295 * static and auto, and omit any use of the static var when using Tls. We
296 * rely on TlsGetValue() to return 0 if the value is not yet set.
297 */
300
303 else
305
306 if (id_return) /* first time through? */
307 {
308 /* same buffer, just wipe contents */
310 }
311 else
312 {
313 /* new buffer */
317 else
319 }
320
321 return id_return;
322}
323#endif /* WIN32 */
324
325/*
326 * pg_dump and pg_restore call this to register the cleanup handler
327 * as soon as they've created the ArchiveHandle.
328 */
329void
335
336/*
337 * on_exit_nicely handler for shutting down database connections and
338 * worker processes cleanly.
339 */
340static void
342{
344
345 if (si->pstate)
346 {
347 /* In parallel mode, must figure out who we are */
348 ParallelSlot *slot = GetMyPSlot(si->pstate);
349
350 if (!slot)
351 {
352 /*
353 * We're the leader. Forcibly shut down workers, then close our
354 * own database connection, if any.
355 */
356 ShutdownWorkersHard(si->pstate);
357
358 if (si->AHX)
360 }
361 else
362 {
363 /*
364 * We're a worker. Shut down our own DB connection if any. On
365 * Windows, we also have to close our communication sockets, to
366 * emulate what will happen on Unix when the worker process exits.
367 * (Without this, if this is a premature exit, the leader would
368 * fail to detect it because there would be no EOF condition on
369 * the other end of the pipe.)
370 */
371 if (slot->AH)
372 DisconnectDatabase(&(slot->AH->public));
373
374#ifdef WIN32
377#endif
378 }
379 }
380 else
381 {
382 /* Non-parallel operation: just kill the leader DB connection */
383 if (si->AHX)
385 }
386}
387
388/*
389 * Forcibly shut down any remaining workers, waiting for them to finish.
390 *
391 * Note that we don't expect to come here during normal exit (the workers
392 * should be long gone, and the ParallelState too). We're only here in a
393 * pg_fatal() situation, so intervening to cancel active commands is
394 * appropriate.
395 */
396static void
398{
399 int i;
400
401 /*
402 * Close our write end of the sockets so that any workers waiting for
403 * commands know they can exit. (Note: some of the pipeWrite fields might
404 * still be zero, if we failed to initialize all the workers. Hence, just
405 * ignore errors here.)
406 */
407 for (i = 0; i < pstate->numWorkers; i++)
409
410 /*
411 * Force early termination of any commands currently in progress.
412 */
413#ifndef WIN32
414 /* On non-Windows, send SIGTERM to each worker process. */
415 for (i = 0; i < pstate->numWorkers; i++)
416 {
417 pid_t pid = pstate->parallelSlot[i].pid;
418
419 if (pid != 0)
420 kill(pid, SIGTERM);
421 }
422#else
423
424 /*
425 * On Windows, send query cancels directly to the workers' backends. Use
426 * a critical section to ensure worker threads don't change state.
427 */
429 for (i = 0; i < pstate->numWorkers; i++)
430 {
431 ArchiveHandle *AH = pstate->parallelSlot[i].AH;
432 char errbuf[1];
433
434 if (AH != NULL && AH->connCancel != NULL)
435 (void) PQcancel(AH->connCancel, errbuf, sizeof(errbuf));
436 }
438#endif
439
440 /* Now wait for them to terminate. */
442}
443
444/*
445 * Wait for all workers to terminate.
446 */
447static void
449{
450 while (!HasEveryWorkerTerminated(pstate))
451 {
452 ParallelSlot *slot = NULL;
453 int j;
454
455#ifndef WIN32
456 /* On non-Windows, use wait() to wait for next worker to end */
457 int status;
458 pid_t pid = wait(&status);
459
460 /* Find dead worker's slot, and clear the PID field */
461 for (j = 0; j < pstate->numWorkers; j++)
462 {
463 slot = &(pstate->parallelSlot[j]);
464 if (slot->pid == pid)
465 {
466 slot->pid = 0;
467 break;
468 }
469 }
470#else /* WIN32 */
471 /* On Windows, we must use WaitForMultipleObjects() */
473 int nrun = 0;
474 DWORD ret;
476
477 for (j = 0; j < pstate->numWorkers; j++)
478 {
480 {
481 lpHandles[nrun] = (HANDLE) pstate->parallelSlot[j].hThread;
482 nrun++;
483 }
484 }
486 Assert(ret != WAIT_FAILED);
489
490 /* Find dead worker's slot, and clear the hThread field */
491 for (j = 0; j < pstate->numWorkers; j++)
492 {
493 slot = &(pstate->parallelSlot[j]);
494 if (slot->hThread == hThread)
495 {
496 /* For cleanliness, close handles for dead threads */
497 CloseHandle((HANDLE) slot->hThread);
498 slot->hThread = (uintptr_t) INVALID_HANDLE_VALUE;
499 break;
500 }
501 }
502#endif /* WIN32 */
503
504 /* On all platforms, update workerStatus and te[] as well */
505 Assert(j < pstate->numWorkers);
507 pstate->te[j] = NULL;
508 }
509}
510
511
512/*
513 * Code for responding to cancel interrupts (SIGINT, control-C, etc)
514 *
515 * This doesn't quite belong in this module, but it needs access to the
516 * ParallelState data, so there's not really a better place either.
517 *
518 * When we get a cancel interrupt, we could just die, but in pg_restore that
519 * could leave a SQL command (e.g., CREATE INDEX on a large table) running
520 * for a long time. Instead, we try to send a cancel request and then die.
521 * pg_dump probably doesn't really need this, but we might as well use it
522 * there too. Note that sending the cancel directly from the signal handler
523 * is safe because PQcancel() is written to make it so.
524 *
525 * In parallel operation on Unix, each process is responsible for canceling
526 * its own connection (this must be so because nobody else has access to it).
527 * Furthermore, the leader process should attempt to forward its signal to
528 * each child. In simple manual use of pg_dump/pg_restore, forwarding isn't
529 * needed because typing control-C at the console would deliver SIGINT to
530 * every member of the terminal process group --- but in other scenarios it
531 * might be that only the leader gets signaled.
532 *
533 * On Windows, the cancel handler runs in a separate thread, because that's
534 * how SetConsoleCtrlHandler works. Because the workers are threads in this
535 * same process, we set a flag (is_cancel_in_progress()) so they stay quiet
536 * about the query cancellations instead of cluttering the screen, then send
537 * cancels on all active connections and return FALSE, which will allow the
538 * process to die. For safety's sake, we use a critical section to protect
539 * the PGcancel structures against being changed while the signal thread runs.
540 */
541
542#ifndef WIN32
543
544/*
545 * Signal handler (Unix only)
546 */
547static void
549{
550 int i;
551 char errbuf[1];
552
553 /*
554 * Some platforms allow delivery of new signals to interrupt an active
555 * signal handler. That could muck up our attempt to send PQcancel, so
556 * disable the signals that set_cancel_handler enabled.
557 */
561
562 /*
563 * If we're in the leader, forward signal to all workers. (It seems best
564 * to do this before PQcancel; killing the leader transaction will result
565 * in invalid-snapshot errors from active workers, which maybe we can
566 * quiet by killing workers first.) Ignore any errors.
567 */
568 if (signal_info.pstate != NULL)
569 {
570 for (i = 0; i < signal_info.pstate->numWorkers; i++)
571 {
573
574 if (pid != 0)
575 kill(pid, SIGTERM);
576 }
577 }
578
579 /*
580 * Send QueryCancel if we have a connection to send to. Ignore errors,
581 * there's not much we can do about them anyway.
582 */
584 (void) PQcancel(signal_info.myAH->connCancel, errbuf, sizeof(errbuf));
585
586 /*
587 * Report we're quitting, using nothing more complicated than write(2).
588 * When in parallel operation, only the leader process should do this.
589 */
591 {
592 if (progname)
593 {
595 write_stderr(": ");
596 }
597 write_stderr("terminated by user\n");
598 }
599
600 /*
601 * And die, using _exit() not exit() because the latter will invoke atexit
602 * handlers that can fail if we interrupted related code.
603 */
604 _exit(1);
605}
606
607/*
608 * Enable cancel interrupt handler, if not already done.
609 */
610static void
612{
613 /*
614 * When forking, signal_info.handler_set will propagate into the new
615 * process, but that's fine because the signal handler state does too.
616 */
618 {
620
624 }
625}
626
627#else /* WIN32 */
628
629/*
630 * Console interrupt handler --- runs in a newly-started thread.
631 *
632 * After stopping other threads and sending cancel requests on all open
633 * connections, we return FALSE which will allow the default ExitProcess()
634 * action to be taken.
635 */
636static BOOL WINAPI
638{
639 int i;
640 char errbuf[1];
641
642 if (dwCtrlType == CTRL_C_EVENT ||
644 {
645 /*
646 * Tell worker threads to stay quiet about the query cancellations
647 * we're about to send them; otherwise they'd report them as errors
648 * and clutter the user's screen. This must be set before we send any
649 * cancel, so that a worker is guaranteed to see it by the time its
650 * query fails as a result.
651 */
653
654 /* Critical section prevents changing data we look at here */
656
657 /*
658 * If in parallel mode, send QueryCancel to each worker's connected
659 * backend. Do this before canceling the main transaction, else we
660 * might get invalid-snapshot errors reported before we can stop the
661 * workers. Ignore errors, there's not much we can do about them
662 * anyway.
663 */
664 if (signal_info.pstate != NULL)
665 {
666 for (i = 0; i < signal_info.pstate->numWorkers; i++)
667 {
669
670 if (AH != NULL && AH->connCancel != NULL)
671 (void) PQcancel(AH->connCancel, errbuf, sizeof(errbuf));
672 }
673 }
674
675 /*
676 * Send QueryCancel to leader connection, if enabled. Ignore errors,
677 * there's not much we can do about them anyway.
678 */
681 errbuf, sizeof(errbuf));
682
684
685 /*
686 * Report we're quitting, using nothing more complicated than
687 * write(2). We should be able to use pg_log_*() here, but for now we
688 * stay aligned with the sigTermHandler behavior.
689 */
690 if (progname)
691 {
693 write_stderr(": ");
694 }
695 write_stderr("terminated by user\n");
696 }
697
698 /* Always return FALSE to allow signal handling to continue */
699 return FALSE;
700}
701
702/*
703 * Enable cancel interrupt handler, if not already done.
704 */
705static void
707{
709 {
711
713
715 }
716}
717
718#endif /* WIN32 */
719
720
721/*
722 * set_archive_cancel_info
723 *
724 * Fill AH->connCancel with cancellation info for the specified database
725 * connection; or clear it if conn is NULL.
726 */
727void
729{
731
732 /*
733 * Activate the interrupt handler if we didn't yet in this process. On
734 * Windows, this also initializes signal_info_lock; therefore it's
735 * important that this happen at least once before we fork off any
736 * threads.
737 */
739
740 /*
741 * On Unix, we assume that storing a pointer value is atomic with respect
742 * to any possible signal interrupt. On Windows, use a critical section.
743 */
744
745#ifdef WIN32
747#endif
748
749 /* Free the old one if we have one */
751 /* be sure interrupt handler doesn't use pointer while freeing */
752 AH->connCancel = NULL;
753
754 if (oldConnCancel != NULL)
756
757 /* Set the new one if specified */
758 if (conn)
760
761 /*
762 * On Unix, there's only ever one active ArchiveHandle per process, so we
763 * can just set signal_info.myAH unconditionally. On Windows, do that
764 * only in the main thread; worker threads have to make sure their
765 * ArchiveHandle appears in the pstate data, which is dealt with in
766 * RunWorker().
767 */
768#ifndef WIN32
769 signal_info.myAH = AH;
770#else
772 signal_info.myAH = AH;
773#endif
774
775#ifdef WIN32
777#endif
778}
779
780/*
781 * set_cancel_pstate
782 *
783 * Set signal_info.pstate to point to the specified ParallelState, if any.
784 * We need this mainly to have an interlock against Windows signal thread.
785 */
786static void
788{
789#ifdef WIN32
791#endif
792
793 signal_info.pstate = pstate;
794
795#ifdef WIN32
797#endif
798}
799
800/*
801 * set_cancel_slot_archive
802 *
803 * Set ParallelSlot's AH field to point to the specified archive, if any.
804 * We need this mainly to have an interlock against Windows signal thread.
805 */
806static void
808{
809#ifdef WIN32
811#endif
812
813 slot->AH = AH;
814
815#ifdef WIN32
817#endif
818}
819
820
821/*
822 * This function is called by both Unix and Windows variants to set up
823 * and run a worker process. Caller should exit the process (or thread)
824 * upon return.
825 */
826static void
828{
829 int pipefd[2];
830
831 /* fetch child ends of pipes */
834
835 /*
836 * Clone the archive so that we have our own state to work with, and in
837 * particular our own database connection.
838 *
839 * We clone on Unix as well as Windows, even though technically we don't
840 * need to because fork() gives us a copy in our own address space
841 * already. But CloneArchive resets the state information and also clones
842 * the database connection which both seem kinda helpful.
843 */
844 AH = CloneArchive(AH);
845
846 /* Remember cloned archive where signal handler can find it */
847 set_cancel_slot_archive(slot, AH);
848
849 /*
850 * Call the setup worker function that's defined in the ArchiveHandle.
851 */
852 (AH->SetupWorkerPtr) ((Archive *) AH);
853
854 /*
855 * Execute commands until done.
856 */
858
859 /*
860 * Disconnect from database and clean up.
861 */
864 DeCloneArchive(AH);
865}
866
867/*
868 * Thread base function for Windows
869 */
870#ifdef WIN32
871static unsigned __stdcall
873{
874 ArchiveHandle *AH = wi->AH;
875 ParallelSlot *slot = wi->slot;
876
877 /* Don't need WorkerInfo anymore */
878 free(wi);
879
880 /* Run the worker ... */
881 RunWorker(AH, slot);
882
883 /* Exit the thread */
884 _endthreadex(0);
885 return 0;
886}
887#endif /* WIN32 */
888
889/*
890 * This function starts a parallel dump or restore by spawning off the worker
891 * processes. For Windows, it creates a number of threads; on Unix the
892 * workers are created with fork().
893 */
896{
897 ParallelState *pstate;
898 int i;
899
900 Assert(AH->public.numWorkers > 0);
901
903
904 pstate->numWorkers = AH->public.numWorkers;
905 pstate->te = NULL;
906 pstate->parallelSlot = NULL;
907
908 if (AH->public.numWorkers == 1)
909 return pstate;
910
911 /* Create status arrays, being sure to initialize all fields to 0 */
912 pstate->te =
914 pstate->parallelSlot =
916
917#ifdef WIN32
918 /* Make fmtId() and fmtQualifiedId() use thread-local storage */
920#endif
921
922 /*
923 * Set the pstate in shutdown_info, to tell the exit handler that it must
924 * clean up workers as well as the main database connection. But we don't
925 * set this in signal_info yet, because we don't want child processes to
926 * inherit non-NULL signal_info.pstate.
927 */
928 shutdown_info.pstate = pstate;
929
930 /*
931 * Temporarily disable query cancellation on the leader connection. This
932 * ensures that child processes won't inherit valid AH->connCancel
933 * settings and thus won't try to issue cancels against the leader's
934 * connection. No harm is done if we fail while it's disabled, because
935 * the leader connection is idle at this point anyway.
936 */
938
939 /* Ensure stdio state is quiesced before forking */
940 fflush(NULL);
941
942 /* Create desired number of workers */
943 for (i = 0; i < pstate->numWorkers; i++)
944 {
945#ifdef WIN32
946 WorkerInfo *wi;
947 uintptr_t handle;
948#else
949 pid_t pid;
950#endif
951 ParallelSlot *slot = &(pstate->parallelSlot[i]);
952 int pipeMW[2],
953 pipeWM[2];
954
955 /* Create communication pipes for this worker */
956 if (pgpipe(pipeMW) < 0 || pgpipe(pipeWM) < 0)
957 pg_fatal("could not create communication channels: %m");
958
959 /* leader's ends of the pipes */
960 slot->pipeRead = pipeWM[PIPE_READ];
961 slot->pipeWrite = pipeMW[PIPE_WRITE];
962 /* child's ends of the pipes */
965
966#ifdef WIN32
967 /* Create transient structure to pass args to worker function */
969
970 wi->AH = AH;
971 wi->slot = slot;
972
973 handle = _beginthreadex(NULL, 0, (void *) &init_spawned_worker_win32,
974 wi, 0, &(slot->threadId));
975 if (handle == 0)
976 pg_fatal("could not create worker thread: %m");
977 slot->hThread = handle;
978 slot->workerStatus = WRKR_IDLE;
979#else /* !WIN32 */
980 pid = fork();
981 if (pid == 0)
982 {
983 /* we are the worker */
984 int j;
985
986 /* this is needed for GetMyPSlot() */
987 slot->pid = getpid();
988
989 /* instruct signal handler that we're in a worker now */
990 signal_info.am_worker = true;
991
992 /* close read end of Worker -> Leader */
994 /* close write end of Leader -> Worker */
996
997 /*
998 * Close all inherited fds for communication of the leader with
999 * previously-forked workers.
1000 */
1001 for (j = 0; j < i; j++)
1002 {
1005 }
1006
1007 /* Run the worker ... */
1008 RunWorker(AH, slot);
1009
1010 /* We can just exit(0) when done */
1011 exit(0);
1012 }
1013 else if (pid < 0)
1014 {
1015 /* fork failed */
1016 pg_fatal("could not create worker process: %m");
1017 }
1018
1019 /* In Leader after successful fork */
1020 slot->pid = pid;
1021 slot->workerStatus = WRKR_IDLE;
1022
1023 /* close read end of Leader -> Worker */
1025 /* close write end of Worker -> Leader */
1027#endif /* WIN32 */
1028 }
1029
1030 /*
1031 * Having forked off the workers, disable SIGPIPE so that leader isn't
1032 * killed if it tries to send a command to a dead worker. We don't want
1033 * the workers to inherit this setting, though.
1034 */
1035#ifndef WIN32
1037#endif
1038
1039 /*
1040 * Re-establish query cancellation on the leader connection.
1041 */
1043
1044 /*
1045 * Tell the cancel signal handler to forward signals to worker processes,
1046 * too. (As with query cancel, we did not need this earlier because the
1047 * workers have not yet been given anything to do; if we die before this
1048 * point, any already-started workers will see EOF and quit promptly.)
1049 */
1050 set_cancel_pstate(pstate);
1051
1052 return pstate;
1053}
1054
1055/*
1056 * Close down a parallel dump or restore.
1057 */
1058void
1060{
1061 int i;
1062
1063 /* No work if non-parallel */
1064 if (pstate->numWorkers == 1)
1065 return;
1066
1067 /* There should not be any unfinished jobs */
1068 Assert(IsEveryWorkerIdle(pstate));
1069
1070 /* Close the sockets so that the workers know they can exit */
1071 for (i = 0; i < pstate->numWorkers; i++)
1072 {
1075 }
1076
1077 /* Wait for them to exit */
1079
1080 /*
1081 * Unlink pstate from shutdown_info, so the exit handler will not try to
1082 * use it; and likewise unlink from signal_info.
1083 */
1086
1087 /* Release state (mere neatnik-ism, since we're about to terminate) */
1088 free(pstate->te);
1089 free(pstate->parallelSlot);
1090 free(pstate);
1091}
1092
1093/*
1094 * These next four functions handle construction and parsing of the command
1095 * strings and response strings for parallel workers.
1096 *
1097 * Currently, these can be the same regardless of which archive format we are
1098 * processing. In future, we might want to let format modules override these
1099 * functions to add format-specific data to a command or response.
1100 */
1101
1102/*
1103 * buildWorkerCommand: format a command string to send to a worker.
1104 *
1105 * The string is built in the caller-supplied buffer of size buflen.
1106 */
1107static void
1109 char *buf, int buflen)
1110{
1111 if (act == ACT_DUMP)
1112 snprintf(buf, buflen, "DUMP %d", te->dumpId);
1113 else if (act == ACT_RESTORE)
1114 snprintf(buf, buflen, "RESTORE %d", te->dumpId);
1115 else
1116 Assert(false);
1117}
1118
1119/*
1120 * parseWorkerCommand: interpret a command string in a worker.
1121 */
1122static void
1124 const char *msg)
1125{
1126 DumpId dumpId;
1127 int nBytes;
1128
1129 if (messageStartsWith(msg, "DUMP "))
1130 {
1131 *act = ACT_DUMP;
1132 sscanf(msg, "DUMP %d%n", &dumpId, &nBytes);
1133 Assert(nBytes == strlen(msg));
1134 *te = getTocEntryByDumpId(AH, dumpId);
1135 Assert(*te != NULL);
1136 }
1137 else if (messageStartsWith(msg, "RESTORE "))
1138 {
1139 *act = ACT_RESTORE;
1140 sscanf(msg, "RESTORE %d%n", &dumpId, &nBytes);
1141 Assert(nBytes == strlen(msg));
1142 *te = getTocEntryByDumpId(AH, dumpId);
1143 Assert(*te != NULL);
1144 }
1145 else
1146 pg_fatal("unrecognized command received from leader: \"%s\"",
1147 msg);
1148}
1149
1150/*
1151 * buildWorkerResponse: format a response string to send to the leader.
1152 *
1153 * The string is built in the caller-supplied buffer of size buflen.
1154 */
1155static void
1157 char *buf, int buflen)
1158{
1159 snprintf(buf, buflen, "OK %d %d %d",
1160 te->dumpId,
1161 status,
1162 status == WORKER_IGNORED_ERRORS ? AH->public.n_errors : 0);
1163}
1164
1165/*
1166 * parseWorkerResponse: parse the status message returned by a worker.
1167 *
1168 * Returns the integer status code, and may update fields of AH and/or te.
1169 */
1170static int
1172 const char *msg)
1173{
1174 DumpId dumpId;
1175 int nBytes,
1176 n_errors;
1177 int status = 0;
1178
1179 if (messageStartsWith(msg, "OK "))
1180 {
1181 sscanf(msg, "OK %d %d %d%n", &dumpId, &status, &n_errors, &nBytes);
1182
1183 Assert(dumpId == te->dumpId);
1184 Assert(nBytes == strlen(msg));
1185
1186 AH->public.n_errors += n_errors;
1187 }
1188 else
1189 pg_fatal("invalid message received from worker: \"%s\"",
1190 msg);
1191
1192 return status;
1193}
1194
1195/*
1196 * Dispatch a job to some free worker.
1197 *
1198 * te is the TocEntry to be processed, act is the action to be taken on it.
1199 * callback is the function to call on completion of the job.
1200 *
1201 * If no worker is currently available, this will block, and previously
1202 * registered callback functions may be called.
1203 */
1204void
1206 ParallelState *pstate,
1207 TocEntry *te,
1208 T_Action act,
1210 void *callback_data)
1211{
1212 int worker;
1213 char buf[256];
1214
1215 /* Get a worker, waiting if none are idle */
1216 while ((worker = GetIdleWorker(pstate)) == NO_SLOT)
1217 WaitForWorkers(AH, pstate, WFW_ONE_IDLE);
1218
1219 /* Construct and send command string */
1220 buildWorkerCommand(AH, te, act, buf, sizeof(buf));
1221
1222 sendMessageToWorker(pstate, worker, buf);
1223
1224 /* Remember worker is busy, and which TocEntry it's working on */
1225 pstate->parallelSlot[worker].workerStatus = WRKR_WORKING;
1226 pstate->parallelSlot[worker].callback = callback;
1227 pstate->parallelSlot[worker].callback_data = callback_data;
1228 pstate->te[worker] = te;
1229}
1230
1231/*
1232 * Find an idle worker and return its slot number.
1233 * Return NO_SLOT if none are idle.
1234 */
1235static int
1237{
1238 int i;
1239
1240 for (i = 0; i < pstate->numWorkers; i++)
1241 {
1242 if (pstate->parallelSlot[i].workerStatus == WRKR_IDLE)
1243 return i;
1244 }
1245 return NO_SLOT;
1246}
1247
1248/*
1249 * Return true iff no worker is running.
1250 */
1251static bool
1253{
1254 int i;
1255
1256 for (i = 0; i < pstate->numWorkers; i++)
1257 {
1259 return false;
1260 }
1261 return true;
1262}
1263
1264/*
1265 * Return true iff every worker is in the WRKR_IDLE state.
1266 */
1267bool
1269{
1270 int i;
1271
1272 for (i = 0; i < pstate->numWorkers; i++)
1273 {
1274 if (pstate->parallelSlot[i].workerStatus != WRKR_IDLE)
1275 return false;
1276 }
1277 return true;
1278}
1279
1280/*
1281 * Acquire lock on a table to be dumped by a worker process.
1282 *
1283 * The leader process is already holding an ACCESS SHARE lock. Ordinarily
1284 * it's no problem for a worker to get one too, but if anything else besides
1285 * pg_dump is running, there's a possible deadlock:
1286 *
1287 * 1) Leader dumps the schema and locks all tables in ACCESS SHARE mode.
1288 * 2) Another process requests an ACCESS EXCLUSIVE lock (which is not granted
1289 * because the leader holds a conflicting ACCESS SHARE lock).
1290 * 3) A worker process also requests an ACCESS SHARE lock to read the table.
1291 * The worker is enqueued behind the ACCESS EXCLUSIVE lock request.
1292 * 4) Now we have a deadlock, since the leader is effectively waiting for
1293 * the worker. The server cannot detect that, however.
1294 *
1295 * To prevent an infinite wait, prior to touching a table in a worker, request
1296 * a lock in ACCESS SHARE mode but with NOWAIT. If we don't get the lock,
1297 * then we know that somebody else has requested an ACCESS EXCLUSIVE lock and
1298 * so we have a deadlock. We must fail the backup in that case.
1299 */
1300static void
1302{
1303 const char *qualId;
1304 PQExpBuffer query;
1305 PGresult *res;
1306
1307 /* Nothing to do for BLOBS */
1308 if (strcmp(te->desc, "BLOBS") == 0)
1309 return;
1310
1311 query = createPQExpBuffer();
1312
1313 qualId = fmtQualifiedId(te->namespace, te->tag);
1314
1315 appendPQExpBuffer(query, "LOCK TABLE %s IN ACCESS SHARE MODE NOWAIT",
1316 qualId);
1317
1318 res = PQexec(AH->connection, query->data);
1319
1320 if (!res || PQresultStatus(res) != PGRES_COMMAND_OK)
1321 pg_fatal("could not obtain lock on relation \"%s\"\n"
1322 "This usually means that someone requested an ACCESS EXCLUSIVE lock "
1323 "on the table after the pg_dump parent process had gotten the "
1324 "initial ACCESS SHARE lock on the table.", qualId);
1325
1326 PQclear(res);
1327 destroyPQExpBuffer(query);
1328}
1329
1330/*
1331 * WaitForCommands: main routine for a worker process.
1332 *
1333 * Read and execute commands from the leader until we see EOF on the pipe.
1334 */
1335static void
1337{
1338 char *command;
1339 TocEntry *te;
1340 T_Action act;
1341 int status = 0;
1342 char buf[256];
1343
1344 for (;;)
1345 {
1346 if (!(command = getMessageFromLeader(pipefd)))
1347 {
1348 /* EOF, so done */
1349 return;
1350 }
1351
1352 /* Decode the command */
1353 parseWorkerCommand(AH, &te, &act, command);
1354
1355 if (act == ACT_DUMP)
1356 {
1357 /* Acquire lock on this table within the worker's session */
1358 lockTableForWorker(AH, te);
1359
1360 /* Perform the dump command */
1361 status = (AH->WorkerJobDumpPtr) (AH, te);
1362 }
1363 else if (act == ACT_RESTORE)
1364 {
1365 /* Perform the restore command */
1366 status = (AH->WorkerJobRestorePtr) (AH, te);
1367 }
1368 else
1369 Assert(false);
1370
1371 /* Return status to leader */
1372 buildWorkerResponse(AH, te, act, status, buf, sizeof(buf));
1373
1375
1376 /* command was pg_malloc'd and we are responsible for free()ing it. */
1377 free(command);
1378 }
1379}
1380
1381/*
1382 * Check for status messages from workers.
1383 *
1384 * If do_wait is true, wait to get a status message; otherwise, just return
1385 * immediately if there is none available.
1386 *
1387 * When we get a status message, we pass the status code to the callback
1388 * function that was specified to DispatchJobForTocEntry, then reset the
1389 * worker status to IDLE.
1390 *
1391 * Returns true if we collected a status message, else false.
1392 *
1393 * XXX is it worth checking for more than one status message per call?
1394 * It seems somewhat unlikely that multiple workers would finish at exactly
1395 * the same time.
1396 */
1397static bool
1399{
1400 int worker;
1401 char *msg;
1402
1403 /* Try to collect a status message */
1404 msg = getMessageFromWorker(pstate, do_wait, &worker);
1405
1406 if (!msg)
1407 {
1408 /* If do_wait is true, we must have detected EOF on some socket */
1409 if (do_wait)
1410 pg_fatal("a worker process died unexpectedly");
1411 return false;
1412 }
1413
1414 /* Process it and update our idea of the worker's status */
1415 if (messageStartsWith(msg, "OK "))
1416 {
1417 ParallelSlot *slot = &pstate->parallelSlot[worker];
1418 TocEntry *te = pstate->te[worker];
1419 int status;
1420
1421 status = parseWorkerResponse(AH, te, msg);
1422 slot->callback(AH, te, status, slot->callback_data);
1423 slot->workerStatus = WRKR_IDLE;
1424 pstate->te[worker] = NULL;
1425 }
1426 else
1427 pg_fatal("invalid message received from worker: \"%s\"",
1428 msg);
1429
1430 /* Free the string returned from getMessageFromWorker */
1431 free(msg);
1432
1433 return true;
1434}
1435
1436/*
1437 * Check for status results from workers, waiting if necessary.
1438 *
1439 * Available wait modes are:
1440 * WFW_NO_WAIT: reap any available status, but don't block
1441 * WFW_GOT_STATUS: wait for at least one more worker to finish
1442 * WFW_ONE_IDLE: wait for at least one worker to be idle
1443 * WFW_ALL_IDLE: wait for all workers to be idle
1444 *
1445 * Any received results are passed to the callback specified to
1446 * DispatchJobForTocEntry.
1447 *
1448 * This function is executed in the leader process.
1449 */
1450void
1452{
1453 bool do_wait = false;
1454
1455 /*
1456 * In GOT_STATUS mode, always block waiting for a message, since we can't
1457 * return till we get something. In other modes, we don't block the first
1458 * time through the loop.
1459 */
1460 if (mode == WFW_GOT_STATUS)
1461 {
1462 /* Assert that caller knows what it's doing */
1463 Assert(!IsEveryWorkerIdle(pstate));
1464 do_wait = true;
1465 }
1466
1467 for (;;)
1468 {
1469 /*
1470 * Check for status messages, even if we don't need to block. We do
1471 * not try very hard to reap all available messages, though, since
1472 * there's unlikely to be more than one.
1473 */
1474 if (ListenToWorkers(AH, pstate, do_wait))
1475 {
1476 /*
1477 * If we got a message, we are done by definition for GOT_STATUS
1478 * mode, and we can also be certain that there's at least one idle
1479 * worker. So we're done in all but ALL_IDLE mode.
1480 */
1481 if (mode != WFW_ALL_IDLE)
1482 return;
1483 }
1484
1485 /* Check whether we must wait for new status messages */
1486 switch (mode)
1487 {
1488 case WFW_NO_WAIT:
1489 return; /* never wait */
1490 case WFW_GOT_STATUS:
1491 Assert(false); /* can't get here, because we waited */
1492 break;
1493 case WFW_ONE_IDLE:
1494 if (GetIdleWorker(pstate) != NO_SLOT)
1495 return;
1496 break;
1497 case WFW_ALL_IDLE:
1498 if (IsEveryWorkerIdle(pstate))
1499 return;
1500 break;
1501 }
1502
1503 /* Loop back, and this time wait for something to happen */
1504 do_wait = true;
1505 }
1506}
1507
1508/*
1509 * Read one command message from the leader, blocking if necessary
1510 * until one is available, and return it as a malloc'd string.
1511 * On EOF, return NULL.
1512 *
1513 * This function is executed in worker processes.
1514 */
1515static char *
1520
1521/*
1522 * Send a status message to the leader.
1523 *
1524 * This function is executed in worker processes.
1525 */
1526static void
1527sendMessageToLeader(int pipefd[2], const char *str)
1528{
1529 size_t len = strlen(str) + 1;
1530
1531 if (pipewrite(pipefd[PIPE_WRITE], str, len) != len)
1532 pg_fatal("could not write to the communication channel: %m");
1533}
1534
1535/*
1536 * Wait until some descriptor in "workerset" becomes readable.
1537 * Returns -1 on error, else the number of readable descriptors.
1538 */
1539static int
1540select_loop(int maxFd, fd_set *workerset)
1541{
1542 int i;
1543 fd_set saveSet = *workerset;
1544
1545 for (;;)
1546 {
1547 *workerset = saveSet;
1548 i = select(maxFd + 1, workerset, NULL, NULL, NULL);
1549
1550#ifndef WIN32
1551 if (i < 0 && errno == EINTR)
1552 continue;
1553#else
1554 if (i == SOCKET_ERROR && WSAGetLastError() == WSAEINTR)
1555 continue;
1556#endif
1557 break;
1558 }
1559
1560 return i;
1561}
1562
1563
1564/*
1565 * Check for messages from worker processes.
1566 *
1567 * If a message is available, return it as a malloc'd string, and put the
1568 * index of the sending worker in *worker.
1569 *
1570 * If nothing is available, wait if "do_wait" is true, else return NULL.
1571 *
1572 * If we detect EOF on any socket, we'll return NULL. It's not great that
1573 * that's hard to distinguish from the no-data-available case, but for now
1574 * our one caller is okay with that.
1575 *
1576 * This function is executed in the leader process.
1577 */
1578static char *
1579getMessageFromWorker(ParallelState *pstate, bool do_wait, int *worker)
1580{
1581 int i;
1582 fd_set workerset;
1583 int maxFd = -1;
1584 struct timeval nowait = {0, 0};
1585
1586 /* construct bitmap of socket descriptors for select() */
1587 FD_ZERO(&workerset);
1588 for (i = 0; i < pstate->numWorkers; i++)
1589 {
1591 continue;
1592 FD_SET(pstate->parallelSlot[i].pipeRead, &workerset);
1593 if (pstate->parallelSlot[i].pipeRead > maxFd)
1594 maxFd = pstate->parallelSlot[i].pipeRead;
1595 }
1596
1597 if (do_wait)
1598 {
1599 i = select_loop(maxFd, &workerset);
1600 Assert(i != 0);
1601 }
1602 else
1603 {
1604 if ((i = select(maxFd + 1, &workerset, NULL, NULL, &nowait)) == 0)
1605 return NULL;
1606 }
1607
1608 if (i < 0)
1609 pg_fatal("%s() failed: %m", "select");
1610
1611 for (i = 0; i < pstate->numWorkers; i++)
1612 {
1613 char *msg;
1614
1616 continue;
1617 if (!FD_ISSET(pstate->parallelSlot[i].pipeRead, &workerset))
1618 continue;
1619
1620 /*
1621 * Read the message if any. If the socket is ready because of EOF,
1622 * we'll return NULL instead (and the socket will stay ready, so the
1623 * condition will persist).
1624 *
1625 * Note: because this is a blocking read, we'll wait if only part of
1626 * the message is available. Waiting a long time would be bad, but
1627 * since worker status messages are short and are always sent in one
1628 * operation, it shouldn't be a problem in practice.
1629 */
1631 *worker = i;
1632 return msg;
1633 }
1634 Assert(false);
1635 return NULL;
1636}
1637
1638/*
1639 * Send a command message to the specified worker process.
1640 *
1641 * This function is executed in the leader process.
1642 */
1643static void
1644sendMessageToWorker(ParallelState *pstate, int worker, const char *str)
1645{
1646 size_t len = strlen(str) + 1;
1647
1648 if (pipewrite(pstate->parallelSlot[worker].pipeWrite, str, len) != len)
1649 {
1650 pg_fatal("could not write to the communication channel: %m");
1651 }
1652}
1653
1654/*
1655 * Read one message from the specified pipe (fd), blocking if necessary
1656 * until one is available, and return it as a malloc'd string.
1657 * On EOF, return NULL.
1658 *
1659 * A "message" on the channel is just a null-terminated string.
1660 */
1661static char *
1663{
1664 char *msg;
1665 int msgsize,
1666 bufsize;
1667 int ret;
1668
1669 /*
1670 * In theory, if we let piperead() read multiple bytes, it might give us
1671 * back fragments of multiple messages. (That can't actually occur, since
1672 * neither leader nor workers send more than one message without waiting
1673 * for a reply, but we don't wish to assume that here.) For simplicity,
1674 * read a byte at a time until we get the terminating '\0'. This method
1675 * is a bit inefficient, but since this is only used for relatively short
1676 * command and status strings, it shouldn't matter.
1677 */
1678 bufsize = 64; /* could be any number */
1679 msg = (char *) pg_malloc(bufsize);
1680 msgsize = 0;
1681 for (;;)
1682 {
1684 ret = piperead(fd, msg + msgsize, 1);
1685 if (ret <= 0)
1686 break; /* error or connection closure */
1687
1688 Assert(ret == 1);
1689
1690 if (msg[msgsize] == '\0')
1691 return msg; /* collected whole message */
1692
1693 msgsize++;
1694 if (msgsize == bufsize) /* enlarge buffer if needed */
1695 {
1696 bufsize += 16; /* could be any number */
1697 msg = (char *) pg_realloc(msg, bufsize);
1698 }
1699 }
1700
1701 /* Other end has closed the connection */
1702 pg_free(msg);
1703 return NULL;
1704}
1705
1706#ifdef WIN32
1707
1708/*
1709 * This is a replacement version of pipe(2) for Windows which allows the pipe
1710 * handles to be used in select().
1711 *
1712 * Reads and writes on the pipe must go through piperead()/pipewrite().
1713 *
1714 * For consistency with Unix we declare the returned handles as "int".
1715 * This is okay even on WIN64 because system handles are not more than
1716 * 32 bits wide, but we do have to do some casting.
1717 */
1718static int
1719pgpipe(int handles[2])
1720{
1721 pgsocket s,
1722 tmp_sock;
1723 struct sockaddr_in serv_addr;
1724 int len = sizeof(serv_addr);
1725
1726 /* We have to use the Unix socket invalid file descriptor value here. */
1727 handles[0] = handles[1] = -1;
1728
1729 /*
1730 * setup listen socket
1731 */
1732 if ((s = socket(AF_INET, SOCK_STREAM, 0)) == PGINVALID_SOCKET)
1733 {
1734 pg_log_error("pgpipe: could not create socket: error code %d",
1735 WSAGetLastError());
1736 return -1;
1737 }
1738
1739 memset(&serv_addr, 0, sizeof(serv_addr));
1740 serv_addr.sin_family = AF_INET;
1741 serv_addr.sin_port = pg_hton16(0);
1742 serv_addr.sin_addr.s_addr = pg_hton32(INADDR_LOOPBACK);
1743 if (bind(s, (SOCKADDR *) &serv_addr, len) == SOCKET_ERROR)
1744 {
1745 pg_log_error("pgpipe: could not bind: error code %d",
1746 WSAGetLastError());
1747 closesocket(s);
1748 return -1;
1749 }
1750 if (listen(s, 1) == SOCKET_ERROR)
1751 {
1752 pg_log_error("pgpipe: could not listen: error code %d",
1753 WSAGetLastError());
1754 closesocket(s);
1755 return -1;
1756 }
1757 if (getsockname(s, (SOCKADDR *) &serv_addr, &len) == SOCKET_ERROR)
1758 {
1759 pg_log_error("pgpipe: %s() failed: error code %d", "getsockname",
1760 WSAGetLastError());
1761 closesocket(s);
1762 return -1;
1763 }
1764
1765 /*
1766 * setup pipe handles
1767 */
1769 {
1770 pg_log_error("pgpipe: could not create second socket: error code %d",
1771 WSAGetLastError());
1772 closesocket(s);
1773 return -1;
1774 }
1775 handles[1] = (int) tmp_sock;
1776
1778 {
1779 pg_log_error("pgpipe: could not connect socket: error code %d",
1780 WSAGetLastError());
1781 closesocket(handles[1]);
1782 handles[1] = -1;
1783 closesocket(s);
1784 return -1;
1785 }
1786 if ((tmp_sock = accept(s, (SOCKADDR *) &serv_addr, &len)) == PGINVALID_SOCKET)
1787 {
1788 pg_log_error("pgpipe: could not accept connection: error code %d",
1789 WSAGetLastError());
1790 closesocket(handles[1]);
1791 handles[1] = -1;
1792 closesocket(s);
1793 return -1;
1794 }
1795 handles[0] = (int) tmp_sock;
1796
1797 closesocket(s);
1798 return 0;
1799}
1800
1801#endif /* WIN32 */
struct WorkerInfoData * WorkerInfo
Definition autovacuum.c:251
void ParallelBackupEnd(ArchiveHandle *AH, ParallelState *pstate)
Definition parallel.c:1059
static void sendMessageToLeader(int pipefd[2], const char *str)
Definition parallel.c:1527
static ParallelSlot * GetMyPSlot(ParallelState *pstate)
Definition parallel.c:266
static void WaitForCommands(ArchiveHandle *AH, int pipefd[2])
Definition parallel.c:1336
void WaitForWorkers(ArchiveHandle *AH, ParallelState *pstate, WFW_WaitOption mode)
Definition parallel.c:1451
T_WorkerStatus
Definition parallel.c:78
@ WRKR_WORKING
Definition parallel.c:81
@ WRKR_IDLE
Definition parallel.c:80
@ WRKR_TERMINATED
Definition parallel.c:82
@ WRKR_NOT_STARTED
Definition parallel.c:79
static bool HasEveryWorkerTerminated(ParallelState *pstate)
Definition parallel.c:1252
#define pgpipe(a)
Definition parallel.c:139
static bool ListenToWorkers(ArchiveHandle *AH, ParallelState *pstate, bool do_wait)
Definition parallel.c:1398
static void sigTermHandler(SIGNAL_ARGS)
Definition parallel.c:548
#define PIPE_READ
Definition parallel.c:71
ParallelState * ParallelBackupStart(ArchiveHandle *AH)
Definition parallel.c:895
static char * readMessageFromPipe(int fd)
Definition parallel.c:1662
static int select_loop(int maxFd, fd_set *workerset)
Definition parallel.c:1540
static int parseWorkerResponse(ArchiveHandle *AH, TocEntry *te, const char *msg)
Definition parallel.c:1171
static int GetIdleWorker(ParallelState *pstate)
Definition parallel.c:1236
static void set_cancel_pstate(ParallelState *pstate)
Definition parallel.c:787
static void RunWorker(ArchiveHandle *AH, ParallelSlot *slot)
Definition parallel.c:827
static void set_cancel_slot_archive(ParallelSlot *slot, ArchiveHandle *AH)
Definition parallel.c:807
static void buildWorkerCommand(ArchiveHandle *AH, TocEntry *te, T_Action act, char *buf, int buflen)
Definition parallel.c:1108
static char * getMessageFromWorker(ParallelState *pstate, bool do_wait, int *worker)
Definition parallel.c:1579
static void archive_close_connection(int code, void *arg)
Definition parallel.c:341
#define NO_SLOT
Definition parallel.c:74
static void sendMessageToWorker(ParallelState *pstate, int worker, const char *str)
Definition parallel.c:1644
#define PIPE_WRITE
Definition parallel.c:72
static ShutdownInformation shutdown_info
Definition parallel.c:154
void on_exit_close_archive(Archive *AHX)
Definition parallel.c:330
void DispatchJobForTocEntry(ArchiveHandle *AH, ParallelState *pstate, TocEntry *te, T_Action act, ParallelCompletionPtr callback, void *callback_data)
Definition parallel.c:1205
#define WORKER_IS_RUNNING(workerStatus)
Definition parallel.c:85
static char * getMessageFromLeader(int pipefd[2])
Definition parallel.c:1516
static void lockTableForWorker(ArchiveHandle *AH, TocEntry *te)
Definition parallel.c:1301
#define piperead(a, b, c)
Definition parallel.c:140
#define pipewrite(a, b, c)
Definition parallel.c:141
void init_parallel_dump_utils(void)
Definition parallel.c:238
static void set_cancel_handler(void)
Definition parallel.c:611
static void buildWorkerResponse(ArchiveHandle *AH, TocEntry *te, T_Action act, int status, char *buf, int buflen)
Definition parallel.c:1156
static volatile DumpSignalInformation signal_info
Definition parallel.c:175
bool IsEveryWorkerIdle(ParallelState *pstate)
Definition parallel.c:1268
#define write_stderr(str)
Definition parallel.c:186
static void parseWorkerCommand(ArchiveHandle *AH, TocEntry **te, T_Action *act, const char *msg)
Definition parallel.c:1123
#define messageStartsWith(msg, prefix)
Definition parallel.c:228
static void ShutdownWorkersHard(ParallelState *pstate)
Definition parallel.c:397
static void WaitForTerminatingWorkers(ParallelState *pstate)
Definition parallel.c:448
void set_archive_cancel_info(ArchiveHandle *AH, PGconn *conn)
Definition parallel.c:728
void(* ParallelCompletionPtr)(ArchiveHandle *AH, TocEntry *te, int status, void *callback_data)
Definition parallel.h:24
WFW_WaitOption
Definition parallel.h:31
@ WFW_ALL_IDLE
Definition parallel.h:35
@ WFW_GOT_STATUS
Definition parallel.h:33
@ WFW_NO_WAIT
Definition parallel.h:32
@ WFW_ONE_IDLE
Definition parallel.h:34
#define SIGNAL_ARGS
Definition c.h:1519
#define Assert(condition)
Definition c.h:1002
Datum arg
Definition elog.c:1323
void err(int eval, const char *fmt,...)
Definition err.c:43
PGcancel * PQgetCancel(PGconn *conn)
Definition fe-cancel.c:368
int PQcancel(PGcancel *cancel, char *errbuf, int errbufsize)
Definition fe-cancel.c:548
void PQfreeCancel(PGcancel *cancel)
Definition fe-cancel.c:502
PGresult * PQexec(PGconn *conn, const char *query)
Definition fe-exec.c:2279
void * pg_malloc(size_t size)
Definition fe_memutils.c:53
void pg_free(void *ptr)
void * pg_realloc(void *ptr, size_t size)
Definition fe_memutils.c:71
#define pg_malloc_array(type, count)
Definition fe_memutils.h:66
#define pg_malloc_object(type)
Definition fe_memutils.h:60
#define pg_malloc0_array(type, count)
Definition fe_memutils.h:67
const char * str
#define bufsize
int j
Definition isn.c:78
int i
Definition isn.c:77
#define PQclear
#define PQresultStatus
@ PGRES_COMMAND_OK
Definition libpq-fe.h:131
#define pg_log_error(...)
Definition logging.h:108
const char * progname
Definition main.c:44
int DumpId
Definition pg_backup.h:285
void DisconnectDatabase(Archive *AHX)
void DeCloneArchive(ArchiveHandle *AH)
ArchiveHandle * CloneArchive(ArchiveHandle *AH)
TocEntry * getTocEntryByDumpId(ArchiveHandle *AH, DumpId id)
#define WORKER_IGNORED_ERRORS
@ ACT_RESTORE
void on_exit_nicely(on_exit_nicely_callback function, void *arg)
#define pg_fatal(...)
#define pg_hton32(x)
Definition pg_bswap.h:121
#define pg_hton16(x)
Definition pg_bswap.h:120
static PgChecksumMode mode
const void size_t len
static bool do_wait
Definition pg_ctl.c:76
static char buf[DEFAULT_XLOG_SEG_SIZE]
#define pqsignal
Definition port.h:548
#define PG_SIG_IGN
Definition port.h:552
int pgsocket
Definition port.h:29
#define snprintf
Definition port.h:261
#define PGINVALID_SOCKET
Definition port.h:31
#define closesocket
Definition port.h:398
PQExpBuffer createPQExpBuffer(void)
Definition pqexpbuffer.c:72
void resetPQExpBuffer(PQExpBuffer str)
void appendPQExpBuffer(PQExpBuffer str, const char *fmt,...)
void destroyPQExpBuffer(PQExpBuffer str)
PQExpBufferData * PQExpBuffer
Definition pqexpbuffer.h:51
static int fd(const char *x, int i)
static int fb(int x)
#define free(a)
PGconn * conn
Definition streamutil.c:52
const char * fmtQualifiedId(const char *schema, const char *id)
PQExpBuffer(* getLocalPQExpBuffer)(void)
int n_errors
Definition pg_backup.h:253
int numWorkers
Definition pg_backup.h:240
ArchiveHandle * myAH
Definition parallel.c:167
ParallelState * pstate
Definition parallel.c:168
ParallelCompletionPtr callback
Definition parallel.c:100
ArchiveHandle * AH
Definition parallel.c:103
void * callback_data
Definition parallel.c:101
T_WorkerStatus workerStatus
Definition parallel.c:97
int pipeRevRead
Definition parallel.c:107
int pipeRevWrite
Definition parallel.c:108
TocEntry ** te
Definition parallel.h:59
ParallelSlot * parallelSlot
Definition parallel.h:60
ParallelState * pstate
Definition parallel.c:150
WorkerJobDumpPtrType WorkerJobDumpPtr
PGcancel *volatile connCancel
WorkerJobRestorePtrType WorkerJobRestorePtr
SetupWorkerPtrType SetupWorkerPtr
static void callback(struct sockaddr *addr, struct sockaddr *mask, void *unused)
#define bind(s, addr, addrlen)
Definition win32_port.h:513
#define EINTR
Definition win32_port.h:378
#define SIGPIPE
Definition win32_port.h:163
#define SIGQUIT
Definition win32_port.h:159
#define kill(pid, sig)
Definition win32_port.h:507
#define socket(af, type, protocol)
Definition win32_port.h:512
#define accept(s, addr, addrlen)
Definition win32_port.h:515
#define connect(s, name, namelen)
Definition win32_port.h:516
#define listen(s, backlog)
Definition win32_port.h:514
#define select(n, r, w, e, timeout)
Definition win32_port.h:517