Coverage Report

Created: 2026-09-28 06:55

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/src/postgres/src/backend/commands/copyto.c
Line
Count
Source
1
/*-------------------------------------------------------------------------
2
 *
3
 * copyto.c
4
 *    COPY <table> TO file/program/client
5
 *
6
 * Portions Copyright (c) 1996-2026, PostgreSQL Global Development Group
7
 * Portions Copyright (c) 1994, Regents of the University of California
8
 *
9
 *
10
 * IDENTIFICATION
11
 *    src/backend/commands/copyto.c
12
 *
13
 *-------------------------------------------------------------------------
14
 */
15
#include "postgres.h"
16
17
#include <ctype.h>
18
#include <unistd.h>
19
#include <sys/stat.h>
20
21
#include "access/table.h"
22
#include "access/tableam.h"
23
#include "access/tupconvert.h"
24
#include "catalog/pg_inherits.h"
25
#include "commands/copyapi.h"
26
#include "commands/progress.h"
27
#include "executor/execdesc.h"
28
#include "executor/executor.h"
29
#include "executor/tuptable.h"
30
#include "funcapi.h"
31
#include "libpq/libpq.h"
32
#include "libpq/pqformat.h"
33
#include "mb/pg_wchar.h"
34
#include "miscadmin.h"
35
#include "pgstat.h"
36
#include "storage/fd.h"
37
#include "tcop/tcopprot.h"
38
#include "utils/json.h"
39
#include "utils/lsyscache.h"
40
#include "utils/memutils.h"
41
#include "utils/rel.h"
42
#include "utils/snapmgr.h"
43
#include "utils/wait_event.h"
44
45
/*
46
 * Represents the different dest cases we need to worry about at
47
 * the bottom level
48
 */
49
typedef enum CopyDest
50
{
51
  COPY_FILE,          /* to file (or a piped program) */
52
  COPY_FRONTEND,        /* to frontend */
53
  COPY_CALLBACK,        /* to callback function */
54
} CopyDest;
55
56
/*
57
 * This struct contains all the state variables used throughout a COPY TO
58
 * operation.
59
 *
60
 * Multi-byte encodings: all supported client-side encodings encode multi-byte
61
 * characters by having the first byte's high bit set. Subsequent bytes of the
62
 * character can have the high bit not set. When scanning data in such an
63
 * encoding to look for a match to a single-byte (ie ASCII) character, we must
64
 * use the full pg_encoding_mblen() machinery to skip over multibyte
65
 * characters, else we might find a false match to a trailing byte. In
66
 * supported server encodings, there is no possibility of a false match, and
67
 * it's faster to make useless comparisons to trailing bytes than it is to
68
 * invoke pg_encoding_mblen() to skip over them. encoding_embeds_ascii is true
69
 * when we have to do it the hard way.
70
 */
71
typedef struct CopyToStateData
72
{
73
  /* format-specific routines */
74
  const CopyToRoutine *routine;
75
76
  /* low-level state data */
77
  CopyDest  copy_dest;    /* type of copy source/destination */
78
  FILE     *copy_file;    /* used if copy_dest == COPY_FILE */
79
  StringInfo  fe_msgbuf;    /* used for all dests during COPY TO */
80
81
  int     file_encoding;  /* file or remote side's character encoding */
82
  bool    need_transcoding; /* file encoding diff from server? */
83
  bool    encoding_embeds_ascii;  /* ASCII can be non-first byte? */
84
85
  /* parameters from the COPY command */
86
  Relation  rel;      /* relation to copy to */
87
  QueryDesc  *queryDesc;    /* executable query to copy from */
88
  List     *attnumlist;   /* integer list of attnums to copy */
89
  char     *filename;   /* filename, or NULL for STDOUT */
90
  bool    is_program;   /* is 'filename' a program to popen? */
91
  bool    json_row_delim_needed;  /* need delimiter before next row */
92
  StringInfo  json_buf;   /* reusable buffer for JSON output,
93
                 * initialized in BeginCopyTo */
94
  TupleDesc tupDesc;    /* Descriptor for JSON output; for a column
95
                 * list this is a projected descriptor */
96
  Datum    *json_projvalues;  /* pre-allocated projection values, or
97
                   * NULL */
98
  bool     *json_projnulls; /* pre-allocated projection nulls, or NULL */
99
  copy_data_dest_cb data_dest_cb; /* function for writing data */
100
101
  CopyFormatOptions opts;
102
  Node     *whereClause;  /* WHERE condition (or NULL) */
103
  List     *partitions;   /* OID list of partitions to copy data from */
104
105
  /*
106
   * Working state
107
   */
108
  MemoryContext copycontext;  /* per-copy execution context */
109
110
  FmgrInfo   *out_functions;  /* lookup info for output functions */
111
  MemoryContext rowcontext; /* per-row evaluation context */
112
  uint64    bytes_processed;  /* number of bytes processed so far */
113
} CopyToStateData;
114
115
/* DestReceiver for COPY (query) TO */
116
typedef struct
117
{
118
  DestReceiver pub;     /* publicly-known function pointers */
119
  CopyToState cstate;     /* CopyToStateData for the command */
120
  uint64    processed;    /* # of tuples processed */
121
} DR_copy;
122
123
/* NOTE: there's a copy of this in copyfromparse.c */
124
static const char BinarySignature[11] = "PGCOPY\n\377\r\n\0";
125
126
127
/* non-export function prototypes */
128
static void EndCopy(CopyToState cstate);
129
static void ClosePipeToProgram(CopyToState cstate);
130
static void CopyOneRowTo(CopyToState cstate, TupleTableSlot *slot);
131
static void CopyAttributeOutText(CopyToState cstate, const char *string);
132
static void CopyAttributeOutCSV(CopyToState cstate, const char *string,
133
                bool use_quote);
134
static void CopyRelationTo(CopyToState cstate, Relation rel, Relation root_rel,
135
               uint64 *processed);
136
137
/* built-in format-specific routines */
138
static void CopyToTextLikeStart(CopyToState cstate, TupleDesc tupDesc);
139
static void CopyToTextLikeOutFunc(CopyToState cstate, Oid atttypid, FmgrInfo *finfo);
140
static void CopyToTextOneRow(CopyToState cstate, TupleTableSlot *slot);
141
static void CopyToCSVOneRow(CopyToState cstate, TupleTableSlot *slot);
142
static void CopyToTextLikeOneRow(CopyToState cstate, TupleTableSlot *slot,
143
                 bool is_csv);
144
static void CopyToTextLikeEnd(CopyToState cstate);
145
static void CopyToJsonOneRow(CopyToState cstate, TupleTableSlot *slot);
146
static void CopyToJsonEnd(CopyToState cstate);
147
static void CopyToBinaryStart(CopyToState cstate, TupleDesc tupDesc);
148
static void CopyToBinaryOutFunc(CopyToState cstate, Oid atttypid, FmgrInfo *finfo);
149
static void CopyToBinaryOneRow(CopyToState cstate, TupleTableSlot *slot);
150
static void CopyToBinaryEnd(CopyToState cstate);
151
152
/* Low-level communications functions */
153
static void SendCopyBegin(CopyToState cstate);
154
static void SendCopyEnd(CopyToState cstate);
155
static void CopySendData(CopyToState cstate, const void *databuf, int datasize);
156
static void CopySendString(CopyToState cstate, const char *str);
157
static void CopySendChar(CopyToState cstate, char c);
158
static void CopySendEndOfRow(CopyToState cstate);
159
static void CopySendTextLikeEndOfRow(CopyToState cstate);
160
static void CopySendInt32(CopyToState cstate, int32 val);
161
static void CopySendInt16(CopyToState cstate, int16 val);
162
163
/*
164
 * COPY TO routines for built-in formats.
165
 */
166
167
/* text format */
168
static const CopyToRoutine CopyToRoutineText = {
169
  .CopyToStart = CopyToTextLikeStart,
170
  .CopyToOutFunc = CopyToTextLikeOutFunc,
171
  .CopyToOneRow = CopyToTextOneRow,
172
  .CopyToEnd = CopyToTextLikeEnd,
173
};
174
175
/* CSV format */
176
static const CopyToRoutine CopyToRoutineCSV = {
177
  .CopyToStart = CopyToTextLikeStart,
178
  .CopyToOutFunc = CopyToTextLikeOutFunc,
179
  .CopyToOneRow = CopyToCSVOneRow,
180
  .CopyToEnd = CopyToTextLikeEnd,
181
};
182
183
/* json format */
184
static const CopyToRoutine CopyToRoutineJson = {
185
  .CopyToStart = CopyToTextLikeStart,
186
  .CopyToOutFunc = CopyToTextLikeOutFunc,
187
  .CopyToOneRow = CopyToJsonOneRow,
188
  .CopyToEnd = CopyToJsonEnd,
189
};
190
191
/* binary format */
192
static const CopyToRoutine CopyToRoutineBinary = {
193
  .CopyToStart = CopyToBinaryStart,
194
  .CopyToOutFunc = CopyToBinaryOutFunc,
195
  .CopyToOneRow = CopyToBinaryOneRow,
196
  .CopyToEnd = CopyToBinaryEnd,
197
};
198
199
/* Return a COPY TO routine for the given options */
200
static const CopyToRoutine *
201
CopyToGetRoutine(const CopyFormatOptions *opts)
202
0
{
203
0
  if (opts->format == COPY_FORMAT_CSV)
204
0
    return &CopyToRoutineCSV;
205
0
  else if (opts->format == COPY_FORMAT_BINARY)
206
0
    return &CopyToRoutineBinary;
207
0
  else if (opts->format == COPY_FORMAT_JSON)
208
0
    return &CopyToRoutineJson;
209
210
  /* default is text */
211
0
  return &CopyToRoutineText;
212
0
}
213
214
/* Implementation of the start callback for text, CSV, and json formats */
215
static void
216
CopyToTextLikeStart(CopyToState cstate, TupleDesc tupDesc)
217
0
{
218
  /*
219
   * For non-binary copy, we need to convert null_print to file encoding,
220
   * because it will be sent directly with CopySendString.
221
   */
222
0
  if (cstate->need_transcoding)
223
0
    cstate->opts.null_print_client = pg_server_to_any(cstate->opts.null_print,
224
0
                              cstate->opts.null_print_len,
225
0
                              cstate->file_encoding);
226
227
  /* if a header has been requested send the line */
228
0
  if (cstate->opts.header_line == COPY_HEADER_TRUE)
229
0
  {
230
0
    ListCell   *cur;
231
0
    bool    hdr_delim = false;
232
233
0
    Assert(cstate->opts.format != COPY_FORMAT_JSON);
234
235
0
    foreach(cur, cstate->attnumlist)
236
0
    {
237
0
      int     attnum = lfirst_int(cur);
238
0
      char     *colname;
239
240
0
      if (hdr_delim)
241
0
        CopySendChar(cstate, cstate->opts.delim[0]);
242
0
      hdr_delim = true;
243
244
0
      colname = NameStr(TupleDescAttr(tupDesc, attnum - 1)->attname);
245
246
0
      if (cstate->opts.format == COPY_FORMAT_CSV)
247
0
        CopyAttributeOutCSV(cstate, colname, false);
248
0
      else
249
0
        CopyAttributeOutText(cstate, colname);
250
0
    }
251
252
0
    CopySendTextLikeEndOfRow(cstate);
253
0
  }
254
255
  /*
256
   * If FORCE_ARRAY has been specified, send the opening bracket.
257
   */
258
0
  if (cstate->opts.format == COPY_FORMAT_JSON && cstate->opts.force_array)
259
0
  {
260
0
    CopySendChar(cstate, '[');
261
0
    CopySendTextLikeEndOfRow(cstate);
262
0
  }
263
0
}
264
265
/*
266
 * Implementation of the outfunc callback for text, CSV, and json formats. Assign
267
 * the output function data to the given *finfo.
268
 */
269
static void
270
CopyToTextLikeOutFunc(CopyToState cstate, Oid atttypid, FmgrInfo *finfo)
271
0
{
272
0
  Oid     func_oid;
273
0
  bool    is_varlena;
274
275
  /* Set output function for an attribute */
276
0
  getTypeOutputInfo(atttypid, &func_oid, &is_varlena);
277
0
  fmgr_info(func_oid, finfo);
278
0
}
279
280
/* Implementation of the per-row callback for text format */
281
static void
282
CopyToTextOneRow(CopyToState cstate, TupleTableSlot *slot)
283
0
{
284
0
  CopyToTextLikeOneRow(cstate, slot, false);
285
0
}
286
287
/* Implementation of the per-row callback for CSV format */
288
static void
289
CopyToCSVOneRow(CopyToState cstate, TupleTableSlot *slot)
290
0
{
291
0
  CopyToTextLikeOneRow(cstate, slot, true);
292
0
}
293
294
/*
295
 * Workhorse for CopyToTextOneRow() and CopyToCSVOneRow().
296
 *
297
 * We use pg_always_inline to reduce function call overhead
298
 * and to help compilers to optimize away the 'is_csv' condition.
299
 */
300
static pg_always_inline void
301
CopyToTextLikeOneRow(CopyToState cstate,
302
           TupleTableSlot *slot,
303
           bool is_csv)
304
0
{
305
0
  bool    need_delim = false;
306
0
  FmgrInfo   *out_functions = cstate->out_functions;
307
308
0
  foreach_int(attnum, cstate->attnumlist)
309
0
  {
310
0
    Datum   value = slot->tts_values[attnum - 1];
311
0
    bool    isnull = slot->tts_isnull[attnum - 1];
312
313
0
    if (need_delim)
314
0
      CopySendChar(cstate, cstate->opts.delim[0]);
315
0
    need_delim = true;
316
317
0
    if (isnull)
318
0
    {
319
0
      CopySendString(cstate, cstate->opts.null_print_client);
320
0
    }
321
0
    else
322
0
    {
323
0
      char     *string;
324
325
0
      string = OutputFunctionCall(&out_functions[attnum - 1],
326
0
                    value);
327
328
0
      if (is_csv)
329
0
        CopyAttributeOutCSV(cstate, string,
330
0
                  cstate->opts.force_quote_flags[attnum - 1]);
331
0
      else
332
0
        CopyAttributeOutText(cstate, string);
333
0
    }
334
0
  }
335
336
0
  CopySendTextLikeEndOfRow(cstate);
337
0
}
338
339
/* Implementation of the end callback for text and CSV formats */
340
static void
341
CopyToTextLikeEnd(CopyToState cstate)
342
0
{
343
  /* Nothing to do here */
344
0
}
345
346
/* Implementation of the end callback for json format */
347
static void
348
CopyToJsonEnd(CopyToState cstate)
349
0
{
350
0
  if (cstate->opts.force_array)
351
0
  {
352
0
    CopySendChar(cstate, ']');
353
0
    CopySendTextLikeEndOfRow(cstate);
354
0
  }
355
0
}
356
357
/* Implementation of per-row callback for json format */
358
static void
359
CopyToJsonOneRow(CopyToState cstate, TupleTableSlot *slot)
360
0
{
361
0
  Datum   rowdata;
362
363
0
  resetStringInfo(cstate->json_buf);
364
365
0
  if (cstate->json_projvalues != NULL)
366
0
  {
367
    /*
368
     * Column list case: project selected column values into sequential
369
     * positions matching the custom TupleDesc, then form a new tuple.
370
     */
371
0
    HeapTuple tup;
372
0
    int     i = 0;
373
374
0
    foreach_int(attnum, cstate->attnumlist)
375
0
    {
376
0
      cstate->json_projvalues[i] = slot->tts_values[attnum - 1];
377
0
      cstate->json_projnulls[i] = slot->tts_isnull[attnum - 1];
378
0
      i++;
379
0
    }
380
381
0
    tup = heap_form_tuple(cstate->tupDesc,
382
0
                cstate->json_projvalues,
383
0
                cstate->json_projnulls);
384
385
    /*
386
     * heap_form_tuple already stamps the datum-length, type-id, and
387
     * type-mod fields on t_data, so we can use it directly as a composite
388
     * Datum without the extra pallocmemcpy that heap_copy_tuple_as_datum
389
     * would do.  Any TOAST pointers in the projected values will be
390
     * detoasted by the per-column output functions called from
391
     * composite_to_json.
392
     */
393
0
    rowdata = HeapTupleGetDatum(tup);
394
0
  }
395
0
  else
396
0
  {
397
    /*
398
     * Full table or query without column list.  For queries, the slot's
399
     * TupleDesc may carry RECORDOID, which is not registered in the type
400
     * cache and would cause composite_to_json's lookup_rowtype_tupdesc
401
     * call to fail.  Build a HeapTuple stamped with the blessed
402
     * descriptor so the type can be looked up correctly.
403
     */
404
0
    if (!cstate->rel && slot->tts_tupleDescriptor->tdtypeid == RECORDOID)
405
0
    {
406
0
      HeapTuple tup = heap_form_tuple(cstate->tupDesc,
407
0
                        slot->tts_values,
408
0
                        slot->tts_isnull);
409
410
0
      rowdata = HeapTupleGetDatum(tup);
411
0
    }
412
0
    else
413
0
      rowdata = ExecFetchSlotHeapTupleDatum(slot);
414
0
  }
415
416
0
  composite_to_json(rowdata, cstate->json_buf, false);
417
418
0
  if (cstate->opts.force_array)
419
0
  {
420
0
    if (cstate->json_row_delim_needed)
421
0
      CopySendChar(cstate, ',');
422
0
    else
423
0
    {
424
      /* first row needs no delimiter */
425
0
      CopySendChar(cstate, ' ');
426
0
      cstate->json_row_delim_needed = true;
427
0
    }
428
0
  }
429
430
  /*
431
   * Convert the JSON output to the target encoding if needed.  Unlike the
432
   * text and CSV paths which convert per-attribute via CopyAttributeOut*,
433
   * composite_to_json() emits the whole row as one buffer, so we transcode
434
   * it here in a single call before sending.
435
   */
436
0
  if (cstate->need_transcoding)
437
0
  {
438
0
    char     *converted;
439
440
0
    converted = pg_server_to_any(cstate->json_buf->data,
441
0
                   cstate->json_buf->len,
442
0
                   cstate->file_encoding);
443
0
    CopySendData(cstate, converted, strlen(converted));
444
0
    if (converted != cstate->json_buf->data)
445
0
      pfree(converted);
446
0
  }
447
0
  else
448
0
    CopySendData(cstate, cstate->json_buf->data, cstate->json_buf->len);
449
450
0
  CopySendTextLikeEndOfRow(cstate);
451
0
}
452
453
/*
454
 * Implementation of the start callback for binary format. Send a header
455
 * for a binary copy.
456
 */
457
static void
458
CopyToBinaryStart(CopyToState cstate, TupleDesc tupDesc)
459
0
{
460
0
  int32   tmp;
461
462
  /* Signature */
463
0
  CopySendData(cstate, BinarySignature, 11);
464
  /* Flags field */
465
0
  tmp = 0;
466
0
  CopySendInt32(cstate, tmp);
467
  /* No header extension */
468
0
  tmp = 0;
469
0
  CopySendInt32(cstate, tmp);
470
0
}
471
472
/*
473
 * Implementation of the outfunc callback for binary format. Assign
474
 * the binary output function to the given *finfo.
475
 */
476
static void
477
CopyToBinaryOutFunc(CopyToState cstate, Oid atttypid, FmgrInfo *finfo)
478
0
{
479
0
  Oid     func_oid;
480
0
  bool    is_varlena;
481
482
  /* Set output function for an attribute */
483
0
  getTypeBinaryOutputInfo(atttypid, &func_oid, &is_varlena);
484
0
  fmgr_info(func_oid, finfo);
485
0
}
486
487
/* Implementation of the per-row callback for binary format */
488
static void
489
CopyToBinaryOneRow(CopyToState cstate, TupleTableSlot *slot)
490
0
{
491
0
  FmgrInfo   *out_functions = cstate->out_functions;
492
493
  /* Binary per-tuple header */
494
0
  CopySendInt16(cstate, list_length(cstate->attnumlist));
495
496
0
  foreach_int(attnum, cstate->attnumlist)
497
0
  {
498
0
    Datum   value = slot->tts_values[attnum - 1];
499
0
    bool    isnull = slot->tts_isnull[attnum - 1];
500
501
0
    if (isnull)
502
0
    {
503
0
      CopySendInt32(cstate, -1);
504
0
    }
505
0
    else
506
0
    {
507
0
      bytea    *outputbytes;
508
509
0
      outputbytes = SendFunctionCall(&out_functions[attnum - 1],
510
0
                       value);
511
0
      CopySendInt32(cstate, VARSIZE(outputbytes) - VARHDRSZ);
512
0
      CopySendData(cstate, VARDATA(outputbytes),
513
0
             VARSIZE(outputbytes) - VARHDRSZ);
514
0
    }
515
0
  }
516
517
0
  CopySendEndOfRow(cstate);
518
0
}
519
520
/* Implementation of the end callback for binary format */
521
static void
522
CopyToBinaryEnd(CopyToState cstate)
523
0
{
524
  /* Generate trailer for a binary copy */
525
0
  CopySendInt16(cstate, -1);
526
  /* Need to flush out the trailer */
527
0
  CopySendEndOfRow(cstate);
528
0
}
529
530
/*
531
 * Send copy start/stop messages for frontend copies.  These have changed
532
 * in past protocol redesigns.
533
 */
534
static void
535
SendCopyBegin(CopyToState cstate)
536
0
{
537
0
  StringInfoData buf;
538
0
  int     natts = list_length(cstate->attnumlist);
539
0
  int16   format = (cstate->opts.format == COPY_FORMAT_BINARY ? 1 : 0);
540
0
  int     i;
541
542
0
  pq_beginmessage(&buf, PqMsg_CopyOutResponse);
543
0
  pq_sendbyte(&buf, format);  /* overall format */
544
0
  if (cstate->opts.format != COPY_FORMAT_JSON)
545
0
  {
546
0
    pq_sendint16(&buf, natts);
547
0
    for (i = 0; i < natts; i++)
548
0
      pq_sendint16(&buf, format); /* per-column formats */
549
0
  }
550
0
  else
551
0
  {
552
    /*
553
     * For JSON format, report one text-format column.  Each CopyData
554
     * message contains one complete JSON object, not individual column
555
     * values, so the per-column count is always 1.
556
     */
557
0
    pq_sendint16(&buf, 1);
558
0
    pq_sendint16(&buf, 0);
559
0
  }
560
561
0
  pq_endmessage(&buf);
562
0
  cstate->copy_dest = COPY_FRONTEND;
563
0
}
564
565
static void
566
SendCopyEnd(CopyToState cstate)
567
0
{
568
  /* Shouldn't have any unsent data */
569
0
  Assert(cstate->fe_msgbuf->len == 0);
570
  /* Send Copy Done message */
571
0
  pq_putemptymessage(PqMsg_CopyDone);
572
0
}
573
574
/*----------
575
 * CopySendData sends output data to the destination (file or frontend)
576
 * CopySendString does the same for null-terminated strings
577
 * CopySendChar does the same for single characters
578
 * CopySendEndOfRow does the appropriate thing at end of each data row
579
 *  (data is not actually flushed except by CopySendEndOfRow)
580
 *
581
 * NB: no data conversion is applied by these functions
582
 *----------
583
 */
584
static void
585
CopySendData(CopyToState cstate, const void *databuf, int datasize)
586
0
{
587
0
  appendBinaryStringInfo(cstate->fe_msgbuf, databuf, datasize);
588
0
}
589
590
static void
591
CopySendString(CopyToState cstate, const char *str)
592
0
{
593
0
  appendBinaryStringInfo(cstate->fe_msgbuf, str, strlen(str));
594
0
}
595
596
static void
597
CopySendChar(CopyToState cstate, char c)
598
0
{
599
0
  appendStringInfoCharMacro(cstate->fe_msgbuf, c);
600
0
}
601
602
static void
603
CopySendEndOfRow(CopyToState cstate)
604
0
{
605
0
  StringInfo  fe_msgbuf = cstate->fe_msgbuf;
606
607
0
  switch (cstate->copy_dest)
608
0
  {
609
0
    case COPY_FILE:
610
0
      pgstat_report_wait_start(WAIT_EVENT_COPY_TO_WRITE);
611
0
      if (fwrite(fe_msgbuf->data, fe_msgbuf->len, 1,
612
0
             cstate->copy_file) != 1 ||
613
0
        ferror(cstate->copy_file))
614
0
      {
615
0
        if (cstate->is_program)
616
0
        {
617
0
          if (errno == EPIPE)
618
0
          {
619
            /*
620
             * The pipe will be closed automatically on error at
621
             * the end of transaction, but we might get a better
622
             * error message from the subprocess' exit code than
623
             * just "Broken Pipe"
624
             */
625
0
            ClosePipeToProgram(cstate);
626
627
            /*
628
             * If ClosePipeToProgram() didn't throw an error, the
629
             * program terminated normally, but closed the pipe
630
             * first. Restore errno, and throw an error.
631
             */
632
0
            errno = EPIPE;
633
0
          }
634
0
          ereport(ERROR,
635
0
              (errcode_for_file_access(),
636
0
               errmsg("could not write to COPY program: %m")));
637
0
        }
638
0
        else
639
0
          ereport(ERROR,
640
0
              (errcode_for_file_access(),
641
0
               errmsg("could not write to COPY file: %m")));
642
0
      }
643
0
      pgstat_report_wait_end();
644
0
      break;
645
0
    case COPY_FRONTEND:
646
      /* Dump the accumulated row as one CopyData message */
647
0
      (void) pq_putmessage(PqMsg_CopyData, fe_msgbuf->data, fe_msgbuf->len);
648
0
      break;
649
0
    case COPY_CALLBACK:
650
0
      cstate->data_dest_cb(fe_msgbuf->data, fe_msgbuf->len);
651
0
      break;
652
0
  }
653
654
  /* Update the progress */
655
0
  cstate->bytes_processed += fe_msgbuf->len;
656
0
  pgstat_progress_update_param(PROGRESS_COPY_BYTES_PROCESSED, cstate->bytes_processed);
657
658
0
  resetStringInfo(fe_msgbuf);
659
0
}
660
661
/*
662
 * Wrapper function of CopySendEndOfRow for text, CSV, and json formats. Sends the
663
 * line termination and do common appropriate things for the end of row.
664
 */
665
static inline void
666
CopySendTextLikeEndOfRow(CopyToState cstate)
667
0
{
668
0
  switch (cstate->copy_dest)
669
0
  {
670
0
    case COPY_FILE:
671
      /* Default line termination depends on platform */
672
0
#ifndef WIN32
673
0
      CopySendChar(cstate, '\n');
674
#else
675
      CopySendString(cstate, "\r\n");
676
#endif
677
0
      break;
678
0
    case COPY_FRONTEND:
679
      /* The FE/BE protocol uses \n as newline for all platforms */
680
0
      CopySendChar(cstate, '\n');
681
0
      break;
682
0
    default:
683
0
      break;
684
0
  }
685
686
  /* Now take the actions related to the end of a row */
687
0
  CopySendEndOfRow(cstate);
688
0
}
689
690
/*
691
 * These functions do apply some data conversion
692
 */
693
694
/*
695
 * CopySendInt32 sends an int32 in network byte order
696
 */
697
static inline void
698
CopySendInt32(CopyToState cstate, int32 val)
699
0
{
700
0
  uint32    buf;
701
702
0
  buf = pg_hton32((uint32) val);
703
0
  CopySendData(cstate, &buf, sizeof(buf));
704
0
}
705
706
/*
707
 * CopySendInt16 sends an int16 in network byte order
708
 */
709
static inline void
710
CopySendInt16(CopyToState cstate, int16 val)
711
0
{
712
0
  uint16    buf;
713
714
0
  buf = pg_hton16((uint16) val);
715
0
  CopySendData(cstate, &buf, sizeof(buf));
716
0
}
717
718
/*
719
 * Closes the pipe to an external program, checking the pclose() return code.
720
 */
721
static void
722
ClosePipeToProgram(CopyToState cstate)
723
0
{
724
0
  int     pclose_rc;
725
726
0
  Assert(cstate->is_program);
727
728
0
  pclose_rc = ClosePipeStream(cstate->copy_file);
729
0
  if (pclose_rc == -1)
730
0
    ereport(ERROR,
731
0
        (errcode_for_file_access(),
732
0
         errmsg("could not close pipe to external command: %m")));
733
0
  else if (pclose_rc != 0)
734
0
  {
735
0
    ereport(ERROR,
736
0
        (errcode(ERRCODE_EXTERNAL_ROUTINE_EXCEPTION),
737
0
         errmsg("program \"%s\" failed",
738
0
            cstate->filename),
739
0
         errdetail_internal("%s", wait_result_to_str(pclose_rc))));
740
0
  }
741
0
}
742
743
/*
744
 * Release resources allocated in a cstate for COPY TO.
745
 */
746
static void
747
EndCopy(CopyToState cstate)
748
0
{
749
0
  if (cstate->is_program)
750
0
  {
751
0
    ClosePipeToProgram(cstate);
752
0
  }
753
0
  else
754
0
  {
755
0
    if (cstate->filename != NULL && FreeFile(cstate->copy_file))
756
0
      ereport(ERROR,
757
0
          (errcode_for_file_access(),
758
0
           errmsg("could not close file \"%s\": %m",
759
0
              cstate->filename)));
760
0
  }
761
762
0
  pgstat_progress_end_command();
763
764
0
  MemoryContextDelete(cstate->copycontext);
765
766
0
  if (cstate->partitions)
767
0
    list_free(cstate->partitions);
768
769
0
  pfree(cstate);
770
0
}
771
772
/*
773
 * Setup CopyToState to read tuples from a table or a query for COPY TO.
774
 *
775
 * 'rel': Relation to be copied
776
 * 'raw_query': Query whose results are to be copied
777
 * 'queryRelId': OID of base relation to convert to a query (for RLS)
778
 * 'filename': Name of server-local file to write, NULL for STDOUT
779
 * 'is_program': true if 'filename' is program to execute
780
 * 'data_dest_cb': Callback that processes the output data
781
 * 'attnamelist': List of char *, columns to include. NIL selects all cols.
782
 * 'options': List of DefElem. See copy_opt_item in gram.y for selections.
783
 *
784
 * Returns a CopyToState, to be passed to DoCopyTo() and related functions.
785
 */
786
CopyToState
787
BeginCopyTo(ParseState *pstate,
788
      Relation rel,
789
      RawStmt *raw_query,
790
      Oid queryRelId,
791
      const char *filename,
792
      bool is_program,
793
      copy_data_dest_cb data_dest_cb,
794
      List *attnamelist,
795
      List *options)
796
0
{
797
0
  CopyToState cstate;
798
0
  bool    pipe = (filename == NULL && data_dest_cb == NULL);
799
0
  TupleDesc tupDesc;
800
0
  int     num_phys_attrs;
801
0
  MemoryContext oldcontext;
802
0
  const int progress_cols[] = {
803
0
    PROGRESS_COPY_COMMAND,
804
0
    PROGRESS_COPY_TYPE
805
0
  };
806
0
  int64   progress_vals[] = {
807
0
    PROGRESS_COPY_COMMAND_TO,
808
0
    0
809
0
  };
810
0
  List     *children = NIL;
811
812
0
  if (rel != NULL && rel->rd_rel->relkind != RELKIND_RELATION)
813
0
  {
814
0
    if (rel->rd_rel->relkind == RELKIND_VIEW)
815
0
      ereport(ERROR,
816
0
          (errcode(ERRCODE_WRONG_OBJECT_TYPE),
817
0
           errmsg("cannot copy from view \"%s\"",
818
0
              RelationGetRelationName(rel)),
819
0
           errhint("Try the COPY (SELECT ...) TO variant.")));
820
0
    else if (rel->rd_rel->relkind == RELKIND_MATVIEW)
821
0
    {
822
0
      if (!RelationIsPopulated(rel))
823
0
        ereport(ERROR,
824
0
            errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
825
0
            errmsg("cannot copy from unpopulated materialized view \"%s\"",
826
0
                 RelationGetRelationName(rel)),
827
0
            errhint("Use the REFRESH MATERIALIZED VIEW command."));
828
0
    }
829
0
    else if (rel->rd_rel->relkind == RELKIND_FOREIGN_TABLE)
830
0
      ereport(ERROR,
831
0
          (errcode(ERRCODE_WRONG_OBJECT_TYPE),
832
0
           errmsg("cannot copy from foreign table \"%s\"",
833
0
              RelationGetRelationName(rel)),
834
0
           errhint("Try the COPY (SELECT ...) TO variant.")));
835
0
    else if (rel->rd_rel->relkind == RELKIND_SEQUENCE)
836
0
      ereport(ERROR,
837
0
          (errcode(ERRCODE_WRONG_OBJECT_TYPE),
838
0
           errmsg("cannot copy from sequence \"%s\"",
839
0
              RelationGetRelationName(rel))));
840
0
    else if (rel->rd_rel->relkind == RELKIND_PARTITIONED_TABLE)
841
0
    {
842
      /*
843
       * Collect OIDs of relation containing data, so that later
844
       * DoCopyTo can copy the data from them.
845
       */
846
0
      children = find_all_inheritors(RelationGetRelid(rel), AccessShareLock, NULL);
847
848
0
      foreach_oid(child, children)
849
0
      {
850
0
        char    relkind = get_rel_relkind(child);
851
852
0
        if (relkind == RELKIND_FOREIGN_TABLE)
853
0
        {
854
0
          char     *relation_name = get_rel_name(child);
855
856
0
          ereport(ERROR,
857
0
              errcode(ERRCODE_WRONG_OBJECT_TYPE),
858
0
              errmsg("cannot copy from foreign table \"%s\"", relation_name),
859
0
              errdetail("Partition \"%s\" is a foreign table in partitioned table \"%s\".",
860
0
                    relation_name, RelationGetRelationName(rel)),
861
0
              errhint("Try the COPY (SELECT ...) TO variant."));
862
0
        }
863
864
        /* Exclude tables with no data */
865
0
        if (RELKIND_HAS_PARTITIONS(relkind))
866
0
          children = foreach_delete_current(children, child);
867
0
      }
868
0
    }
869
0
    else
870
0
      ereport(ERROR,
871
0
          (errcode(ERRCODE_WRONG_OBJECT_TYPE),
872
0
           errmsg("cannot copy from non-table relation \"%s\"",
873
0
              RelationGetRelationName(rel))));
874
0
  }
875
876
877
  /* Allocate workspace and zero all fields */
878
0
  cstate = palloc0_object(CopyToStateData);
879
880
  /*
881
   * We allocate everything used by a cstate in a new memory context. This
882
   * avoids memory leaks during repeated use of COPY in a query.
883
   */
884
0
  cstate->copycontext = AllocSetContextCreate(CurrentMemoryContext,
885
0
                        "COPY",
886
0
                        ALLOCSET_DEFAULT_SIZES);
887
888
0
  oldcontext = MemoryContextSwitchTo(cstate->copycontext);
889
890
  /* Extract options from the statement node tree */
891
0
  ProcessCopyOptions(pstate, &cstate->opts, false /* is_from */ , options);
892
893
  /* Set format routine */
894
0
  cstate->routine = CopyToGetRoutine(&cstate->opts);
895
896
  /* Process the source/target relation or query */
897
0
  if (rel)
898
0
  {
899
0
    Assert(!raw_query);
900
901
0
    cstate->rel = rel;
902
903
0
    tupDesc = RelationGetDescr(cstate->rel);
904
0
    cstate->partitions = children;
905
0
    cstate->tupDesc = tupDesc;
906
0
  }
907
0
  else
908
0
  {
909
0
    List     *rewritten;
910
0
    Query    *query;
911
0
    PlannedStmt *plan;
912
0
    DestReceiver *dest;
913
914
0
    cstate->rel = NULL;
915
0
    cstate->partitions = NIL;
916
917
    /*
918
     * Run parse analysis and rewrite.  Note this also acquires sufficient
919
     * locks on the source table(s).
920
     */
921
0
    rewritten = pg_analyze_and_rewrite_fixedparams(raw_query,
922
0
                             pstate->p_sourcetext, NULL, 0,
923
0
                             NULL);
924
925
    /* check that we got back something we can work with */
926
0
    if (rewritten == NIL)
927
0
    {
928
0
      ereport(ERROR,
929
0
          (errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
930
0
           errmsg("DO INSTEAD NOTHING rules are not supported for COPY")));
931
0
    }
932
0
    else if (list_length(rewritten) > 1)
933
0
    {
934
0
      ListCell   *lc;
935
936
      /* examine queries to determine which error message to issue */
937
0
      foreach(lc, rewritten)
938
0
      {
939
0
        Query    *q = lfirst_node(Query, lc);
940
941
0
        if (q->querySource == QSRC_QUAL_INSTEAD_RULE)
942
0
          ereport(ERROR,
943
0
              (errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
944
0
               errmsg("conditional DO INSTEAD rules are not supported for COPY")));
945
0
        if (q->querySource == QSRC_NON_INSTEAD_RULE)
946
0
          ereport(ERROR,
947
0
              (errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
948
0
               errmsg("DO ALSO rules are not supported for COPY")));
949
0
      }
950
951
0
      ereport(ERROR,
952
0
          (errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
953
0
           errmsg("multi-statement DO INSTEAD rules are not supported for COPY")));
954
0
    }
955
956
0
    query = linitial_node(Query, rewritten);
957
958
    /* The grammar allows SELECT INTO, but we don't support that */
959
0
    if (query->utilityStmt != NULL &&
960
0
      IsA(query->utilityStmt, CreateTableAsStmt))
961
0
      ereport(ERROR,
962
0
          (errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
963
0
           errmsg("COPY (SELECT INTO) is not supported")));
964
965
    /* The only other utility command we could see is NOTIFY */
966
0
    if (query->utilityStmt != NULL)
967
0
      ereport(ERROR,
968
0
          (errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
969
0
           errmsg("COPY query must not be a utility command")));
970
971
    /*
972
     * Similarly the grammar doesn't enforce the presence of a RETURNING
973
     * clause, but this is required here.
974
     */
975
0
    if (query->commandType != CMD_SELECT &&
976
0
      query->returningList == NIL)
977
0
    {
978
0
      Assert(query->commandType == CMD_INSERT ||
979
0
           query->commandType == CMD_UPDATE ||
980
0
           query->commandType == CMD_DELETE ||
981
0
           query->commandType == CMD_MERGE);
982
983
0
      ereport(ERROR,
984
0
          (errcode(ERRCODE_FEATURE_NOT_SUPPORTED),
985
0
           errmsg("COPY query must have a RETURNING clause")));
986
0
    }
987
988
    /* plan the query */
989
0
    plan = pg_plan_query(query, pstate->p_sourcetext,
990
0
               CURSOR_OPT_PARALLEL_OK, NULL, NULL);
991
992
    /*
993
     * With row-level security and a user using "COPY relation TO", we
994
     * have to convert the "COPY relation TO" to a query-based COPY (eg:
995
     * "COPY (SELECT * FROM ONLY relation) TO"), to allow the rewriter to
996
     * add in any RLS clauses.
997
     *
998
     * When this happens, we are passed in the relid of the originally
999
     * found relation (which we have locked).  As the planner will look up
1000
     * the relation again, we double-check here to make sure it found the
1001
     * same one that we have locked.
1002
     */
1003
0
    if (queryRelId != InvalidOid)
1004
0
    {
1005
      /*
1006
       * Note that with RLS involved there may be multiple relations,
1007
       * and while the one we need is almost certainly first, we don't
1008
       * make any guarantees of that in the planner, so check the whole
1009
       * list and make sure we find the original relation.
1010
       */
1011
0
      if (!list_member_oid(plan->relationOids, queryRelId))
1012
0
        ereport(ERROR,
1013
0
            (errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE),
1014
0
             errmsg("relation referenced by COPY statement has changed")));
1015
0
    }
1016
1017
    /*
1018
     * Use a snapshot with an updated command ID to ensure this query sees
1019
     * results of any previously executed queries.
1020
     */
1021
0
    PushCopiedSnapshot(GetActiveSnapshot());
1022
0
    UpdateActiveSnapshotCommandId();
1023
1024
    /* Create dest receiver for COPY OUT */
1025
0
    dest = CreateDestReceiver(DestCopyOut);
1026
0
    ((DR_copy *) dest)->cstate = cstate;
1027
1028
    /* Create a QueryDesc requesting no output */
1029
0
    cstate->queryDesc = CreateQueryDesc(plan, pstate->p_sourcetext,
1030
0
                      GetActiveSnapshot(),
1031
0
                      InvalidSnapshot,
1032
0
                      dest, NULL, NULL, 0);
1033
1034
    /*
1035
     * Call ExecutorStart to prepare the plan for execution.
1036
     *
1037
     * ExecutorStart computes a result tupdesc for us
1038
     */
1039
0
    ExecutorStart(cstate->queryDesc, 0);
1040
1041
0
    tupDesc = cstate->queryDesc->tupDesc;
1042
0
    tupDesc = BlessTupleDesc(tupDesc);
1043
0
    cstate->tupDesc = tupDesc;
1044
0
  }
1045
1046
  /* Generate or convert list of attributes to process */
1047
0
  cstate->attnumlist = CopyGetAttnums(tupDesc, cstate->rel, attnamelist);
1048
1049
  /* Set up JSON-specific state */
1050
0
  if (cstate->opts.format == COPY_FORMAT_JSON)
1051
0
  {
1052
0
    cstate->json_buf = makeStringInfo();
1053
1054
    /*
1055
     * Build a projected TupleDesc describing only the selected columns so
1056
     * that composite_to_json() emits the right column names and types;
1057
     * needed when an explicit column list was given (possibly with a
1058
     * different column order) or when generated columns are excluded from
1059
     * the output.
1060
     */
1061
0
    if (rel && (attnamelist != NIL ||
1062
0
          list_length(cstate->attnumlist) < tupDesc->natts))
1063
0
    {
1064
0
      int     natts = list_length(cstate->attnumlist);
1065
0
      TupleDesc resultDesc;
1066
1067
0
      resultDesc = CreateTemplateTupleDesc(natts);
1068
1069
0
      foreach_int(attnum, cstate->attnumlist)
1070
0
      {
1071
0
        Form_pg_attribute attr = TupleDescAttr(tupDesc, attnum - 1);
1072
1073
0
        TupleDescInitEntry(resultDesc,
1074
0
                   foreach_current_index(attnum) + 1,
1075
0
                   NameStr(attr->attname),
1076
0
                   attr->atttypid,
1077
0
                   attr->atttypmod,
1078
0
                   attr->attndims);
1079
0
      }
1080
1081
0
      TupleDescFinalize(resultDesc);
1082
0
      cstate->tupDesc = BlessTupleDesc(resultDesc);
1083
1084
      /*
1085
       * Pre-allocate arrays for projecting selected column values into
1086
       * sequential positions matching the custom TupleDesc.
1087
       */
1088
0
      cstate->json_projvalues = palloc_array(Datum, natts);
1089
0
      cstate->json_projnulls = palloc_array(bool, natts);
1090
0
    }
1091
0
  }
1092
1093
0
  num_phys_attrs = tupDesc->natts;
1094
1095
  /* Convert FORCE_QUOTE name list to per-column flags, check validity */
1096
0
  cstate->opts.force_quote_flags = palloc0_array(bool, num_phys_attrs);
1097
0
  if (cstate->opts.force_quote_all)
1098
0
  {
1099
0
    MemSet(cstate->opts.force_quote_flags, true, num_phys_attrs * sizeof(bool));
1100
0
  }
1101
0
  else if (cstate->opts.force_quote)
1102
0
  {
1103
0
    List     *attnums;
1104
0
    ListCell   *cur;
1105
1106
0
    attnums = CopyGetAttnums(tupDesc, cstate->rel, cstate->opts.force_quote);
1107
1108
0
    foreach(cur, attnums)
1109
0
    {
1110
0
      int     attnum = lfirst_int(cur);
1111
0
      Form_pg_attribute attr = TupleDescAttr(tupDesc, attnum - 1);
1112
1113
0
      if (!list_member_int(cstate->attnumlist, attnum))
1114
0
        ereport(ERROR,
1115
0
            (errcode(ERRCODE_INVALID_COLUMN_REFERENCE),
1116
        /*- translator: %s is the name of a COPY option, e.g. FORCE_NOT_NULL */
1117
0
             errmsg("%s column \"%s\" not referenced by COPY",
1118
0
                "FORCE_QUOTE", NameStr(attr->attname))));
1119
0
      cstate->opts.force_quote_flags[attnum - 1] = true;
1120
0
    }
1121
0
  }
1122
1123
  /* Use client encoding when ENCODING option is not specified. */
1124
0
  if (cstate->opts.file_encoding < 0)
1125
0
    cstate->file_encoding = pg_get_client_encoding();
1126
0
  else
1127
0
    cstate->file_encoding = cstate->opts.file_encoding;
1128
1129
  /*
1130
   * Set up encoding conversion info if the file and server encodings differ
1131
   * (see also pg_server_to_any).
1132
   */
1133
0
  if (cstate->file_encoding == GetDatabaseEncoding() ||
1134
0
    cstate->file_encoding == PG_SQL_ASCII)
1135
0
    cstate->need_transcoding = false;
1136
0
  else
1137
0
    cstate->need_transcoding = true;
1138
1139
  /* See Multibyte encoding comment above */
1140
0
  cstate->encoding_embeds_ascii = PG_ENCODING_IS_CLIENT_ONLY(cstate->file_encoding);
1141
1142
0
  cstate->copy_dest = COPY_FILE;  /* default */
1143
1144
0
  if (data_dest_cb)
1145
0
  {
1146
0
    progress_vals[1] = PROGRESS_COPY_TYPE_CALLBACK;
1147
0
    cstate->copy_dest = COPY_CALLBACK;
1148
0
    cstate->data_dest_cb = data_dest_cb;
1149
0
  }
1150
0
  else if (pipe)
1151
0
  {
1152
0
    progress_vals[1] = PROGRESS_COPY_TYPE_PIPE;
1153
1154
0
    Assert(!is_program);  /* the grammar does not allow this */
1155
0
    if (whereToSendOutput != DestRemote)
1156
0
      cstate->copy_file = stdout;
1157
0
  }
1158
0
  else
1159
0
  {
1160
0
    cstate->filename = pstrdup(filename);
1161
0
    cstate->is_program = is_program;
1162
1163
0
    if (is_program)
1164
0
    {
1165
0
      progress_vals[1] = PROGRESS_COPY_TYPE_PROGRAM;
1166
0
      cstate->copy_file = OpenPipeStream(cstate->filename, PG_BINARY_W);
1167
0
      if (cstate->copy_file == NULL)
1168
0
        ereport(ERROR,
1169
0
            (errcode_for_file_access(),
1170
0
             errmsg("could not execute command \"%s\": %m",
1171
0
                cstate->filename)));
1172
0
    }
1173
0
    else
1174
0
    {
1175
0
      mode_t    oumask; /* Pre-existing umask value */
1176
0
      struct stat st;
1177
1178
0
      progress_vals[1] = PROGRESS_COPY_TYPE_FILE;
1179
1180
      /*
1181
       * Prevent write to relative path ... too easy to shoot oneself in
1182
       * the foot by overwriting a database file ...
1183
       */
1184
0
      if (!is_absolute_path(filename))
1185
0
        ereport(ERROR,
1186
0
            (errcode(ERRCODE_INVALID_NAME),
1187
0
             errmsg("relative path not allowed for COPY to file")));
1188
1189
0
      oumask = umask(S_IWGRP | S_IWOTH);
1190
0
      PG_TRY();
1191
0
      {
1192
0
        cstate->copy_file = AllocateFile(cstate->filename, PG_BINARY_W);
1193
0
      }
1194
0
      PG_FINALLY();
1195
0
      {
1196
0
        umask(oumask);
1197
0
      }
1198
0
      PG_END_TRY();
1199
0
      if (cstate->copy_file == NULL)
1200
0
      {
1201
        /* copy errno because ereport subfunctions might change it */
1202
0
        int     save_errno = errno;
1203
1204
0
        ereport(ERROR,
1205
0
            (errcode_for_file_access(),
1206
0
             errmsg("could not open file \"%s\" for writing: %m",
1207
0
                cstate->filename),
1208
0
             (save_errno == ENOENT || save_errno == EACCES) ?
1209
0
             errhint("COPY TO instructs the PostgreSQL server process to write a file. "
1210
0
                 "You may want a client-side facility such as psql's \\copy.") : 0));
1211
0
      }
1212
1213
0
      if (fstat(fileno(cstate->copy_file), &st))
1214
0
        ereport(ERROR,
1215
0
            (errcode_for_file_access(),
1216
0
             errmsg("could not stat file \"%s\": %m",
1217
0
                cstate->filename)));
1218
1219
0
      if (S_ISDIR(st.st_mode))
1220
0
        ereport(ERROR,
1221
0
            (errcode(ERRCODE_WRONG_OBJECT_TYPE),
1222
0
             errmsg("\"%s\" is a directory", cstate->filename)));
1223
0
    }
1224
0
  }
1225
1226
  /* initialize progress */
1227
0
  pgstat_progress_start_command(PROGRESS_COMMAND_COPY,
1228
0
                  cstate->rel ? RelationGetRelid(cstate->rel) : InvalidOid);
1229
0
  pgstat_progress_update_multi_param(2, progress_cols, progress_vals);
1230
1231
0
  cstate->bytes_processed = 0;
1232
1233
0
  MemoryContextSwitchTo(oldcontext);
1234
1235
0
  return cstate;
1236
0
}
1237
1238
/*
1239
 * Clean up storage and release resources for COPY TO.
1240
 */
1241
void
1242
EndCopyTo(CopyToState cstate)
1243
0
{
1244
0
  if (cstate->queryDesc != NULL)
1245
0
  {
1246
    /* Close down the query and free resources. */
1247
0
    ExecutorFinish(cstate->queryDesc);
1248
0
    ExecutorEnd(cstate->queryDesc);
1249
0
    FreeQueryDesc(cstate->queryDesc);
1250
0
    PopActiveSnapshot();
1251
0
  }
1252
1253
  /* Clean up storage */
1254
0
  EndCopy(cstate);
1255
0
}
1256
1257
/*
1258
 * Copy from relation or query TO file.
1259
 *
1260
 * Returns the number of rows processed.
1261
 */
1262
uint64
1263
DoCopyTo(CopyToState cstate)
1264
0
{
1265
0
  bool    pipe = (cstate->filename == NULL && cstate->data_dest_cb == NULL);
1266
0
  bool    fe_copy = (pipe && whereToSendOutput == DestRemote);
1267
0
  TupleDesc tupDesc;
1268
0
  int     num_phys_attrs;
1269
0
  ListCell   *cur;
1270
0
  uint64    processed = 0;
1271
1272
0
  if (fe_copy)
1273
0
    SendCopyBegin(cstate);
1274
1275
0
  if (cstate->rel)
1276
0
    tupDesc = RelationGetDescr(cstate->rel);
1277
0
  else
1278
0
    tupDesc = cstate->queryDesc->tupDesc;
1279
0
  num_phys_attrs = tupDesc->natts;
1280
0
  cstate->opts.null_print_client = cstate->opts.null_print; /* default */
1281
1282
  /* We use fe_msgbuf as a per-row buffer regardless of copy_dest */
1283
0
  cstate->fe_msgbuf = makeStringInfo();
1284
1285
  /* Get info about the columns we need to process. */
1286
0
  cstate->out_functions = palloc_array(FmgrInfo, num_phys_attrs);
1287
0
  foreach(cur, cstate->attnumlist)
1288
0
  {
1289
0
    int     attnum = lfirst_int(cur);
1290
0
    Form_pg_attribute attr = TupleDescAttr(tupDesc, attnum - 1);
1291
1292
0
    cstate->routine->CopyToOutFunc(cstate, attr->atttypid,
1293
0
                     &cstate->out_functions[attnum - 1]);
1294
0
  }
1295
1296
  /*
1297
   * Create a temporary memory context that we can reset once per row to
1298
   * recover palloc'd memory.  This avoids any problems with leaks inside
1299
   * datatype output routines, and should be faster than retail pfree's
1300
   * anyway.  (We don't need a whole econtext as CopyFrom does.)
1301
   */
1302
0
  cstate->rowcontext = AllocSetContextCreate(CurrentMemoryContext,
1303
0
                         "COPY TO",
1304
0
                         ALLOCSET_DEFAULT_SIZES);
1305
1306
0
  cstate->routine->CopyToStart(cstate, tupDesc);
1307
1308
0
  if (cstate->rel)
1309
0
  {
1310
    /*
1311
     * If COPY TO source table is a partitioned table, then open each
1312
     * partition and process each individual partition.
1313
     */
1314
0
    if (cstate->rel->rd_rel->relkind == RELKIND_PARTITIONED_TABLE)
1315
0
    {
1316
0
      foreach_oid(child, cstate->partitions)
1317
0
      {
1318
0
        Relation  scan_rel;
1319
1320
        /* We already got the lock in BeginCopyTo */
1321
0
        scan_rel = table_open(child, NoLock);
1322
0
        CopyRelationTo(cstate, scan_rel, cstate->rel, &processed);
1323
0
        table_close(scan_rel, NoLock);
1324
0
      }
1325
0
    }
1326
0
    else
1327
0
      CopyRelationTo(cstate, cstate->rel, NULL, &processed);
1328
0
  }
1329
0
  else
1330
0
  {
1331
    /* run the plan --- the dest receiver will send tuples */
1332
0
    ExecutorRun(cstate->queryDesc, ForwardScanDirection, 0);
1333
0
    processed = ((DR_copy *) cstate->queryDesc->dest)->processed;
1334
0
  }
1335
1336
0
  cstate->routine->CopyToEnd(cstate);
1337
1338
0
  MemoryContextDelete(cstate->rowcontext);
1339
1340
0
  if (fe_copy)
1341
0
    SendCopyEnd(cstate);
1342
1343
0
  return processed;
1344
0
}
1345
1346
/*
1347
 * Scans a single table and exports its rows to the COPY destination.
1348
 *
1349
 * root_rel can be set to the root table of rel if rel is a partition
1350
 * table so that we can send tuples in root_rel's rowtype, which might
1351
 * differ from individual partitions.
1352
 */
1353
static void
1354
CopyRelationTo(CopyToState cstate, Relation rel, Relation root_rel, uint64 *processed)
1355
0
{
1356
0
  TupleTableSlot *slot;
1357
0
  TableScanDesc scandesc;
1358
0
  AttrMap    *map = NULL;
1359
0
  TupleTableSlot *root_slot = NULL;
1360
1361
0
  scandesc = table_beginscan(rel, GetActiveSnapshot(), 0, NULL,
1362
0
                 SO_NONE);
1363
0
  slot = table_slot_create(rel, NULL);
1364
1365
  /*
1366
   * If we are exporting partition data here, we check if converting tuples
1367
   * to the root table's rowtype, because a partition might have column
1368
   * order different than its root table.
1369
   */
1370
0
  if (root_rel != NULL)
1371
0
  {
1372
0
    root_slot = table_slot_create(root_rel, NULL);
1373
0
    map = build_attrmap_by_name_if_req(RelationGetDescr(rel),
1374
0
                       RelationGetDescr(root_rel),
1375
0
                       false);
1376
0
  }
1377
1378
0
  while (table_scan_getnextslot(scandesc, ForwardScanDirection, slot))
1379
0
  {
1380
0
    TupleTableSlot *copyslot;
1381
1382
0
    CHECK_FOR_INTERRUPTS();
1383
1384
0
    if (map != NULL)
1385
0
      copyslot = execute_attr_map_slot(map, slot, root_slot);
1386
0
    else
1387
0
    {
1388
      /* Deconstruct the tuple */
1389
0
      slot_getallattrs(slot);
1390
0
      copyslot = slot;
1391
0
    }
1392
1393
    /* Format and send the data */
1394
0
    CopyOneRowTo(cstate, copyslot);
1395
1396
    /*
1397
     * Increment the number of processed tuples, and report the progress.
1398
     */
1399
0
    pgstat_progress_update_param(PROGRESS_COPY_TUPLES_PROCESSED,
1400
0
                   ++(*processed));
1401
0
  }
1402
1403
0
  ExecDropSingleTupleTableSlot(slot);
1404
1405
0
  if (root_slot != NULL)
1406
0
    ExecDropSingleTupleTableSlot(root_slot);
1407
1408
0
  if (map != NULL)
1409
0
    free_attrmap(map);
1410
1411
0
  table_endscan(scandesc);
1412
0
}
1413
1414
/*
1415
 * Emit one row during DoCopyTo().
1416
 */
1417
static inline void
1418
CopyOneRowTo(CopyToState cstate, TupleTableSlot *slot)
1419
0
{
1420
0
  MemoryContext oldcontext;
1421
1422
0
  MemoryContextReset(cstate->rowcontext);
1423
0
  oldcontext = MemoryContextSwitchTo(cstate->rowcontext);
1424
1425
  /* Make sure the tuple is fully deconstructed */
1426
0
  slot_getallattrs(slot);
1427
1428
0
  cstate->routine->CopyToOneRow(cstate, slot);
1429
1430
0
  MemoryContextSwitchTo(oldcontext);
1431
0
}
1432
1433
/*
1434
 * Send text representation of one attribute, with conversion and escaping
1435
 */
1436
#define DUMPSOFAR() \
1437
0
  do { \
1438
0
    if (ptr > start) \
1439
0
      CopySendData(cstate, start, ptr - start); \
1440
0
  } while (0)
1441
1442
static void
1443
CopyAttributeOutText(CopyToState cstate, const char *string)
1444
0
{
1445
0
  const char *ptr;
1446
0
  const char *start;
1447
0
  char    c;
1448
0
  char    delimc = cstate->opts.delim[0];
1449
1450
0
  if (cstate->need_transcoding)
1451
0
    ptr = pg_server_to_any(string, strlen(string), cstate->file_encoding);
1452
0
  else
1453
0
    ptr = string;
1454
1455
  /*
1456
   * We have to grovel through the string searching for control characters
1457
   * and instances of the delimiter character.  In most cases, though, these
1458
   * are infrequent.  To avoid overhead from calling CopySendData once per
1459
   * character, we dump out all characters between escaped characters in a
1460
   * single call.  The loop invariant is that the data from "start" to "ptr"
1461
   * can be sent literally, but hasn't yet been.
1462
   *
1463
   * We can skip pg_encoding_mblen() overhead when encoding is safe, because
1464
   * in valid backend encodings, extra bytes of a multibyte character never
1465
   * look like ASCII.  This loop is sufficiently performance-critical that
1466
   * it's worth making two copies of it to get the IS_HIGHBIT_SET() test out
1467
   * of the normal safe-encoding path.
1468
   */
1469
0
  if (cstate->encoding_embeds_ascii)
1470
0
  {
1471
0
    start = ptr;
1472
0
    while ((c = *ptr) != '\0')
1473
0
    {
1474
0
      if ((unsigned char) c < (unsigned char) 0x20)
1475
0
      {
1476
        /*
1477
         * \r and \n must be escaped, the others are traditional. We
1478
         * prefer to dump these using the C-like notation, rather than
1479
         * a backslash and the literal character, because it makes the
1480
         * dump file a bit more proof against Microsoftish data
1481
         * mangling.
1482
         */
1483
0
        switch (c)
1484
0
        {
1485
0
          case '\b':
1486
0
            c = 'b';
1487
0
            break;
1488
0
          case '\f':
1489
0
            c = 'f';
1490
0
            break;
1491
0
          case '\n':
1492
0
            c = 'n';
1493
0
            break;
1494
0
          case '\r':
1495
0
            c = 'r';
1496
0
            break;
1497
0
          case '\t':
1498
0
            c = 't';
1499
0
            break;
1500
0
          case '\v':
1501
0
            c = 'v';
1502
0
            break;
1503
0
          default:
1504
            /* If it's the delimiter, must backslash it */
1505
0
            if (c == delimc)
1506
0
              break;
1507
            /* All ASCII control chars are length 1 */
1508
0
            ptr++;
1509
0
            continue; /* fall to end of loop */
1510
0
        }
1511
        /* if we get here, we need to convert the control char */
1512
0
        DUMPSOFAR();
1513
0
        CopySendChar(cstate, '\\');
1514
0
        CopySendChar(cstate, c);
1515
0
        start = ++ptr;  /* do not include char in next run */
1516
0
      }
1517
0
      else if (c == '\\' || c == delimc)
1518
0
      {
1519
0
        DUMPSOFAR();
1520
0
        CopySendChar(cstate, '\\');
1521
0
        start = ptr++;  /* we include char in next run */
1522
0
      }
1523
0
      else if (IS_HIGHBIT_SET(c))
1524
0
        ptr += pg_encoding_mblen(cstate->file_encoding, ptr);
1525
0
      else
1526
0
        ptr++;
1527
0
    }
1528
0
  }
1529
0
  else
1530
0
  {
1531
0
    start = ptr;
1532
0
    while ((c = *ptr) != '\0')
1533
0
    {
1534
0
      if ((unsigned char) c < (unsigned char) 0x20)
1535
0
      {
1536
        /*
1537
         * \r and \n must be escaped, the others are traditional. We
1538
         * prefer to dump these using the C-like notation, rather than
1539
         * a backslash and the literal character, because it makes the
1540
         * dump file a bit more proof against Microsoftish data
1541
         * mangling.
1542
         */
1543
0
        switch (c)
1544
0
        {
1545
0
          case '\b':
1546
0
            c = 'b';
1547
0
            break;
1548
0
          case '\f':
1549
0
            c = 'f';
1550
0
            break;
1551
0
          case '\n':
1552
0
            c = 'n';
1553
0
            break;
1554
0
          case '\r':
1555
0
            c = 'r';
1556
0
            break;
1557
0
          case '\t':
1558
0
            c = 't';
1559
0
            break;
1560
0
          case '\v':
1561
0
            c = 'v';
1562
0
            break;
1563
0
          default:
1564
            /* If it's the delimiter, must backslash it */
1565
0
            if (c == delimc)
1566
0
              break;
1567
            /* All ASCII control chars are length 1 */
1568
0
            ptr++;
1569
0
            continue; /* fall to end of loop */
1570
0
        }
1571
        /* if we get here, we need to convert the control char */
1572
0
        DUMPSOFAR();
1573
0
        CopySendChar(cstate, '\\');
1574
0
        CopySendChar(cstate, c);
1575
0
        start = ++ptr;  /* do not include char in next run */
1576
0
      }
1577
0
      else if (c == '\\' || c == delimc)
1578
0
      {
1579
0
        DUMPSOFAR();
1580
0
        CopySendChar(cstate, '\\');
1581
0
        start = ptr++;  /* we include char in next run */
1582
0
      }
1583
0
      else
1584
0
        ptr++;
1585
0
    }
1586
0
  }
1587
1588
0
  DUMPSOFAR();
1589
0
}
1590
1591
/*
1592
 * Send text representation of one attribute, with conversion and
1593
 * CSV-style escaping
1594
 */
1595
static void
1596
CopyAttributeOutCSV(CopyToState cstate, const char *string,
1597
          bool use_quote)
1598
0
{
1599
0
  const char *ptr;
1600
0
  const char *start;
1601
0
  char    c;
1602
0
  char    delimc = cstate->opts.delim[0];
1603
0
  char    quotec = cstate->opts.quote[0];
1604
0
  char    escapec = cstate->opts.escape[0];
1605
0
  bool    single_attr = (list_length(cstate->attnumlist) == 1);
1606
1607
  /* force quoting if it matches null_print (before conversion!) */
1608
0
  if (!use_quote && strcmp(string, cstate->opts.null_print) == 0)
1609
0
    use_quote = true;
1610
1611
0
  if (cstate->need_transcoding)
1612
0
    ptr = pg_server_to_any(string, strlen(string), cstate->file_encoding);
1613
0
  else
1614
0
    ptr = string;
1615
1616
  /*
1617
   * Make a preliminary pass to discover if it needs quoting
1618
   */
1619
0
  if (!use_quote)
1620
0
  {
1621
    /*
1622
     * Quote '\.' if it appears alone on a line, so that it will not be
1623
     * interpreted as an end-of-data marker.  (PG 18 and up will not
1624
     * interpret '\.' in CSV that way, except in embedded-in-SQL data; but
1625
     * we want the data to be loadable by older versions too.  Also, this
1626
     * avoids breaking clients that are still using PQgetline().)
1627
     */
1628
0
    if (single_attr && strcmp(ptr, "\\.") == 0)
1629
0
      use_quote = true;
1630
0
    else
1631
0
    {
1632
0
      const char *tptr = ptr;
1633
1634
0
      while ((c = *tptr) != '\0')
1635
0
      {
1636
0
        if (c == delimc || c == quotec || c == '\n' || c == '\r')
1637
0
        {
1638
0
          use_quote = true;
1639
0
          break;
1640
0
        }
1641
0
        if (IS_HIGHBIT_SET(c) && cstate->encoding_embeds_ascii)
1642
0
          tptr += pg_encoding_mblen(cstate->file_encoding, tptr);
1643
0
        else
1644
0
          tptr++;
1645
0
      }
1646
0
    }
1647
0
  }
1648
1649
0
  if (use_quote)
1650
0
  {
1651
0
    CopySendChar(cstate, quotec);
1652
1653
    /*
1654
     * We adopt the same optimization strategy as in CopyAttributeOutText
1655
     */
1656
0
    start = ptr;
1657
0
    while ((c = *ptr) != '\0')
1658
0
    {
1659
0
      if (c == quotec || c == escapec)
1660
0
      {
1661
0
        DUMPSOFAR();
1662
0
        CopySendChar(cstate, escapec);
1663
0
        start = ptr;  /* we include char in next run */
1664
0
      }
1665
0
      if (IS_HIGHBIT_SET(c) && cstate->encoding_embeds_ascii)
1666
0
        ptr += pg_encoding_mblen(cstate->file_encoding, ptr);
1667
0
      else
1668
0
        ptr++;
1669
0
    }
1670
0
    DUMPSOFAR();
1671
1672
0
    CopySendChar(cstate, quotec);
1673
0
  }
1674
0
  else
1675
0
  {
1676
    /* If it doesn't need quoting, we can just dump it as-is */
1677
0
    CopySendString(cstate, ptr);
1678
0
  }
1679
0
}
1680
1681
/*
1682
 * copy_dest_startup --- executor startup
1683
 */
1684
static void
1685
copy_dest_startup(DestReceiver *self, int operation, TupleDesc typeinfo)
1686
0
{
1687
  /* no-op */
1688
0
}
1689
1690
/*
1691
 * copy_dest_receive --- receive one tuple
1692
 */
1693
static bool
1694
copy_dest_receive(TupleTableSlot *slot, DestReceiver *self)
1695
0
{
1696
0
  DR_copy    *myState = (DR_copy *) self;
1697
0
  CopyToState cstate = myState->cstate;
1698
1699
  /* Send the data */
1700
0
  CopyOneRowTo(cstate, slot);
1701
1702
  /* Increment the number of processed tuples, and report the progress */
1703
0
  pgstat_progress_update_param(PROGRESS_COPY_TUPLES_PROCESSED,
1704
0
                 ++myState->processed);
1705
1706
0
  return true;
1707
0
}
1708
1709
/*
1710
 * copy_dest_shutdown --- executor end
1711
 */
1712
static void
1713
copy_dest_shutdown(DestReceiver *self)
1714
0
{
1715
  /* no-op */
1716
0
}
1717
1718
/*
1719
 * copy_dest_destroy --- release DestReceiver object
1720
 */
1721
static void
1722
copy_dest_destroy(DestReceiver *self)
1723
0
{
1724
0
  pfree(self);
1725
0
}
1726
1727
/*
1728
 * CreateCopyDestReceiver -- create a suitable DestReceiver object
1729
 */
1730
DestReceiver *
1731
CreateCopyDestReceiver(void)
1732
0
{
1733
0
  DR_copy    *self = palloc_object(DR_copy);
1734
1735
0
  self->pub.receiveSlot = copy_dest_receive;
1736
0
  self->pub.rStartup = copy_dest_startup;
1737
0
  self->pub.rShutdown = copy_dest_shutdown;
1738
0
  self->pub.rDestroy = copy_dest_destroy;
1739
0
  self->pub.mydest = DestCopyOut;
1740
1741
0
  self->cstate = NULL;   /* will be set later */
1742
0
  self->processed = 0;
1743
1744
0
  return (DestReceiver *) self;
1745
0
}