/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 | } |