/src/postgres/src/backend/postmaster/datachecksum_state.c
Line | Count | Source |
1 | | /*------------------------------------------------------------------------- |
2 | | * |
3 | | * datachecksum_state.c |
4 | | * Background worker for enabling or disabling data checksums online as |
5 | | * well as functionality for manipulating data checksum state |
6 | | * |
7 | | * When enabling data checksums on a cluster at initdb time or when shut down |
8 | | * with pg_checksums, no extra process is required as each page is checksummed, |
9 | | * and verified, when accessed. When enabling checksums on an already running |
10 | | * cluster, this worker will ensure that all pages are checksummed before |
11 | | * verification of the checksums is turned on. In the case of disabling |
12 | | * checksums, the state transition is performed only in the control file, no |
13 | | * changes are performed on the data pages. |
14 | | * |
15 | | * Checksums can be either enabled or disabled cluster-wide, with on/off being |
16 | | * the end state for data_checksums. |
17 | | * |
18 | | * 1. Enabling checksums |
19 | | * --------------------- |
20 | | * When enabling checksums in an online cluster, data_checksums will be set to |
21 | | * "inprogress-on" which signals that write operations MUST compute and write |
22 | | * the checksum on the data page, but during reading the checksum SHALL NOT be |
23 | | * verified. This ensures that all objects created while checksums are being |
24 | | * enabled will have checksums set, but reads won't fail due to missing or |
25 | | * invalid checksums. Invalid checksums can be present in case the cluster had |
26 | | * checksums enabled, then disabled them and updated the page while they were |
27 | | * disabled. |
28 | | * |
29 | | * The DataChecksumsWorker will compile a list of all databases at the start, |
30 | | * any databases created concurrently will see the in-progress state and will |
31 | | * be checksummed automatically. All databases from the original list MUST BE |
32 | | * successfully processed in order for data checksums to be enabled, the only |
33 | | * exception are databases which are dropped before having been processed. |
34 | | * |
35 | | * For each database, all relations which have storage are read and every data |
36 | | * page is marked dirty to force a write with the checksum. This will generate |
37 | | * a lot of WAL as the entire database is read and written. |
38 | | * |
39 | | * If the processing is interrupted by a cluster crash or restart, it needs to |
40 | | * be restarted from the beginning again as state isn't persisted. |
41 | | * |
42 | | * 2. Disabling checksums |
43 | | * ---------------------- |
44 | | * When disabling checksums, data_checksums will be set to "inprogress-off" |
45 | | * which signals that checksums are written but no longer need to be verified. |
46 | | * This ensures that backends which have not yet transitioned to the |
47 | | * "inprogress-off" state will still see valid checksums on pages. |
48 | | * |
49 | | * 3. Synchronization and Correctness |
50 | | * ---------------------------------- |
51 | | * The processes involved in enabling or disabling data checksums in an |
52 | | * online cluster must be properly synchronized with the normal backends |
53 | | * serving concurrent queries to ensure correctness. Correctness is defined |
54 | | * as the following: |
55 | | * |
56 | | * - Backends SHALL NOT violate the data_checksums state they have agreed to |
57 | | * by acknowledging the procsignalbarrier: This means that all backends |
58 | | * MUST calculate and write data checksums during all states except off; |
59 | | * MUST validate checksums only in the 'on' state. |
60 | | * - Data checksums SHALL NOT be considered enabled cluster-wide until all |
61 | | * currently connected backends have state "on": This means that all |
62 | | * backends must wait on the procsignalbarrier to be acknowledged by all |
63 | | * before proceeding to validate data checksums. |
64 | | * |
65 | | * There are two steps of synchronization required for changing data_checksums |
66 | | * in an online cluster: (i) changing state in the active backends ("on", |
67 | | * "off", "inprogress-on" and "inprogress-off"), and (ii) ensuring no |
68 | | * incompatible objects and processes are left in a database when workers end. |
69 | | * The former deals with cluster-wide agreement on data checksum state and the |
70 | | * latter with ensuring that any concurrent activity cannot break the data |
71 | | * checksum contract during processing. |
72 | | * |
73 | | * Synchronizing the state change is done with procsignal barriers. Before |
74 | | * updating the data_checksums state in the control file, all other backends must absorb the |
75 | | * barrier. Barrier absorption will happen during interrupt processing, which |
76 | | * means that connected backends will change state at different times. If |
77 | | * waiting for a barrier is done during startup, for example during replay, it |
78 | | * is important to realize that any locks held by the startup process might |
79 | | * cause deadlocks if backends end up waiting for those locks while startup |
80 | | * is waiting for a procsignalbarrier. |
81 | | * |
82 | | * 3.1 When Enabling Data Checksums |
83 | | * -------------------------------- |
84 | | * A process which fails to observe data checksums being enabled can induce two |
85 | | * types of errors: failing to write the checksum when modifying the page and |
86 | | * failing to validate the data checksum on the page when reading it. |
87 | | * |
88 | | * When processing starts all backends belong to one of the below sets, with |
89 | | * one of Bd and Bi being empty: |
90 | | * |
91 | | * Bg: Backend updating the global state and emitting the procsignalbarrier |
92 | | * Bd: Backends in "off" state |
93 | | * Bi: Backends in "inprogress-on" state |
94 | | * |
95 | | * If processing is started in an online cluster then all backends are in Bd. |
96 | | * If processing was halted by the cluster shutting down (due to a crash or |
97 | | * intentional restart), the control file state "inprogress-on" will be |
98 | | * observed on system startup and all backends will be placed in Bd. The |
99 | | * control file state will also be set to "off". |
100 | | * |
101 | | * Backends transition Bd -> Bi via a procsignalbarrier which is emitted by the |
102 | | * DataChecksumsWorkerLauncherMain. When all backends have acknowledged the |
103 | | * barrier then Bd will be empty and the next phase can begin: calculating and |
104 | | * writing data checksums with DataChecksumsWorkers. When the |
105 | | * DataChecksumsWorker processes have finished writing checksums on all pages, |
106 | | * data checksums are enabled cluster-wide via another procsignalbarrier. |
107 | | * There are four sets of backends where Bd shall be an empty set: |
108 | | * |
109 | | * Bg: Backend updating the global state and emitting the procsignalbarrier |
110 | | * Bd: Backends in "off" state |
111 | | * Be: Backends in "on" state |
112 | | * Bi: Backends in "inprogress-on" state |
113 | | * |
114 | | * Backends in Bi and Be will write checksums when modifying a page, but only |
115 | | * backends in Be will verify the checksum during reading. The Bg backend is |
116 | | * blocked waiting for all backends in Bi to process interrupts and move to |
117 | | * Be. Any backend starting while Bg is waiting on the procsignalbarrier will |
118 | | * observe the global state being "on" and will thus automatically belong to |
119 | | * Be. Checksums are enabled cluster-wide when Bi is an empty set. Bi and Be |
120 | | * are compatible sets while still operating based on their local state as |
121 | | * both write data checksums. |
122 | | * |
123 | | * 3.2 When Disabling Data Checksums |
124 | | * --------------------------------- |
125 | | * A process which fails to observe that data checksums have been disabled |
126 | | * can induce two types of errors: writing the checksum when modifying the |
127 | | * page and validating a data checksum which is no longer correct due to |
128 | | * modifications to the page. The former is not an error per se as data |
129 | | * integrity is maintained, but it is wasteful. The latter will cause errors |
130 | | * in user operations. Assuming the following sets of backends: |
131 | | * |
132 | | * Bg: Backend updating the global state and emitting the procsignalbarrier |
133 | | * Bd: Backends in "off" state |
134 | | * Be: Backends in "on" state |
135 | | * Bo: Backends in "inprogress-off" state |
136 | | * Bi: Backends in "inprogress-on" state |
137 | | * |
138 | | * Backends transition from the Be state to Bd like so: Be -> Bo -> Bd. From |
139 | | * all other states, the transition can be straight to Bd. |
140 | | * |
141 | | * The goal is to transition all backends to Bd making the others empty sets. |
142 | | * Backends in Bo write data checksums, but don't validate them, such that |
143 | | * backends still in Be can continue to validate pages until the barrier has |
144 | | * been absorbed such that they are in Bo. Once all backends are in Bo, the |
145 | | * barrier to transition to "off" can be raised and all backends can safely |
146 | | * stop writing data checksums as no backend is enforcing data checksum |
147 | | * validation any longer. |
148 | | * |
149 | | * 4. Interaction with offline data checksum changes |
150 | | * ------------------------------------------------- |
151 | | * Enabling or disabling checksums offline with pg_checksums uses none of the |
152 | | * machinery in this file, but the two mechanisms share the state kept in the |
153 | | * control file, so their interaction is documented here. |
154 | | * |
155 | | * pg_checksums writes the new state to the control file and sets |
156 | | * data_checksum_is_local, marking a state that no WAL record accounts for. |
157 | | * Recovery then does not adopt the state carried by a replayed checkpoint |
158 | | * record over it. The control file also carries a watermark, the WAL |
159 | | * position through which data checksum transitions are covered. Replay skips |
160 | | * transition records ending at or below the watermark, as their effect is |
161 | | * already contained in the control file, and applies records above it as |
162 | | * usual, whether they were written before or after an offline change. This |
163 | | * is why an offline change in a replicated setup must be made on every node |
164 | | * while all of them are stopped and caught up; see the pg_checksums |
165 | | * documentation for the procedure. |
166 | | * |
167 | | * 5. Future opportunities for optimizations |
168 | | * ----------------------------------------- |
169 | | * Below are some potential optimizations and improvements which were brought |
170 | | * up during reviews of this feature, but which weren't implemented in the |
171 | | * initial version. These are ideas listed without any validation on their |
172 | | * feasibility or potential payoff. More discussion on (most of) these can be |
173 | | * found on the -hackers threads linked to in the commit message of this |
174 | | * feature. |
175 | | * |
176 | | * * Launching datachecksumsworker for resuming operation from the startup |
177 | | * process: Currently users have to restart processing manually after a |
178 | | * restart since dynamic background worker cannot be started from the |
179 | | * postmaster. Changing the startup process could make restarting the |
180 | | * processing automatic on cluster restart. |
181 | | * * Avoid dirtying the page when checksums already match: Iff the checksum |
182 | | * on the page happens to already match we still dirty the page. It should |
183 | | * be enough to only do the log_newpage_buffer() call in that case. |
184 | | * * Teach pg_checksums to avoid checksummed pages when pg_checksums is used |
185 | | * to enable checksums on a cluster which is in inprogress-on state and |
186 | | * may have checksummed pages (make pg_checksums be able to resume an |
187 | | * online operation). This should only be attempted for wal_level minimal. |
188 | | * * Restartability (not necessarily with page granularity). |
189 | | * * Avoid processing databases which were created during inprogress-on. |
190 | | * Right now all databases are processed regardless to be safe. |
191 | | * * Teach CREATE DATABASE to calculate checksums for databases created |
192 | | * during inprogress-on with a template database which has yet to be |
193 | | * processed. |
194 | | * |
195 | | * |
196 | | * Portions Copyright (c) 1996-2026, PostgreSQL Global Development Group |
197 | | * Portions Copyright (c) 1994, Regents of the University of California |
198 | | * |
199 | | * |
200 | | * IDENTIFICATION |
201 | | * src/backend/postmaster/datachecksum_state.c |
202 | | * |
203 | | *------------------------------------------------------------------------- |
204 | | */ |
205 | | #include "postgres.h" |
206 | | |
207 | | #include "access/genam.h" |
208 | | #include "access/heapam.h" |
209 | | #include "access/htup_details.h" |
210 | | #include "access/xact.h" |
211 | | #include "access/xlog.h" |
212 | | #include "access/xloginsert.h" |
213 | | #include "catalog/indexing.h" |
214 | | #include "catalog/pg_class.h" |
215 | | #include "catalog/pg_database.h" |
216 | | #include "commands/progress.h" |
217 | | #include "commands/vacuum.h" |
218 | | #include "common/relpath.h" |
219 | | #include "miscadmin.h" |
220 | | #include "pgstat.h" |
221 | | #include "postmaster/bgworker.h" |
222 | | #include "postmaster/bgwriter.h" |
223 | | #include "postmaster/datachecksum_state.h" |
224 | | #include "storage/bufmgr.h" |
225 | | #include "storage/checksum.h" |
226 | | #include "storage/ipc.h" |
227 | | #include "storage/latch.h" |
228 | | #include "storage/lmgr.h" |
229 | | #include "storage/lwlock.h" |
230 | | #include "storage/procarray.h" |
231 | | #include "storage/smgr.h" |
232 | | #include "storage/subsystems.h" |
233 | | #include "tcop/tcopprot.h" |
234 | | #include "utils/builtins.h" |
235 | | #include "utils/fmgroids.h" |
236 | | #include "utils/injection_point.h" |
237 | | #include "utils/lsyscache.h" |
238 | | #include "utils/ps_status.h" |
239 | | #include "utils/syscache.h" |
240 | | #include "utils/wait_event.h" |
241 | | |
242 | | /* |
243 | | * Configuration of conditions which must match when absorbing a procsignal |
244 | | * barrier during data checksum enable/disable operations. A single function |
245 | | * is used for absorbing all barriers, and the current and target states must |
246 | | * be defined as a from/to tuple in the checksum_barriers struct. |
247 | | */ |
248 | | typedef struct ChecksumBarrierCondition |
249 | | { |
250 | | /* Current state of data checksums */ |
251 | | int from; |
252 | | /* Target state for data checksums */ |
253 | | int to; |
254 | | } ChecksumBarrierCondition; |
255 | | |
256 | | static const ChecksumBarrierCondition checksum_barriers[9] = |
257 | | { |
258 | | /* |
259 | | * Disabling checksums: If checksums are currently enabled, disabling must |
260 | | * go through the 'inprogress-off' state. |
261 | | */ |
262 | | {PG_DATA_CHECKSUM_VERSION, PG_DATA_CHECKSUM_INPROGRESS_OFF}, |
263 | | {PG_DATA_CHECKSUM_INPROGRESS_OFF, PG_DATA_CHECKSUM_OFF}, |
264 | | |
265 | | /* |
266 | | * If checksums are in the process of being enabled, but are not yet being |
267 | | * verified, we can abort by going back to 'off' state. |
268 | | */ |
269 | | {PG_DATA_CHECKSUM_INPROGRESS_ON, PG_DATA_CHECKSUM_OFF}, |
270 | | |
271 | | /* |
272 | | * Enabling checksums must normally go through the 'inprogress-on' state. |
273 | | */ |
274 | | {PG_DATA_CHECKSUM_OFF, PG_DATA_CHECKSUM_INPROGRESS_ON}, |
275 | | {PG_DATA_CHECKSUM_INPROGRESS_ON, PG_DATA_CHECKSUM_VERSION}, |
276 | | |
277 | | /* |
278 | | * If checksums are being disabled but all backends are still computing |
279 | | * checksums, we can go straight back to 'on' |
280 | | */ |
281 | | {PG_DATA_CHECKSUM_INPROGRESS_OFF, PG_DATA_CHECKSUM_VERSION}, |
282 | | |
283 | | /* |
284 | | * If checksums are being enabled when launcher_exit is executed, state is |
285 | | * set to off since we cannot reach on at that point. |
286 | | */ |
287 | | {PG_DATA_CHECKSUM_INPROGRESS_ON, PG_DATA_CHECKSUM_INPROGRESS_OFF}, |
288 | | |
289 | | /* |
290 | | * Transitions that can happen when a new request is made while another is |
291 | | * currently being processed. |
292 | | */ |
293 | | {PG_DATA_CHECKSUM_INPROGRESS_OFF, PG_DATA_CHECKSUM_INPROGRESS_ON}, |
294 | | {PG_DATA_CHECKSUM_OFF, PG_DATA_CHECKSUM_INPROGRESS_OFF}, |
295 | | }; |
296 | | |
297 | | /* Possible operations the DataChecksumsWorker can perform */ |
298 | | typedef enum DataChecksumsWorkerOperation |
299 | | { |
300 | | ENABLE_DATACHECKSUMS, |
301 | | DISABLE_DATACHECKSUMS, |
302 | | } DataChecksumsWorkerOperation; |
303 | | |
304 | | /* Possible states for a database entry which has been processed */ |
305 | | typedef enum |
306 | | { |
307 | | DATACHECKSUMSWORKER_SUCCESSFUL = 0, |
308 | | DATACHECKSUMSWORKER_ABORTED, |
309 | | DATACHECKSUMSWORKER_FAILED, |
310 | | DATACHECKSUMSWORKER_DROPDB, |
311 | | } DataChecksumsWorkerResult; |
312 | | |
313 | | /* |
314 | | * Signaling between backends calling pg_enable/disable_data_checksums, the |
315 | | * checksums launcher process, and the checksums worker process. |
316 | | * |
317 | | * This struct is protected by DataChecksumsWorkerLock |
318 | | */ |
319 | | typedef struct DataChecksumsStateStruct |
320 | | { |
321 | | /* |
322 | | * These are set by pg_{enable|disable}_data_checksums, to tell the |
323 | | * launcher what the target state is. |
324 | | */ |
325 | | DataChecksumsWorkerOperation launch_operation; |
326 | | int launch_cost_delay; |
327 | | int launch_cost_limit; |
328 | | |
329 | | /* |
330 | | * Is a launcher process currently running? This is set by the main |
331 | | * launcher process, after it has read the above launch_* parameters. |
332 | | */ |
333 | | bool launcher_running; |
334 | | |
335 | | /* |
336 | | * Every time a new worker is launched, it's assigned a unique invocation |
337 | | * number by incrementing this counter. |
338 | | */ |
339 | | uint64 worker_invocation_counter; |
340 | | |
341 | | /* |
342 | | * Information about the current worker, if it's currently running. These |
343 | | * are set by the worker launcher. |
344 | | */ |
345 | | uint64 worker_invocation; /* unique invocation number */ |
346 | | Oid database_oid; /* database it's processing */ |
347 | | pid_t worker_pid; /* worker process's PID */ |
348 | | |
349 | | /* |
350 | | * These fields indicate the target state that the worker is currently |
351 | | * running with. They can be different from the corresponding launch_* |
352 | | * fields, if a new pg_enable/disable_data_checksums() call was made while |
353 | | * the launcher/worker was already running. The worker will periodically |
354 | | * check if new cost settings have been requested, and if so will copy |
355 | | * them from the launch_* fields and reset cost throttling to match the |
356 | | * new values. |
357 | | */ |
358 | | DataChecksumsWorkerOperation operation; |
359 | | int cost_delay; |
360 | | int cost_limit; |
361 | | |
362 | | /* |
363 | | * Signaling between the launcher and the worker process. Protected by |
364 | | * DataChecksumsWorkerLock. |
365 | | */ |
366 | | |
367 | | /* result, set by worker before exiting */ |
368 | | DataChecksumsWorkerResult worker_result; |
369 | | |
370 | | /* |
371 | | * Tells the worker process whether it should also process the shared |
372 | | * catalogs |
373 | | */ |
374 | | bool process_shared_catalogs; |
375 | | } DataChecksumsStateStruct; |
376 | | |
377 | | /* Shared memory segment for datachecksumsworker */ |
378 | | static DataChecksumsStateStruct *DataChecksumState; |
379 | | |
380 | | typedef struct DataChecksumsWorkerDatabase |
381 | | { |
382 | | Oid dboid; |
383 | | char *dbname; |
384 | | } DataChecksumsWorkerDatabase; |
385 | | |
386 | | /* Flag set by the interrupt handler */ |
387 | | static volatile sig_atomic_t abort_requested = false; |
388 | | |
389 | | static uint64 worker_invocation; |
390 | | |
391 | | /* |
392 | | * Have we set the DataChecksumsStateStruct->launcher_running flag? |
393 | | * If we have, we need to clear it before exiting! |
394 | | */ |
395 | | static volatile sig_atomic_t launcher_running = false; |
396 | | |
397 | | /* Are we enabling data checksums, or disabling them? */ |
398 | | static DataChecksumsWorkerOperation operation; |
399 | | |
400 | | /* Prototypes */ |
401 | | static void StartDataChecksumsWorkerLauncher(DataChecksumsWorkerOperation op, |
402 | | int cost_delay, |
403 | | int cost_limit); |
404 | | static void ErrorOnInvalidDatabases(void); |
405 | | static bool ProcessSingleRelationFork(Relation reln, ForkNumber forkNum, BufferAccessStrategy strategy); |
406 | | static void ResetDataChecksumsProgressCounters(void); |
407 | | static bool ProcessSingleRelationByOid(Oid relationId, BufferAccessStrategy strategy); |
408 | | static BgwHandleStatus WaitForDataChecksumsWorkerState(BackgroundWorkerHandle *handle, |
409 | | bool wait_for_startup, |
410 | | pid_t *pidp, |
411 | | uint32 wait_event); |
412 | | static DataChecksumsWorkerResult ProcessDatabase(DataChecksumsWorkerDatabase *db); |
413 | | static void launcher_exit(int code, Datum arg); |
414 | | static void launcher_cancel_handler(SIGNAL_ARGS); |
415 | | static void WaitForAllTransactionsToFinish(void); |
416 | | static bool ProcessAllDatabases(void); |
417 | | static void DataChecksumsShmemRequest(void *arg); |
418 | | static bool DatabaseExists(Oid dboid); |
419 | | static List *BuildDatabaseList(void); |
420 | | static void FreeDatabaseList(List *dblist); |
421 | | static List *BuildRelationList(bool temp_relations, bool include_shared); |
422 | | |
423 | | const ShmemCallbacks DataChecksumsShmemCallbacks = { |
424 | | .request_fn = DataChecksumsShmemRequest, |
425 | | }; |
426 | | |
427 | | #define CHECK_FOR_LAUNCHER_ABORT_REQUEST() \ |
428 | 0 | do { \ |
429 | 0 | Assert(MyBackendType == B_DATACHECKSUMSWORKER_LAUNCHER); \ |
430 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_SHARED); \ |
431 | 0 | if (DataChecksumState->launch_operation != operation) \ |
432 | 0 | abort_requested = true; \ |
433 | 0 | LWLockRelease(DataChecksumsWorkerLock); \ |
434 | 0 | } while (0) |
435 | | |
436 | | #define CHECK_FOR_WORKER_ABORT_REQUEST() \ |
437 | 0 | do { \ |
438 | 0 | Assert(MyBackendType == B_DATACHECKSUMSWORKER_WORKER); \ |
439 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_SHARED); \ |
440 | 0 | if (DataChecksumState->worker_invocation != worker_invocation || \ |
441 | 0 | DataChecksumState->launch_operation != operation) \ |
442 | 0 | abort_requested = true; \ |
443 | 0 | LWLockRelease(DataChecksumsWorkerLock); \ |
444 | 0 | } while (0) |
445 | | |
446 | | |
447 | | /***************************************************************************** |
448 | | * Functionality for manipulating the data checksum state in the cluster |
449 | | */ |
450 | | |
451 | | void |
452 | | EmitAndWaitDataChecksumsBarrier(uint32 state) |
453 | 0 | { |
454 | 0 | uint64 barrier; |
455 | |
|
456 | 0 | switch (state) |
457 | 0 | { |
458 | 0 | case PG_DATA_CHECKSUM_INPROGRESS_ON: |
459 | 0 | barrier = EmitProcSignalBarrier(PROCSIGNAL_BARRIER_CHECKSUM_INPROGRESS_ON); |
460 | 0 | WaitForProcSignalBarrier(barrier); |
461 | 0 | break; |
462 | | |
463 | 0 | case PG_DATA_CHECKSUM_INPROGRESS_OFF: |
464 | 0 | barrier = EmitProcSignalBarrier(PROCSIGNAL_BARRIER_CHECKSUM_INPROGRESS_OFF); |
465 | 0 | WaitForProcSignalBarrier(barrier); |
466 | 0 | break; |
467 | | |
468 | 0 | case PG_DATA_CHECKSUM_VERSION: |
469 | 0 | barrier = EmitProcSignalBarrier(PROCSIGNAL_BARRIER_CHECKSUM_ON); |
470 | 0 | WaitForProcSignalBarrier(barrier); |
471 | 0 | break; |
472 | | |
473 | 0 | case PG_DATA_CHECKSUM_OFF: |
474 | 0 | barrier = EmitProcSignalBarrier(PROCSIGNAL_BARRIER_CHECKSUM_OFF); |
475 | 0 | WaitForProcSignalBarrier(barrier); |
476 | 0 | break; |
477 | | |
478 | 0 | default: |
479 | 0 | Assert(false); |
480 | 0 | } |
481 | 0 | } |
482 | | |
483 | | /* |
484 | | * AbsorbDataChecksumsBarrier |
485 | | * Generic function for absorbing data checksum state changes |
486 | | * |
487 | | * All procsignalbarriers regarding data checksum state changes are absorbed |
488 | | * with this function. The set of conditions required for the state change to |
489 | | * be accepted are listed in the checksum_barriers struct, target_state is |
490 | | * used to look up the relevant entry. |
491 | | */ |
492 | | bool |
493 | | AbsorbDataChecksumsBarrier(ProcSignalBarrierType barrier) |
494 | 0 | { |
495 | 0 | uint32 target_state; |
496 | 0 | int current = data_checksums; |
497 | 0 | bool found = false; |
498 | | |
499 | | /* |
500 | | * Translate the barrier condition to the target state, doing it here |
501 | | * instead of in the procsignal code saves the latter from knowing about |
502 | | * checksum states. |
503 | | */ |
504 | 0 | switch (barrier) |
505 | 0 | { |
506 | 0 | case PROCSIGNAL_BARRIER_CHECKSUM_INPROGRESS_ON: |
507 | 0 | target_state = PG_DATA_CHECKSUM_INPROGRESS_ON; |
508 | 0 | break; |
509 | 0 | case PROCSIGNAL_BARRIER_CHECKSUM_ON: |
510 | 0 | target_state = PG_DATA_CHECKSUM_VERSION; |
511 | 0 | break; |
512 | 0 | case PROCSIGNAL_BARRIER_CHECKSUM_INPROGRESS_OFF: |
513 | 0 | target_state = PG_DATA_CHECKSUM_INPROGRESS_OFF; |
514 | 0 | break; |
515 | 0 | case PROCSIGNAL_BARRIER_CHECKSUM_OFF: |
516 | 0 | target_state = PG_DATA_CHECKSUM_OFF; |
517 | 0 | break; |
518 | 0 | default: |
519 | 0 | elog(ERROR, "incorrect barrier \"%d\" received", barrier); |
520 | 0 | } |
521 | | |
522 | | /* |
523 | | * If the target state matches the current state then the barrier has been |
524 | | * repeated. |
525 | | */ |
526 | 0 | if (current == target_state) |
527 | 0 | return true; |
528 | | |
529 | | /* |
530 | | * If the cluster is in recovery we skip the validation of current state |
531 | | * since the replay is trusted. |
532 | | */ |
533 | 0 | if (RecoveryInProgress()) |
534 | 0 | { |
535 | 0 | SetLocalDataChecksumState(target_state); |
536 | 0 | return true; |
537 | 0 | } |
538 | | |
539 | | /* |
540 | | * Find the barrier condition definition for the target state. Not finding |
541 | | * a condition would be a grave programmer error as the states are a |
542 | | * discrete set. |
543 | | */ |
544 | 0 | for (size_t i = 0; i < lengthof(checksum_barriers) && !found; i++) |
545 | 0 | { |
546 | 0 | if (checksum_barriers[i].from == current && checksum_barriers[i].to == target_state) |
547 | 0 | found = true; |
548 | 0 | } |
549 | | |
550 | | /* |
551 | | * If the relevant state criteria aren't satisfied, throw an error which |
552 | | * will be caught by the procsignal machinery for a later retry. |
553 | | */ |
554 | 0 | if (!found) |
555 | 0 | ereport(ERROR, |
556 | 0 | errcode(ERRCODE_INVALID_PARAMETER_VALUE), |
557 | 0 | errmsg("incorrect data checksum state %d for target state %d", |
558 | 0 | current, target_state)); |
559 | | |
560 | 0 | SetLocalDataChecksumState(target_state); |
561 | 0 | return true; |
562 | 0 | } |
563 | | |
564 | | |
565 | | /* |
566 | | * Disables data checksums for the cluster, if applicable. Starts a background |
567 | | * worker which turns off the data checksums. |
568 | | */ |
569 | | Datum |
570 | | disable_data_checksums(PG_FUNCTION_ARGS) |
571 | 0 | { |
572 | 0 | PreventCommandDuringRecovery("pg_disable_data_checksums()"); |
573 | |
|
574 | 0 | if (!superuser()) |
575 | 0 | ereport(ERROR, |
576 | 0 | errcode(ERRCODE_INSUFFICIENT_PRIVILEGE), |
577 | 0 | errmsg("must be superuser to change data checksum state")); |
578 | | |
579 | 0 | StartDataChecksumsWorkerLauncher(DISABLE_DATACHECKSUMS, 0, 0); |
580 | 0 | PG_RETURN_VOID(); |
581 | 0 | } |
582 | | |
583 | | /* |
584 | | * Enables data checksums for the cluster, if applicable. Supports vacuum- |
585 | | * like cost based throttling to limit system load. Starts a background worker |
586 | | * which updates data checksums on existing data. |
587 | | */ |
588 | | Datum |
589 | | enable_data_checksums(PG_FUNCTION_ARGS) |
590 | 0 | { |
591 | 0 | int cost_delay = PG_GETARG_INT32(0); |
592 | 0 | int cost_limit = PG_GETARG_INT32(1); |
593 | |
|
594 | 0 | PreventCommandDuringRecovery("pg_enable_data_checksums()"); |
595 | |
|
596 | 0 | if (!superuser()) |
597 | 0 | ereport(ERROR, |
598 | 0 | errcode(ERRCODE_INSUFFICIENT_PRIVILEGE), |
599 | 0 | errmsg("must be superuser to change data checksum state")); |
600 | | |
601 | 0 | if (cost_delay < 0) |
602 | 0 | ereport(ERROR, |
603 | 0 | errcode(ERRCODE_INVALID_PARAMETER_VALUE), |
604 | 0 | errmsg("cost delay cannot be a negative value")); |
605 | | |
606 | 0 | if (cost_limit <= 0) |
607 | 0 | ereport(ERROR, |
608 | 0 | errcode(ERRCODE_INVALID_PARAMETER_VALUE), |
609 | 0 | errmsg("cost limit must be greater than zero")); |
610 | | |
611 | | /* |
612 | | * An invalid database cannot be connected to, so the worker would fail to |
613 | | * process it, and unlike a dropped database its files stay around. Error |
614 | | * out early with a hint rather than failing halfway through processing. A |
615 | | * database which turns invalid after this check, for example from an |
616 | | * interrupted DROP DATABASE, instead makes its worker fail; the launcher |
617 | | * then aborts and leaves checksums disabled, since the invalid database's |
618 | | * files would otherwise be left without valid checksums. |
619 | | */ |
620 | 0 | ErrorOnInvalidDatabases(); |
621 | |
|
622 | 0 | StartDataChecksumsWorkerLauncher(ENABLE_DATACHECKSUMS, cost_delay, cost_limit); |
623 | |
|
624 | 0 | PG_RETURN_VOID(); |
625 | 0 | } |
626 | | |
627 | | |
628 | | /***************************************************************************** |
629 | | * Functionality for running the datachecksumsworker and associated launcher |
630 | | */ |
631 | | |
632 | | /* |
633 | | * StartDataChecksumsWorkerLauncher |
634 | | * Start the datachecksumsworker launcher process, if not running yet |
635 | | * |
636 | | * This is called to start data checksums processing for enabling as well as |
637 | | * disabling. |
638 | | */ |
639 | | static void |
640 | | StartDataChecksumsWorkerLauncher(DataChecksumsWorkerOperation op, |
641 | | int cost_delay, |
642 | | int cost_limit) |
643 | | { |
644 | | BackgroundWorker bgw; |
645 | | BackgroundWorkerHandle *bgw_handle; |
646 | | bool running; |
647 | | |
648 | | #ifdef USE_ASSERT_CHECKING |
649 | | /* The cost delay settings have no effect when disabling */ |
650 | | if (op == DISABLE_DATACHECKSUMS) |
651 | | Assert(cost_delay == 0 && cost_limit == 0); |
652 | | #endif |
653 | | |
654 | | INJECTION_POINT("datachecksumsworker-startup-delay", NULL); |
655 | | |
656 | | /* Store the desired state in shared memory */ |
657 | | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
658 | | |
659 | | DataChecksumState->launch_operation = op; |
660 | | DataChecksumState->launch_cost_delay = cost_delay; |
661 | | DataChecksumState->launch_cost_limit = cost_limit; |
662 | | |
663 | | /* Is the launcher already running? If so, what is it doing? */ |
664 | | running = DataChecksumState->launcher_running; |
665 | | |
666 | | LWLockRelease(DataChecksumsWorkerLock); |
667 | | |
668 | | /* |
669 | | * Launch a new launcher process, if it's not running already. |
670 | | * |
671 | | * If the launcher is currently busy enabling the checksums, and we want |
672 | | * them disabled (or vice versa), the launcher will notice that at latest |
673 | | * when it's about to exit, and will loop back to process the new request. |
674 | | * So if the launcher is already running, we don't need to do anything |
675 | | * more here to abort it. |
676 | | * |
677 | | * If you call pg_enable/disable_data_checksums() twice in a row, before |
678 | | * the launcher has had a chance to start up, we still end up launching it |
679 | | * twice. That's OK, the second invocation will see that a launcher is |
680 | | * already running and exit quickly. |
681 | | */ |
682 | | if (!running) |
683 | | { |
684 | | if ((op == ENABLE_DATACHECKSUMS && DataChecksumsOn()) || |
685 | | (op == DISABLE_DATACHECKSUMS && DataChecksumsOff())) |
686 | | { |
687 | | ereport(LOG, |
688 | | errmsg("data checksums already in desired state, exiting")); |
689 | | return; |
690 | | } |
691 | | |
692 | | /* |
693 | | * Prepare the BackgroundWorker and launch it. |
694 | | */ |
695 | | memset(&bgw, 0, sizeof(bgw)); |
696 | | bgw.bgw_flags = BGWORKER_SHMEM_ACCESS | BGWORKER_BACKEND_DATABASE_CONNECTION; |
697 | | bgw.bgw_start_time = BgWorkerStart_RecoveryFinished; |
698 | | snprintf(bgw.bgw_library_name, BGW_MAXLEN, "postgres"); |
699 | | snprintf(bgw.bgw_function_name, BGW_MAXLEN, "DataChecksumsWorkerLauncherMain"); |
700 | | snprintf(bgw.bgw_name, BGW_MAXLEN, "datachecksums launcher"); |
701 | | snprintf(bgw.bgw_type, BGW_MAXLEN, "datachecksums launcher"); |
702 | | bgw.bgw_restart_time = BGW_NEVER_RESTART; |
703 | | bgw.bgw_notify_pid = MyProcPid; |
704 | | bgw.bgw_main_arg = (Datum) 0; |
705 | | |
706 | | if (!RegisterDynamicBackgroundWorker(&bgw, &bgw_handle)) |
707 | | ereport(ERROR, |
708 | | errcode(ERRCODE_INSUFFICIENT_RESOURCES), |
709 | | errmsg("failed to start background worker to process data checksums")); |
710 | | } |
711 | | else |
712 | | { |
713 | | ereport(LOG, |
714 | | errmsg("data checksum processing already running")); |
715 | | } |
716 | | } |
717 | | |
718 | | /* |
719 | | * ProcessSingleRelationFork |
720 | | * Enable data checksums in a single relation/fork. |
721 | | * |
722 | | * Returns true if successful, and false if *aborted*. On error, an actual |
723 | | * error is raised in the lower levels. |
724 | | */ |
725 | | static bool |
726 | | ProcessSingleRelationFork(Relation reln, ForkNumber forkNum, BufferAccessStrategy strategy) |
727 | 0 | { |
728 | 0 | BlockNumber numblocks = RelationGetNumberOfBlocksInFork(reln, forkNum); |
729 | 0 | char activity[NAMEDATALEN * 2 + 128]; |
730 | 0 | char *relns; |
731 | |
|
732 | 0 | relns = get_namespace_name(RelationGetNamespace(reln)); |
733 | | |
734 | | /* Report the current relation to pg_stat_activity */ |
735 | 0 | snprintf(activity, sizeof(activity) - 1, "processing: %s.%s (%s, %u blocks)", |
736 | 0 | (relns ? relns : ""), RelationGetRelationName(reln), forkNames[forkNum], numblocks); |
737 | 0 | pgstat_report_activity(STATE_RUNNING, activity); |
738 | 0 | { |
739 | 0 | const int index[] = { |
740 | 0 | PROGRESS_DATACHECKSUMS_BLOCKS_TOTAL, |
741 | 0 | PROGRESS_DATACHECKSUMS_BLOCKS_DONE |
742 | 0 | }; |
743 | |
|
744 | 0 | int64 vals[2]; |
745 | |
|
746 | 0 | vals[0] = numblocks; |
747 | 0 | vals[1] = 0; |
748 | |
|
749 | 0 | pgstat_progress_update_multi_param(2, index, vals); |
750 | 0 | } |
751 | 0 | if (relns) |
752 | 0 | pfree(relns); |
753 | | |
754 | | /* |
755 | | * We are looping over the blocks which existed at the time of process |
756 | | * start, which is safe since new blocks are created with checksums set |
757 | | * already due to the state being "inprogress-on". |
758 | | */ |
759 | 0 | for (BlockNumber blknum = 0; blknum < numblocks; blknum++) |
760 | 0 | { |
761 | 0 | Buffer buf = ReadBufferExtended(reln, forkNum, blknum, RBM_NORMAL, strategy); |
762 | | |
763 | | /* Need to get an exclusive lock to mark the buffer as dirty */ |
764 | 0 | LockBuffer(buf, BUFFER_LOCK_EXCLUSIVE); |
765 | | |
766 | | /* |
767 | | * Mark the buffer as dirty and force a full page write. We have to |
768 | | * re-write the page to WAL even if the checksum hasn't changed, |
769 | | * because if there is a replica it might have a slightly different |
770 | | * version of the page with an invalid checksum, caused by unlogged |
771 | | * changes (e.g. hint bits) on the primary happening while checksums |
772 | | * were off. This can happen if there was a valid checksum on the page |
773 | | * at one point in the past, so only when checksums are first on, then |
774 | | * off, and then turned on again. TODO: investigate if this could be |
775 | | * avoided if the checksum is calculated to be correct and wal_level |
776 | | * is set to "minimal". |
777 | | * |
778 | | * Unlogged relations don't need WAL since they are reset to their |
779 | | * init fork on recovery. We still dirty the buffer so that the |
780 | | * checksum is written to disk at the next checkpoint. |
781 | | * |
782 | | * The init fork is an exception: it is WAL-logged so the standby can |
783 | | * materialize the relation after promotion (see |
784 | | * ResetUnloggedRelations()). Skipping it here would leave the |
785 | | * standby with a stale init fork that, once copied to the main fork |
786 | | * on promotion, would fail checksum verification on every read. |
787 | | */ |
788 | 0 | START_CRIT_SECTION(); |
789 | 0 | MarkBufferDirty(buf); |
790 | 0 | if (RelationNeedsWAL(reln) || forkNum == INIT_FORKNUM) |
791 | 0 | log_newpage_buffer(buf, false); |
792 | 0 | END_CRIT_SECTION(); |
793 | |
|
794 | 0 | UnlockReleaseBuffer(buf); |
795 | | |
796 | | /* Check if we are asked to abort, the abortion will bubble up. */ |
797 | 0 | Assert(operation == ENABLE_DATACHECKSUMS); |
798 | 0 | CHECK_FOR_WORKER_ABORT_REQUEST(); |
799 | 0 | if (abort_requested) |
800 | 0 | return false; |
801 | | |
802 | | /* update the block counter */ |
803 | 0 | pgstat_progress_update_param(PROGRESS_DATACHECKSUMS_BLOCKS_DONE, |
804 | 0 | (blknum + 1)); |
805 | | |
806 | | /* |
807 | | * Processing is re-using the vacuum cost delay for process |
808 | | * throttling, hence why we call vacuum APIs here. |
809 | | */ |
810 | 0 | vacuum_delay_point(false); |
811 | 0 | } |
812 | | |
813 | 0 | return true; |
814 | 0 | } |
815 | | |
816 | | /* |
817 | | * Initialize all data checksum progress counters to be displayed as NULL. |
818 | | */ |
819 | | static void |
820 | | ResetDataChecksumsProgressCounters(void) |
821 | 0 | { |
822 | 0 | const int index[] = { |
823 | 0 | PROGRESS_DATACHECKSUMS_DBS_TOTAL, |
824 | 0 | PROGRESS_DATACHECKSUMS_DBS_DONE, |
825 | 0 | PROGRESS_DATACHECKSUMS_RELS_TOTAL, |
826 | 0 | PROGRESS_DATACHECKSUMS_RELS_DONE, |
827 | 0 | PROGRESS_DATACHECKSUMS_BLOCKS_TOTAL, |
828 | 0 | PROGRESS_DATACHECKSUMS_BLOCKS_DONE, |
829 | 0 | }; |
830 | |
|
831 | 0 | int64 vals[lengthof(index)]; |
832 | |
|
833 | 0 | for (size_t i = 0; i < lengthof(index); i++) |
834 | 0 | vals[i] = -1; |
835 | |
|
836 | 0 | pgstat_progress_update_multi_param(lengthof(index), index, vals); |
837 | 0 | } |
838 | | |
839 | | /* |
840 | | * ProcessSingleRelationByOid |
841 | | * Process a single relation based on oid. |
842 | | * |
843 | | * Returns true if successful, and false if *aborted*. On error, an actual |
844 | | * error is raised in the lower levels. |
845 | | */ |
846 | | static bool |
847 | | ProcessSingleRelationByOid(Oid relationId, BufferAccessStrategy strategy) |
848 | 0 | { |
849 | 0 | Relation rel; |
850 | 0 | bool aborted = false; |
851 | |
|
852 | 0 | StartTransactionCommand(); |
853 | |
|
854 | 0 | rel = try_relation_open(relationId, AccessShareLock); |
855 | 0 | if (rel == NULL) |
856 | 0 | { |
857 | | /* |
858 | | * Relation no longer exists. We don't consider this an error since |
859 | | * there are no pages in it that need data checksums, and thus return |
860 | | * true. The worker operates off a list of relations generated at the |
861 | | * start of processing, so relations being dropped in the meantime is |
862 | | * to be expected. |
863 | | */ |
864 | 0 | CommitTransactionCommand(); |
865 | 0 | pgstat_report_activity(STATE_IDLE, NULL); |
866 | 0 | return true; |
867 | 0 | } |
868 | | |
869 | 0 | for (ForkNumber fnum = 0; fnum <= MAX_FORKNUM; fnum++) |
870 | 0 | { |
871 | 0 | if (smgrexists(RelationGetSmgr(rel), fnum)) |
872 | 0 | { |
873 | 0 | if (!ProcessSingleRelationFork(rel, fnum, strategy)) |
874 | 0 | { |
875 | 0 | aborted = true; |
876 | 0 | break; |
877 | 0 | } |
878 | 0 | } |
879 | 0 | } |
880 | 0 | relation_close(rel, AccessShareLock); |
881 | |
|
882 | 0 | CommitTransactionCommand(); |
883 | 0 | pgstat_report_activity(STATE_IDLE, NULL); |
884 | |
|
885 | 0 | return !aborted; |
886 | 0 | } |
887 | | |
888 | | /* |
889 | | * WaitForDataChecksumsWorkerState |
890 | | * Wait for a data checksums worker to start or stop. |
891 | | * |
892 | | * This is like WaitForBackgroundWorkerStartup() and |
893 | | * WaitForBackgroundWorkerShutdown(), except that it also reacts to SIGINT |
894 | | * received by the launcher. The launcher owns the overall checksum |
895 | | * operation, so canceling it should stop the worker it has registered or is |
896 | | * currently running. |
897 | | * |
898 | | * If wait_for_startup is true, wait until the worker is no longer in |
899 | | * BGWH_NOT_YET_STARTED state, like WaitForBackgroundWorkerStartup(). If it |
900 | | * is false, wait until the worker reaches BGWH_STOPPED state, like |
901 | | * WaitForBackgroundWorkerShutdown(). |
902 | | * |
903 | | * pidp is set to the worker's PID when startup succeeds, if it is not NULL. |
904 | | */ |
905 | | static BgwHandleStatus |
906 | | WaitForDataChecksumsWorkerState(BackgroundWorkerHandle *handle, |
907 | | bool wait_for_startup, |
908 | | pid_t *pidp, |
909 | | uint32 wait_event) |
910 | 0 | { |
911 | 0 | BgwHandleStatus status; |
912 | 0 | bool termination_requested = false; |
913 | |
|
914 | 0 | for (;;) |
915 | 0 | { |
916 | 0 | int rc; |
917 | 0 | pid_t pid; |
918 | |
|
919 | 0 | CHECK_FOR_INTERRUPTS(); |
920 | |
|
921 | 0 | status = GetBackgroundWorkerPid(handle, &pid); |
922 | 0 | if (status == BGWH_STARTED && pidp) |
923 | 0 | *pidp = pid; |
924 | |
|
925 | 0 | if (abort_requested && !termination_requested) |
926 | 0 | { |
927 | 0 | TerminateBackgroundWorker(handle); |
928 | 0 | termination_requested = true; |
929 | 0 | } |
930 | | |
931 | | /* |
932 | | * Startup waits for the worker to leave BGWH_NOT_YET_STARTED, while |
933 | | * shutdown waits for it to reach BGWH_STOPPED. |
934 | | */ |
935 | 0 | if (status == BGWH_STOPPED || |
936 | 0 | (wait_for_startup && status == BGWH_STARTED)) |
937 | 0 | break; |
938 | | |
939 | 0 | rc = WaitLatch(MyLatch, |
940 | 0 | WL_LATCH_SET | WL_POSTMASTER_DEATH, 0, |
941 | 0 | wait_event); |
942 | |
|
943 | 0 | if (rc & WL_POSTMASTER_DEATH) |
944 | 0 | { |
945 | 0 | status = BGWH_POSTMASTER_DIED; |
946 | 0 | break; |
947 | 0 | } |
948 | | |
949 | 0 | ResetLatch(MyLatch); |
950 | 0 | } |
951 | |
|
952 | 0 | return status; |
953 | 0 | } |
954 | | |
955 | | /* |
956 | | * ProcessDatabase |
957 | | * Enable data checksums in a single database. |
958 | | * |
959 | | * We do this by launching a dynamic background worker into this database, and |
960 | | * waiting for it to finish. We have to do this in a separate worker, since |
961 | | * each process can only be connected to one database during its lifetime. |
962 | | */ |
963 | | static DataChecksumsWorkerResult |
964 | | ProcessDatabase(DataChecksumsWorkerDatabase *db) |
965 | | { |
966 | | BackgroundWorker bgw; |
967 | | BackgroundWorkerHandle *bgw_handle; |
968 | | BgwHandleStatus status; |
969 | | pid_t pid; |
970 | | uint64 invocation; |
971 | | char activity[NAMEDATALEN + 64]; |
972 | | DataChecksumsWorkerResult result; |
973 | | |
974 | | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
975 | | |
976 | | /* |
977 | | * Initialize result to FAILED. The worker will change it to SUCCESSFUL |
978 | | * if it completes successfully. |
979 | | */ |
980 | | DataChecksumState->worker_result = DATACHECKSUMSWORKER_FAILED; |
981 | | DataChecksumState->worker_pid = InvalidPid; |
982 | | |
983 | | invocation = ++DataChecksumState->worker_invocation_counter; |
984 | | DataChecksumState->worker_invocation = invocation; |
985 | | DataChecksumState->database_oid = db->dboid; |
986 | | |
987 | | LWLockRelease(DataChecksumsWorkerLock); |
988 | | |
989 | | memset(&bgw, 0, sizeof(bgw)); |
990 | | bgw.bgw_flags = BGWORKER_SHMEM_ACCESS | BGWORKER_BACKEND_DATABASE_CONNECTION; |
991 | | bgw.bgw_start_time = BgWorkerStart_RecoveryFinished; |
992 | | snprintf(bgw.bgw_library_name, BGW_MAXLEN, "postgres"); |
993 | | snprintf(bgw.bgw_function_name, BGW_MAXLEN, "%s", "DataChecksumsWorkerMain"); |
994 | | snprintf(bgw.bgw_name, BGW_MAXLEN, "datachecksums worker"); |
995 | | snprintf(bgw.bgw_type, BGW_MAXLEN, "datachecksums worker"); |
996 | | bgw.bgw_restart_time = BGW_NEVER_RESTART; |
997 | | bgw.bgw_notify_pid = MyProcPid; |
998 | | /* pass the invocation number to the worker process */ |
999 | | bgw.bgw_main_arg = UInt64GetDatum(invocation); |
1000 | | |
1001 | | /* |
1002 | | * If there are no worker slots available, there is little we can do. If |
1003 | | * we retry in a bit it's still unlikely that the user has managed to |
1004 | | * reconfigure in the meantime and we'd be run through retries fast. |
1005 | | */ |
1006 | | if (!RegisterDynamicBackgroundWorker(&bgw, &bgw_handle)) |
1007 | | { |
1008 | | ereport(WARNING, |
1009 | | errmsg("could not start background worker for enabling data checksums in database \"%s\"", |
1010 | | db->dbname), |
1011 | | errhint("The \"%s\" setting might be too low.", "max_worker_processes")); |
1012 | | return DATACHECKSUMSWORKER_FAILED; |
1013 | | } |
1014 | | |
1015 | | /* |
1016 | | * While this expects to wait for BGWORKER_STARTUP it may return _STOPPED |
1017 | | * if the worker was terminated in the meantime so we must check status. |
1018 | | */ |
1019 | | status = WaitForDataChecksumsWorkerState(bgw_handle, true, &pid, |
1020 | | WAIT_EVENT_BGWORKER_STARTUP); |
1021 | | if (status == BGWH_STOPPED) |
1022 | | { |
1023 | | if (abort_requested) |
1024 | | { |
1025 | | result = DATACHECKSUMSWORKER_ABORTED; |
1026 | | goto done; |
1027 | | } |
1028 | | |
1029 | | /* |
1030 | | * If the worker managed to start, and stop, before we got to waiting |
1031 | | * for it we can see a STOPPED status here without it being a failure. |
1032 | | */ |
1033 | | LWLockAcquire(DataChecksumsWorkerLock, LW_SHARED); |
1034 | | Assert(DataChecksumState->worker_invocation == invocation); |
1035 | | if (DataChecksumState->worker_result == DATACHECKSUMSWORKER_SUCCESSFUL) |
1036 | | { |
1037 | | LWLockRelease(DataChecksumsWorkerLock); |
1038 | | pgstat_report_activity(STATE_IDLE, NULL); |
1039 | | return DATACHECKSUMSWORKER_SUCCESSFUL; |
1040 | | } |
1041 | | LWLockRelease(DataChecksumsWorkerLock); |
1042 | | |
1043 | | ereport(WARNING, |
1044 | | errmsg("could not start background worker for enabling data checksums in database \"%s\"", |
1045 | | db->dbname), |
1046 | | errhint("More details on the error might be found in the server log.")); |
1047 | | |
1048 | | /* |
1049 | | * Heuristic to see if the database was dropped, and if it was we can |
1050 | | * treat it as not an error, else treat as fatal and error out. |
1051 | | */ |
1052 | | if (DatabaseExists(db->dboid)) |
1053 | | return DATACHECKSUMSWORKER_FAILED; |
1054 | | else |
1055 | | return DATACHECKSUMSWORKER_DROPDB; |
1056 | | } |
1057 | | |
1058 | | /* |
1059 | | * If the postmaster crashed we cannot end up with a processed database so |
1060 | | * we have no alternative other than exiting. When enabling checksums we |
1061 | | * won't at this time have changed the data checksums state in pg_control |
1062 | | * to enabled so when the cluster comes back up processing will have to be |
1063 | | * restarted. |
1064 | | */ |
1065 | | if (status == BGWH_POSTMASTER_DIED) |
1066 | | ereport(FATAL, |
1067 | | errcode(ERRCODE_ADMIN_SHUTDOWN), |
1068 | | errmsg("cannot enable data checksums without the postmaster process"), |
1069 | | errhint("Restart the database and restart data checksum processing by calling pg_enable_data_checksums().")); |
1070 | | |
1071 | | Assert(status == BGWH_STARTED); |
1072 | | ereport(LOG, |
1073 | | errmsg("initiating data checksum processing in database \"%s\"", |
1074 | | db->dbname)); |
1075 | | |
1076 | | /* Save the pid of the worker so we can signal it later */ |
1077 | | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1078 | | Assert(DataChecksumState->worker_invocation == invocation); |
1079 | | DataChecksumState->worker_pid = pid; |
1080 | | LWLockRelease(DataChecksumsWorkerLock); |
1081 | | |
1082 | | snprintf(activity, sizeof(activity) - 1, |
1083 | | "Waiting for worker in database %s (pid %ld)", db->dbname, (long) pid); |
1084 | | pgstat_report_activity(STATE_RUNNING, activity); |
1085 | | |
1086 | | status = WaitForDataChecksumsWorkerState(bgw_handle, false, NULL, |
1087 | | WAIT_EVENT_BGWORKER_SHUTDOWN); |
1088 | | if (status == BGWH_POSTMASTER_DIED) |
1089 | | ereport(FATAL, |
1090 | | errcode(ERRCODE_ADMIN_SHUTDOWN), |
1091 | | errmsg("postmaster exited during data checksum processing in \"%s\"", |
1092 | | db->dbname), |
1093 | | errhint("Restart the database and restart data checksum processing by calling pg_enable_data_checksums().")); |
1094 | | |
1095 | | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1096 | | Assert(DataChecksumState->worker_invocation == invocation); |
1097 | | result = DataChecksumState->worker_result; |
1098 | | DataChecksumState->worker_pid = InvalidPid; |
1099 | | LWLockRelease(DataChecksumsWorkerLock); |
1100 | | |
1101 | | /* |
1102 | | * A worker which started but failed before reporting a result has most |
1103 | | * likely FATALed in InitPostgres. If the database was dropped after we |
1104 | | * built the database list then that is the expected outcome and not an |
1105 | | * error, so apply the same heuristic as when the worker failed to start. |
1106 | | */ |
1107 | | if (result == DATACHECKSUMSWORKER_FAILED && !DatabaseExists(db->dboid)) |
1108 | | result = DATACHECKSUMSWORKER_DROPDB; |
1109 | | |
1110 | | CHECK_FOR_LAUNCHER_ABORT_REQUEST(); |
1111 | | if (abort_requested) |
1112 | | result = DATACHECKSUMSWORKER_ABORTED; |
1113 | | |
1114 | | done: |
1115 | | if (result == DATACHECKSUMSWORKER_ABORTED) |
1116 | | ereport(LOG, |
1117 | | errmsg("data checksums processing was aborted in database \"%s\"", |
1118 | | db->dbname)); |
1119 | | pgstat_report_activity(STATE_IDLE, NULL); |
1120 | | return result; |
1121 | | } |
1122 | | |
1123 | | /* |
1124 | | * launcher_exit |
1125 | | * |
1126 | | * Internal routine for cleaning up state when a launcher process which has |
1127 | | * performed checksum operations exits. A launcher process which is exiting due |
1128 | | * to a duplicate started launcher does not need to perform any cleanup and |
1129 | | * this function should not be called. Otherwise, we need to clean up the abort |
1130 | | * flag to ensure that processing can be started again if it was previously |
1131 | | * aborted (note: started again, *not* restarted from where it left off). |
1132 | | */ |
1133 | | static void |
1134 | | launcher_exit(int code, Datum arg) |
1135 | 0 | { |
1136 | 0 | abort_requested = false; |
1137 | |
|
1138 | 0 | if (launcher_running) |
1139 | 0 | { |
1140 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1141 | 0 | if (DataChecksumState->worker_pid != InvalidPid) |
1142 | 0 | { |
1143 | 0 | ereport(LOG, |
1144 | 0 | errmsg("data checksums launcher exiting while worker is still running, signalling worker")); |
1145 | 0 | kill(DataChecksumState->worker_pid, SIGTERM); |
1146 | 0 | DataChecksumState->worker_pid = InvalidPid; |
1147 | 0 | } |
1148 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1149 | 0 | } |
1150 | | |
1151 | | /* |
1152 | | * If the launcher is exiting before data checksums are enabled then set |
1153 | | * the state to off since processing cannot be resumed. |
1154 | | */ |
1155 | 0 | if (DataChecksumsInProgressOn()) |
1156 | 0 | SetDataChecksumsOff(); |
1157 | |
|
1158 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1159 | 0 | if (launcher_running) |
1160 | 0 | DataChecksumState->launcher_running = false; |
1161 | 0 | launcher_running = false; |
1162 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1163 | 0 | } |
1164 | | |
1165 | | /* |
1166 | | * launcher_cancel_handler |
1167 | | * |
1168 | | * Internal routine for reacting to SIGINT and flagging the launcher to abort. |
1169 | | * If a worker is registered or running, the launcher will request worker |
1170 | | * termination from its normal control flow. |
1171 | | */ |
1172 | | static void |
1173 | | launcher_cancel_handler(SIGNAL_ARGS) |
1174 | 0 | { |
1175 | 0 | int save_errno = errno; |
1176 | |
|
1177 | 0 | abort_requested = true; |
1178 | | |
1179 | | /* |
1180 | | * Wake the launcher if it is waiting for transactions to finish or for a |
1181 | | * worker to start up or shut down. |
1182 | | */ |
1183 | 0 | SetLatch(MyLatch); |
1184 | |
|
1185 | 0 | errno = save_errno; |
1186 | 0 | } |
1187 | | |
1188 | | /* |
1189 | | * WaitForAllTransactionsToFinish |
1190 | | * Blocks awaiting all current transactions to finish |
1191 | | * |
1192 | | * Returns when all transactions which are active at the call of the function |
1193 | | * have ended. |
1194 | | * |
1195 | | * NB: this will return early, if aborted by SIGINT or if the target state |
1196 | | * is changed while we're running. |
1197 | | */ |
1198 | | static void |
1199 | | WaitForAllTransactionsToFinish(void) |
1200 | 0 | { |
1201 | 0 | TransactionId waitforxid; |
1202 | |
|
1203 | 0 | LWLockAcquire(XidGenLock, LW_SHARED); |
1204 | 0 | waitforxid = XidFromFullTransactionId(TransamVariables->nextXid); |
1205 | 0 | LWLockRelease(XidGenLock); |
1206 | |
|
1207 | 0 | while (TransactionIdPrecedes(GetOldestActiveTransactionId(false, true), waitforxid)) |
1208 | 0 | { |
1209 | 0 | char activity[64]; |
1210 | 0 | int rc; |
1211 | | |
1212 | | /* Oldest running xid is older than us, so wait */ |
1213 | 0 | snprintf(activity, |
1214 | 0 | sizeof(activity), |
1215 | 0 | "Waiting for transactions older than %u to end", |
1216 | 0 | waitforxid); |
1217 | 0 | pgstat_report_activity(STATE_RUNNING, activity); |
1218 | | |
1219 | | /* Retry every 3 seconds */ |
1220 | 0 | ResetLatch(MyLatch); |
1221 | 0 | rc = WaitLatch(MyLatch, |
1222 | 0 | WL_LATCH_SET | WL_TIMEOUT | WL_POSTMASTER_DEATH, |
1223 | 0 | 3000, |
1224 | 0 | WAIT_EVENT_CHECKSUM_ENABLE_STARTCONDITION); |
1225 | | |
1226 | | /* |
1227 | | * If the postmaster died, bail out. But first print a log message to |
1228 | | * note that the checksumming didn't complete. |
1229 | | */ |
1230 | 0 | if (rc & WL_POSTMASTER_DEATH) |
1231 | 0 | ereport(FATAL, |
1232 | 0 | errcode(ERRCODE_ADMIN_SHUTDOWN), |
1233 | 0 | errmsg("postmaster exited during data checksums processing"), |
1234 | 0 | errhint("Data checksums processing must be restarted manually after cluster restart.")); |
1235 | | |
1236 | 0 | CHECK_FOR_INTERRUPTS(); |
1237 | 0 | CHECK_FOR_LAUNCHER_ABORT_REQUEST(); |
1238 | |
|
1239 | 0 | if (abort_requested) |
1240 | 0 | break; |
1241 | 0 | } |
1242 | | |
1243 | 0 | pgstat_report_activity(STATE_IDLE, NULL); |
1244 | 0 | return; |
1245 | 0 | } |
1246 | | |
1247 | | /* |
1248 | | * DataChecksumsWorkerLauncherMain |
1249 | | * |
1250 | | * Main function for launching dynamic background workers for processing data |
1251 | | * checksums in databases. This function has the bgworker management, with |
1252 | | * ProcessAllDatabases being responsible for looping over the databases and |
1253 | | * initiating processing. |
1254 | | */ |
1255 | | void |
1256 | | DataChecksumsWorkerLauncherMain(Datum arg) |
1257 | 0 | { |
1258 | |
|
1259 | 0 | ereport(DEBUG1, |
1260 | 0 | errmsg("background worker \"datachecksums launcher\" started")); |
1261 | | |
1262 | 0 | pqsignal(SIGTERM, die); |
1263 | 0 | pqsignal(SIGINT, launcher_cancel_handler); |
1264 | 0 | pqsignal(SIGUSR1, procsignal_sigusr1_handler); |
1265 | 0 | pqsignal(SIGUSR2, PG_SIG_IGN); |
1266 | |
|
1267 | 0 | BackgroundWorkerUnblockSignals(); |
1268 | |
|
1269 | 0 | MyBackendType = B_DATACHECKSUMSWORKER_LAUNCHER; |
1270 | 0 | init_ps_display(NULL); |
1271 | |
|
1272 | 0 | INJECTION_POINT("datachecksumsworker-launcher-delay", NULL); |
1273 | |
|
1274 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1275 | |
|
1276 | 0 | if (DataChecksumState->launcher_running) |
1277 | 0 | { |
1278 | 0 | ereport(LOG, |
1279 | 0 | errmsg("background worker \"datachecksums launcher\" already running, exiting")); |
1280 | | /* Launcher was already running, let it finish */ |
1281 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1282 | 0 | return; |
1283 | 0 | } |
1284 | | |
1285 | 0 | on_shmem_exit(launcher_exit, 0); |
1286 | 0 | launcher_running = true; |
1287 | | |
1288 | | /* Initialize a connection to shared catalogs only */ |
1289 | 0 | BackgroundWorkerInitializeConnectionByOid(InvalidOid, InvalidOid, 0); |
1290 | |
|
1291 | 0 | operation = DataChecksumState->launch_operation; |
1292 | 0 | DataChecksumState->launcher_running = true; |
1293 | 0 | DataChecksumState->operation = operation; |
1294 | 0 | DataChecksumState->cost_delay = DataChecksumState->launch_cost_delay; |
1295 | 0 | DataChecksumState->cost_limit = DataChecksumState->launch_cost_limit; |
1296 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1297 | | |
1298 | | /* |
1299 | | * The target state can change while we are busy enabling/disabling |
1300 | | * checksums, if the user calls pg_disable/enable_data_checksums() before |
1301 | | * we are finished with the previous request. In that case, we will loop |
1302 | | * back here, to process the new request. |
1303 | | */ |
1304 | 0 | again: |
1305 | |
|
1306 | 0 | pgstat_progress_start_command(PROGRESS_COMMAND_DATACHECKSUMS, |
1307 | 0 | InvalidOid); |
1308 | 0 | ResetDataChecksumsProgressCounters(); |
1309 | |
|
1310 | 0 | if (operation == ENABLE_DATACHECKSUMS) |
1311 | 0 | { |
1312 | | /* |
1313 | | * If we are asked to enable checksums in a cluster which already has |
1314 | | * checksums enabled, exit immediately as there is nothing more to do. |
1315 | | */ |
1316 | 0 | if (DataChecksumsNeedVerify()) |
1317 | 0 | goto done; |
1318 | | |
1319 | 0 | ereport(LOG, |
1320 | 0 | errmsg("enabling data checksums requested, starting data checksum calculation")); |
1321 | | |
1322 | | /* |
1323 | | * Set the state to inprogress-on and wait on the procsignal barrier. |
1324 | | */ |
1325 | 0 | pgstat_progress_update_param(PROGRESS_DATACHECKSUMS_PHASE, |
1326 | 0 | PROGRESS_DATACHECKSUMS_PHASE_ENABLING); |
1327 | 0 | SetDataChecksumsOnInProgress(); |
1328 | | |
1329 | | /* |
1330 | | * All backends are now in inprogress-on state and are writing data |
1331 | | * checksums. Start processing all data at rest. |
1332 | | */ |
1333 | 0 | if (!ProcessAllDatabases()) |
1334 | 0 | { |
1335 | | /* |
1336 | | * If processing was canceled, or the target state changed during |
1337 | | * processing, then it's not a failure. In the latter case, the |
1338 | | * launcher will restart processing with the new target state. |
1339 | | */ |
1340 | 0 | CHECK_FOR_LAUNCHER_ABORT_REQUEST(); |
1341 | 0 | if (abort_requested) |
1342 | 0 | goto done; |
1343 | 0 | ereport(ERROR, |
1344 | 0 | errcode(ERRCODE_INSUFFICIENT_RESOURCES), |
1345 | 0 | errmsg("unable to enable data checksums in cluster")); |
1346 | 0 | } |
1347 | | |
1348 | | /* |
1349 | | * Data checksums have been set on all pages, set the state to on in |
1350 | | * order to instruct backends to validate checksums on reading. |
1351 | | */ |
1352 | 0 | SetDataChecksumsOn(); |
1353 | |
|
1354 | 0 | ereport(LOG, |
1355 | 0 | errmsg("data checksums are now enabled")); |
1356 | 0 | } |
1357 | 0 | else if (operation == DISABLE_DATACHECKSUMS) |
1358 | 0 | { |
1359 | 0 | ereport(LOG, |
1360 | 0 | errmsg("disabling data checksums requested")); |
1361 | | |
1362 | 0 | pgstat_progress_update_param(PROGRESS_DATACHECKSUMS_PHASE, |
1363 | 0 | PROGRESS_DATACHECKSUMS_PHASE_DISABLING); |
1364 | 0 | SetDataChecksumsOff(); |
1365 | 0 | ereport(LOG, |
1366 | 0 | errmsg("data checksums are now disabled")); |
1367 | 0 | } |
1368 | 0 | else |
1369 | 0 | Assert(false); |
1370 | | |
1371 | 0 | done: |
1372 | | |
1373 | | /* |
1374 | | * This state will only be displayed for a fleeting moment, but for the |
1375 | | * sake of correctness it is still added before ending the command. |
1376 | | */ |
1377 | 0 | pgstat_progress_update_param(PROGRESS_DATACHECKSUMS_PHASE, |
1378 | 0 | PROGRESS_DATACHECKSUMS_PHASE_DONE); |
1379 | | |
1380 | | /* |
1381 | | * All done. But before we exit, check if the target state was changed |
1382 | | * while we were running. In that case we will have to start all over |
1383 | | * again. |
1384 | | */ |
1385 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1386 | 0 | if (DataChecksumState->launch_operation != operation) |
1387 | 0 | { |
1388 | 0 | DataChecksumState->operation = DataChecksumState->launch_operation; |
1389 | 0 | operation = DataChecksumState->launch_operation; |
1390 | 0 | DataChecksumState->cost_delay = DataChecksumState->launch_cost_delay; |
1391 | 0 | DataChecksumState->cost_limit = DataChecksumState->launch_cost_limit; |
1392 | | |
1393 | | /* |
1394 | | * If the user started, but aborted processing, and then changed their |
1395 | | * mind again before we had time to exit we need to clear the abort |
1396 | | * flag. |
1397 | | */ |
1398 | 0 | abort_requested = false; |
1399 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1400 | 0 | goto again; |
1401 | 0 | } |
1402 | | |
1403 | | /* Shut down progress reporting as we are done */ |
1404 | 0 | pgstat_progress_end_command(); |
1405 | |
|
1406 | 0 | launcher_running = false; |
1407 | 0 | DataChecksumState->launcher_running = false; |
1408 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1409 | 0 | } |
1410 | | |
1411 | | /* |
1412 | | * ProcessAllDatabases |
1413 | | * Compute the list of all databases and process checksums in each |
1414 | | * |
1415 | | * This will generate a list of databases to process for enabling checksums. |
1416 | | * If a database encounters a failure then processing will end immediately and |
1417 | | * return an error. |
1418 | | */ |
1419 | | static bool |
1420 | | ProcessAllDatabases(void) |
1421 | 0 | { |
1422 | 0 | List *DatabaseList; |
1423 | 0 | int cumulative_total = 0; |
1424 | | |
1425 | | /* Set up so first run processes shared catalogs, not once in every db */ |
1426 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1427 | 0 | DataChecksumState->process_shared_catalogs = true; |
1428 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1429 | | |
1430 | | /* Get a list of all databases to process */ |
1431 | 0 | WaitForAllTransactionsToFinish(); |
1432 | 0 | if (abort_requested) |
1433 | 0 | return false; |
1434 | 0 | DatabaseList = BuildDatabaseList(); |
1435 | | |
1436 | | /* |
1437 | | * Update progress reporting with the total number of databases we need to |
1438 | | * process. This number should not be changed during processing, the |
1439 | | * columns for processed databases is instead increased such that it can |
1440 | | * be compared against the total. |
1441 | | */ |
1442 | 0 | { |
1443 | 0 | const int index[] = { |
1444 | 0 | PROGRESS_DATACHECKSUMS_DBS_TOTAL, |
1445 | 0 | PROGRESS_DATACHECKSUMS_DBS_DONE, |
1446 | 0 | }; |
1447 | |
|
1448 | 0 | int64 vals[2]; |
1449 | |
|
1450 | 0 | vals[0] = list_length(DatabaseList); |
1451 | 0 | vals[1] = 0; |
1452 | |
|
1453 | 0 | pgstat_progress_update_multi_param(2, index, vals); |
1454 | 0 | } |
1455 | |
|
1456 | 0 | foreach_ptr(DataChecksumsWorkerDatabase, db, DatabaseList) |
1457 | 0 | { |
1458 | 0 | DataChecksumsWorkerResult result; |
1459 | |
|
1460 | 0 | result = ProcessDatabase(db); |
1461 | |
|
1462 | | #ifdef USE_INJECTION_POINTS |
1463 | | /* Allow a test process to alter the result of the operation */ |
1464 | | if (IS_INJECTION_POINT_ATTACHED("datachecksumsworker-fail-db-result")) |
1465 | | { |
1466 | | result = DATACHECKSUMSWORKER_FAILED; |
1467 | | INJECTION_POINT_CACHED("datachecksumsworker-fail-db-result", |
1468 | | db->dbname); |
1469 | | } |
1470 | | #endif |
1471 | |
|
1472 | 0 | pgstat_progress_update_param(PROGRESS_DATACHECKSUMS_DBS_DONE, |
1473 | 0 | ++cumulative_total); |
1474 | |
|
1475 | 0 | if (result == DATACHECKSUMSWORKER_FAILED) |
1476 | 0 | { |
1477 | | /* |
1478 | | * Disable checksums on cluster, because we failed one of the |
1479 | | * databases and this is an all or nothing process. |
1480 | | */ |
1481 | 0 | SetDataChecksumsOff(); |
1482 | 0 | ereport(ERROR, |
1483 | 0 | errcode(ERRCODE_INSUFFICIENT_RESOURCES), |
1484 | 0 | errmsg("data checksums failed to get enabled in all databases, aborting"), |
1485 | 0 | errhint("The server log might have more information on the cause of the error.")); |
1486 | 0 | } |
1487 | 0 | else if (result == DATACHECKSUMSWORKER_ABORTED || abort_requested) |
1488 | 0 | { |
1489 | | /* Abort flag set, so exit the whole process */ |
1490 | 0 | FreeDatabaseList(DatabaseList); |
1491 | 0 | return false; |
1492 | 0 | } |
1493 | 0 | else if (result == DATACHECKSUMSWORKER_DROPDB) |
1494 | 0 | { |
1495 | | /* |
1496 | | * Ignore databases that were dropped before their worker could |
1497 | | * process them, and continue with the remaining databases. |
1498 | | */ |
1499 | 0 | continue; |
1500 | 0 | } |
1501 | | |
1502 | | /* |
1503 | | * When one database has completed, it will have done shared catalogs |
1504 | | * so we don't have to process them again. |
1505 | | */ |
1506 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1507 | 0 | DataChecksumState->process_shared_catalogs = false; |
1508 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1509 | 0 | } |
1510 | | |
1511 | 0 | FreeDatabaseList(DatabaseList); |
1512 | |
|
1513 | 0 | pgstat_progress_update_param(PROGRESS_DATACHECKSUMS_PHASE, |
1514 | 0 | PROGRESS_DATACHECKSUMS_PHASE_WAITING_BARRIER); |
1515 | 0 | return true; |
1516 | 0 | } |
1517 | | |
1518 | | /* |
1519 | | * DataChecksumsShmemRequest |
1520 | | * Request datachecksumsworker-related shared memory |
1521 | | */ |
1522 | | static void |
1523 | | DataChecksumsShmemRequest(void *arg) |
1524 | 0 | { |
1525 | 0 | ShmemRequestStruct(.name = "DataChecksumsWorker Data", |
1526 | 0 | .size = sizeof(DataChecksumsStateStruct), |
1527 | 0 | .ptr = (void **) &DataChecksumState, |
1528 | 0 | ); |
1529 | 0 | } |
1530 | | |
1531 | | /* |
1532 | | * DatabaseExists |
1533 | | * |
1534 | | * Scans the system catalog to check if a database with the given Oid exists |
1535 | | * and returns true if it is found, even if it is marked invalid. An invalid |
1536 | | * database still has files that need checksums, so only a missing catalog row |
1537 | | * proves that a concurrent DROP DATABASE completed. |
1538 | | */ |
1539 | | static bool |
1540 | | DatabaseExists(Oid dboid) |
1541 | 0 | { |
1542 | 0 | Relation rel; |
1543 | 0 | ScanKeyData skey; |
1544 | 0 | SysScanDesc scan; |
1545 | 0 | bool found; |
1546 | 0 | HeapTuple tuple; |
1547 | |
|
1548 | 0 | StartTransactionCommand(); |
1549 | | |
1550 | | /* |
1551 | | * DROP DATABASE holds an exclusive lock on the database from before it |
1552 | | * terminates the connections to it until it commits, so take a lock which |
1553 | | * conflicts with it to wait out a drop which is in flight. Without this |
1554 | | * we can see a database whose worker was just killed by DROP DATABASE ... |
1555 | | * WITH (FORCE) as still existing, and report a spurious failure. |
1556 | | */ |
1557 | 0 | LockSharedObject(DatabaseRelationId, dboid, 0, AccessShareLock); |
1558 | |
|
1559 | 0 | rel = table_open(DatabaseRelationId, AccessShareLock); |
1560 | 0 | ScanKeyInit(&skey, |
1561 | 0 | Anum_pg_database_oid, |
1562 | 0 | BTEqualStrategyNumber, F_OIDEQ, |
1563 | 0 | ObjectIdGetDatum(dboid)); |
1564 | 0 | scan = systable_beginscan(rel, DatabaseOidIndexId, true, SnapshotSelf, |
1565 | 0 | 1, &skey); |
1566 | 0 | tuple = systable_getnext(scan); |
1567 | 0 | found = HeapTupleIsValid(tuple); |
1568 | |
|
1569 | 0 | systable_endscan(scan); |
1570 | 0 | table_close(rel, AccessShareLock); |
1571 | |
|
1572 | 0 | CommitTransactionCommand(); |
1573 | |
|
1574 | 0 | return found; |
1575 | 0 | } |
1576 | | |
1577 | | /* |
1578 | | * ErrorOnInvalidDatabases |
1579 | | * Error out if the cluster contains an invalid database |
1580 | | * |
1581 | | * A database left invalid by an interrupted DROP DATABASE cannot be connected |
1582 | | * to, so data checksums can never be enabled in it, while its files remain on |
1583 | | * disk where checksum verification will find them. Report it to the caller |
1584 | | * so the user can drop it before retrying. Called from a normal backend, so |
1585 | | * unlike DatabaseExists we are already in a transaction. |
1586 | | * |
1587 | | * A cluster can contain more than one invalid database, but only the first one |
1588 | | * found is reported; collecting them all is not worth the complexity here. A |
1589 | | * user with several of them gets the error again for the next one after |
1590 | | * dropping the reported database, which the hint accounts for. |
1591 | | */ |
1592 | | static void |
1593 | | ErrorOnInvalidDatabases(void) |
1594 | 0 | { |
1595 | 0 | Relation rel; |
1596 | 0 | TableScanDesc scan; |
1597 | 0 | HeapTuple tup; |
1598 | |
|
1599 | 0 | rel = table_open(DatabaseRelationId, AccessShareLock); |
1600 | 0 | scan = table_beginscan_catalog(rel, 0, NULL); |
1601 | |
|
1602 | 0 | while (HeapTupleIsValid(tup = heap_getnext(scan, ForwardScanDirection))) |
1603 | 0 | { |
1604 | 0 | Form_pg_database pgdb = (Form_pg_database) GETSTRUCT(tup); |
1605 | |
|
1606 | 0 | if (database_is_invalid_form(pgdb)) |
1607 | 0 | ereport(ERROR, |
1608 | 0 | errcode(ERRCODE_OBJECT_NOT_IN_PREREQUISITE_STATE), |
1609 | 0 | errmsg("cannot enable data checksums in a cluster with invalid database \"%s\"", |
1610 | 0 | NameStr(pgdb->datname)), |
1611 | 0 | errhint("Use DROP DATABASE to drop invalid databases.")); |
1612 | 0 | } |
1613 | | |
1614 | 0 | table_endscan(scan); |
1615 | 0 | table_close(rel, AccessShareLock); |
1616 | 0 | } |
1617 | | |
1618 | | /* |
1619 | | * BuildDatabaseList |
1620 | | * Compile a list of all currently available databases in the cluster |
1621 | | * |
1622 | | * This creates the list of databases for the datachecksumsworker workers to |
1623 | | * add checksums to. If the caller wants to ensure that no concurrently |
1624 | | * running CREATE DATABASE calls exist, this needs to be preceded by a call |
1625 | | * to WaitForAllTransactionsToFinish(). |
1626 | | */ |
1627 | | static List * |
1628 | | BuildDatabaseList(void) |
1629 | 0 | { |
1630 | 0 | List *DatabaseList = NIL; |
1631 | 0 | Relation rel; |
1632 | 0 | TableScanDesc scan; |
1633 | 0 | HeapTuple tup; |
1634 | 0 | MemoryContext ctx = CurrentMemoryContext; |
1635 | 0 | MemoryContext oldctx; |
1636 | |
|
1637 | 0 | StartTransactionCommand(); |
1638 | |
|
1639 | 0 | rel = table_open(DatabaseRelationId, AccessShareLock); |
1640 | 0 | scan = table_beginscan_catalog(rel, 0, NULL); |
1641 | |
|
1642 | 0 | while (HeapTupleIsValid(tup = heap_getnext(scan, ForwardScanDirection))) |
1643 | 0 | { |
1644 | 0 | Form_pg_database pgdb = (Form_pg_database) GETSTRUCT(tup); |
1645 | 0 | DataChecksumsWorkerDatabase *db; |
1646 | |
|
1647 | 0 | oldctx = MemoryContextSwitchTo(ctx); |
1648 | |
|
1649 | 0 | db = palloc0_object(DataChecksumsWorkerDatabase); |
1650 | |
|
1651 | 0 | db->dboid = pgdb->oid; |
1652 | 0 | db->dbname = pstrdup(NameStr(pgdb->datname)); |
1653 | |
|
1654 | 0 | DatabaseList = lappend(DatabaseList, db); |
1655 | |
|
1656 | 0 | MemoryContextSwitchTo(oldctx); |
1657 | 0 | } |
1658 | |
|
1659 | 0 | table_endscan(scan); |
1660 | 0 | table_close(rel, AccessShareLock); |
1661 | |
|
1662 | 0 | CommitTransactionCommand(); |
1663 | |
|
1664 | 0 | return DatabaseList; |
1665 | 0 | } |
1666 | | |
1667 | | static void |
1668 | | FreeDatabaseList(List *dblist) |
1669 | 0 | { |
1670 | 0 | if (!dblist) |
1671 | 0 | return; |
1672 | | |
1673 | 0 | foreach_ptr(DataChecksumsWorkerDatabase, db, dblist) |
1674 | 0 | { |
1675 | 0 | if (db->dbname != NULL) |
1676 | 0 | pfree(db->dbname); |
1677 | 0 | } |
1678 | |
|
1679 | 0 | list_free_deep(dblist); |
1680 | 0 | } |
1681 | | |
1682 | | /* |
1683 | | * BuildRelationList |
1684 | | * Compile a list of relations in the database |
1685 | | * |
1686 | | * Returns a list of OIDs for the requested relation types. If temp_relations |
1687 | | * is True then only temporary relations with storage are returned. If |
1688 | | * temp_relations is False then non-temporary relations with storage are |
1689 | | * returned. If include_shared is True then shared relations are included as |
1690 | | * well in a non-temporary list. include_shared has no relevance when building |
1691 | | * a list of temporary relations. |
1692 | | */ |
1693 | | static List * |
1694 | | BuildRelationList(bool temp_relations, bool include_shared) |
1695 | 0 | { |
1696 | 0 | List *RelationList = NIL; |
1697 | 0 | Relation rel; |
1698 | 0 | TableScanDesc scan; |
1699 | 0 | HeapTuple tup; |
1700 | 0 | MemoryContext ctx = CurrentMemoryContext; |
1701 | 0 | MemoryContext oldctx; |
1702 | |
|
1703 | 0 | StartTransactionCommand(); |
1704 | |
|
1705 | 0 | rel = table_open(RelationRelationId, AccessShareLock); |
1706 | 0 | scan = table_beginscan_catalog(rel, 0, NULL); |
1707 | |
|
1708 | 0 | while (HeapTupleIsValid(tup = heap_getnext(scan, ForwardScanDirection))) |
1709 | 0 | { |
1710 | 0 | Form_pg_class pgc = (Form_pg_class) GETSTRUCT(tup); |
1711 | |
|
1712 | 0 | if (!RELKIND_HAS_STORAGE(pgc->relkind)) |
1713 | 0 | continue; |
1714 | | |
1715 | | /* Only include temporary relations when explicitly asked to */ |
1716 | 0 | if (pgc->relpersistence == RELPERSISTENCE_TEMP) |
1717 | 0 | { |
1718 | 0 | if (!temp_relations) |
1719 | 0 | continue; |
1720 | 0 | } |
1721 | 0 | else |
1722 | 0 | { |
1723 | | /* |
1724 | | * If we are only interested in temp relations then continue |
1725 | | * immediately as the current relation isn't a temp relation. |
1726 | | */ |
1727 | 0 | if (temp_relations) |
1728 | 0 | continue; |
1729 | | |
1730 | 0 | if (pgc->relisshared && !include_shared) |
1731 | 0 | continue; |
1732 | 0 | } |
1733 | | |
1734 | 0 | oldctx = MemoryContextSwitchTo(ctx); |
1735 | 0 | RelationList = lappend_oid(RelationList, pgc->oid); |
1736 | 0 | MemoryContextSwitchTo(oldctx); |
1737 | 0 | } |
1738 | |
|
1739 | 0 | table_endscan(scan); |
1740 | 0 | table_close(rel, AccessShareLock); |
1741 | |
|
1742 | 0 | CommitTransactionCommand(); |
1743 | |
|
1744 | 0 | return RelationList; |
1745 | 0 | } |
1746 | | |
1747 | | /* |
1748 | | * DataChecksumsWorkerMain |
1749 | | * |
1750 | | * Main function for enabling checksums in a single database. This is the |
1751 | | * function set as the bgw_function_name in the dynamic background worker |
1752 | | * process initiated for each database by the worker launcher. After enabling |
1753 | | * data checksums in each applicable relation in the database, it will wait for |
1754 | | * all temporary relations that were present when the function started to |
1755 | | * disappear before returning. This is required since we cannot rewrite |
1756 | | * existing temporary relations with data checksums. |
1757 | | */ |
1758 | | void |
1759 | | DataChecksumsWorkerMain(Datum arg) |
1760 | 0 | { |
1761 | 0 | Oid dboid; |
1762 | 0 | List *RelationList = NIL; |
1763 | 0 | List *InitialTempTableList = NIL; |
1764 | 0 | BufferAccessStrategy strategy; |
1765 | 0 | bool aborted = false; |
1766 | 0 | int64 rels_done; |
1767 | 0 | bool process_shared; |
1768 | | #ifdef USE_INJECTION_POINTS |
1769 | | bool retried = false; |
1770 | | #endif |
1771 | |
|
1772 | 0 | worker_invocation = DatumGetUInt64(arg); |
1773 | |
|
1774 | 0 | operation = ENABLE_DATACHECKSUMS; |
1775 | |
|
1776 | 0 | pqsignal(SIGTERM, die); |
1777 | 0 | pqsignal(SIGUSR1, procsignal_sigusr1_handler); |
1778 | |
|
1779 | 0 | BackgroundWorkerUnblockSignals(); |
1780 | |
|
1781 | 0 | MyBackendType = B_DATACHECKSUMSWORKER_WORKER; |
1782 | 0 | init_ps_display(NULL); |
1783 | |
|
1784 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_SHARED); |
1785 | 0 | if (DataChecksumState->worker_invocation != worker_invocation) |
1786 | 0 | { |
1787 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1788 | 0 | return; |
1789 | 0 | } |
1790 | 0 | dboid = DataChecksumState->database_oid; |
1791 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1792 | |
|
1793 | 0 | BackgroundWorkerInitializeConnectionByOid(dboid, InvalidOid, |
1794 | 0 | BGWORKER_BYPASS_ALLOWCONN); |
1795 | | |
1796 | | /* worker will have a separate entry in pg_stat_progress_data_checksums */ |
1797 | 0 | pgstat_progress_start_command(PROGRESS_COMMAND_DATACHECKSUMS, |
1798 | 0 | InvalidOid); |
1799 | 0 | ResetDataChecksumsProgressCounters(); |
1800 | | |
1801 | | /* |
1802 | | * Get a list of all temp tables present as we start in this database. We |
1803 | | * need to wait until they are all gone before we exit. For the list of |
1804 | | * relations to enable checksums in, check if shared catalogs have been |
1805 | | * processed already. |
1806 | | */ |
1807 | 0 | InitialTempTableList = BuildRelationList(true, false); |
1808 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1809 | 0 | if (DataChecksumState->worker_invocation != worker_invocation) |
1810 | 0 | { |
1811 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1812 | 0 | return; |
1813 | 0 | } |
1814 | 0 | process_shared = DataChecksumState->process_shared_catalogs; |
1815 | | |
1816 | | /* |
1817 | | * Enable vacuum cost delay, if any. While this process isn't doing any |
1818 | | * vacuuming, we are re-using the infrastructure that vacuum cost delay |
1819 | | * provides rather than inventing something bespoke. This is an internal |
1820 | | * implementation detail and care should be taken to avoid it bleeding |
1821 | | * through to the user to avoid confusion. |
1822 | | * |
1823 | | * VacuumUpdateCosts() propagates the values to the variables actually |
1824 | | * read by vacuum_delay_point(). |
1825 | | */ |
1826 | 0 | VacuumCostDelay = DataChecksumState->cost_delay; |
1827 | 0 | VacuumCostLimit = DataChecksumState->cost_limit; |
1828 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1829 | 0 | VacuumUpdateCosts(); |
1830 | 0 | VacuumCostBalance = 0; |
1831 | | |
1832 | | /* |
1833 | | * Create and set the vacuum strategy as our buffer strategy. |
1834 | | */ |
1835 | 0 | strategy = GetAccessStrategy(BAS_VACUUM); |
1836 | |
|
1837 | 0 | RelationList = BuildRelationList(false, process_shared); |
1838 | | |
1839 | | /* Update the total number of relations to be processed in this DB. */ |
1840 | 0 | { |
1841 | 0 | const int index[] = { |
1842 | 0 | PROGRESS_DATACHECKSUMS_RELS_TOTAL, |
1843 | 0 | PROGRESS_DATACHECKSUMS_RELS_DONE |
1844 | 0 | }; |
1845 | |
|
1846 | 0 | int64 vals[2]; |
1847 | |
|
1848 | 0 | vals[0] = list_length(RelationList); |
1849 | 0 | vals[1] = 0; |
1850 | |
|
1851 | 0 | pgstat_progress_update_multi_param(2, index, vals); |
1852 | 0 | } |
1853 | | |
1854 | | /* Process the relations */ |
1855 | 0 | rels_done = 0; |
1856 | 0 | foreach_oid(reloid, RelationList) |
1857 | 0 | { |
1858 | 0 | bool costs_updated = false; |
1859 | |
|
1860 | 0 | if (!ProcessSingleRelationByOid(reloid, strategy)) |
1861 | 0 | { |
1862 | 0 | aborted = true; |
1863 | 0 | break; |
1864 | 0 | } |
1865 | | |
1866 | 0 | pgstat_progress_update_param(PROGRESS_DATACHECKSUMS_RELS_DONE, |
1867 | 0 | ++rels_done); |
1868 | 0 | CHECK_FOR_INTERRUPTS(); |
1869 | 0 | CHECK_FOR_WORKER_ABORT_REQUEST(); |
1870 | |
|
1871 | 0 | if (abort_requested) |
1872 | 0 | break; |
1873 | | |
1874 | | /* |
1875 | | * Check if the cost settings changed during runtime and if so, update |
1876 | | * to reflect the new values and signal that the access strategy needs |
1877 | | * to be refreshed. |
1878 | | */ |
1879 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1880 | 0 | if (DataChecksumState->worker_invocation != worker_invocation) |
1881 | 0 | { |
1882 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1883 | 0 | break; |
1884 | 0 | } |
1885 | 0 | if ((DataChecksumState->launch_cost_delay != DataChecksumState->cost_delay) |
1886 | 0 | || (DataChecksumState->launch_cost_limit != DataChecksumState->cost_limit)) |
1887 | 0 | { |
1888 | 0 | costs_updated = true; |
1889 | 0 | VacuumCostDelay = DataChecksumState->launch_cost_delay; |
1890 | 0 | VacuumCostLimit = DataChecksumState->launch_cost_limit; |
1891 | 0 | VacuumUpdateCosts(); |
1892 | |
|
1893 | 0 | DataChecksumState->cost_delay = DataChecksumState->launch_cost_delay; |
1894 | 0 | DataChecksumState->cost_limit = DataChecksumState->launch_cost_limit; |
1895 | 0 | } |
1896 | 0 | else |
1897 | 0 | costs_updated = false; |
1898 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1899 | |
|
1900 | 0 | if (costs_updated) |
1901 | 0 | { |
1902 | 0 | FreeAccessStrategy(strategy); |
1903 | 0 | strategy = GetAccessStrategy(BAS_VACUUM); |
1904 | 0 | } |
1905 | 0 | } |
1906 | |
|
1907 | 0 | list_free(RelationList); |
1908 | 0 | FreeAccessStrategy(strategy); |
1909 | |
|
1910 | 0 | if (aborted || abort_requested) |
1911 | 0 | { |
1912 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1913 | 0 | if (DataChecksumState->worker_invocation == worker_invocation) |
1914 | 0 | DataChecksumState->worker_result = DATACHECKSUMSWORKER_ABORTED; |
1915 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1916 | 0 | ereport(DEBUG1, |
1917 | 0 | errmsg("data checksum processing aborted in database OID %u", |
1918 | 0 | dboid)); |
1919 | 0 | return; |
1920 | 0 | } |
1921 | | |
1922 | | /* The worker is about to wait for temporary tables to go away. */ |
1923 | 0 | pgstat_progress_update_param(PROGRESS_DATACHECKSUMS_PHASE, |
1924 | 0 | PROGRESS_DATACHECKSUMS_PHASE_WAITING_TEMPREL); |
1925 | | |
1926 | | /* |
1927 | | * Wait for all temp tables that existed when we started to go away. This |
1928 | | * is necessary since we cannot "reach" them to enable checksums. Any temp |
1929 | | * tables created after we started will already have checksums in them |
1930 | | * (due to the "inprogress-on" state), so no need to wait for those. |
1931 | | */ |
1932 | 0 | for (;;) |
1933 | 0 | { |
1934 | 0 | List *CurrentTempTables; |
1935 | 0 | int numleft; |
1936 | 0 | char activity[64]; |
1937 | |
|
1938 | 0 | CurrentTempTables = BuildRelationList(true, false); |
1939 | 0 | numleft = 0; |
1940 | 0 | foreach_oid(tmptbloid, InitialTempTableList) |
1941 | 0 | { |
1942 | 0 | if (list_member_oid(CurrentTempTables, tmptbloid)) |
1943 | 0 | numleft++; |
1944 | 0 | } |
1945 | 0 | list_free(CurrentTempTables); |
1946 | |
|
1947 | | #ifdef USE_INJECTION_POINTS |
1948 | | if (IS_INJECTION_POINT_ATTACHED("datachecksumsworker-fake-temptable-wait")) |
1949 | | { |
1950 | | /* Make sure to just cause one retry */ |
1951 | | if (!retried && numleft == 0) |
1952 | | { |
1953 | | numleft = 1; |
1954 | | retried = true; |
1955 | | |
1956 | | INJECTION_POINT_CACHED("datachecksumsworker-fake-temptable-wait", NULL); |
1957 | | } |
1958 | | } |
1959 | | #endif |
1960 | |
|
1961 | 0 | if (numleft == 0) |
1962 | 0 | break; |
1963 | | |
1964 | | /* |
1965 | | * At least one temp table is left to wait for, indicate in pgstat |
1966 | | * activity and progress reporting. |
1967 | | */ |
1968 | 0 | snprintf(activity, |
1969 | 0 | sizeof(activity), |
1970 | 0 | "Waiting for %d temp tables to be removed", numleft); |
1971 | 0 | pgstat_report_activity(STATE_RUNNING, activity); |
1972 | | |
1973 | | /* Retry every 3 seconds */ |
1974 | 0 | ResetLatch(MyLatch); |
1975 | 0 | (void) WaitLatch(MyLatch, |
1976 | 0 | WL_LATCH_SET | WL_TIMEOUT | WL_EXIT_ON_PM_DEATH, |
1977 | 0 | 3000, |
1978 | 0 | WAIT_EVENT_CHECKSUM_ENABLE_TEMPTABLE_WAIT); |
1979 | |
|
1980 | 0 | CHECK_FOR_INTERRUPTS(); |
1981 | 0 | CHECK_FOR_WORKER_ABORT_REQUEST(); |
1982 | |
|
1983 | 0 | if (aborted || abort_requested) |
1984 | 0 | { |
1985 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
1986 | 0 | if (DataChecksumState->worker_invocation == worker_invocation) |
1987 | 0 | DataChecksumState->worker_result = DATACHECKSUMSWORKER_ABORTED; |
1988 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
1989 | 0 | ereport(LOG, |
1990 | 0 | errmsg("data checksum processing aborted in database OID %u", |
1991 | 0 | dboid)); |
1992 | 0 | return; |
1993 | 0 | } |
1994 | 0 | } |
1995 | | |
1996 | 0 | list_free(InitialTempTableList); |
1997 | | |
1998 | | /* worker done */ |
1999 | 0 | pgstat_progress_end_command(); |
2000 | |
|
2001 | 0 | LWLockAcquire(DataChecksumsWorkerLock, LW_EXCLUSIVE); |
2002 | 0 | if (DataChecksumState->worker_invocation == worker_invocation) |
2003 | 0 | DataChecksumState->worker_result = DATACHECKSUMSWORKER_SUCCESSFUL; |
2004 | 0 | LWLockRelease(DataChecksumsWorkerLock); |
2005 | 0 | } |